Epic: Pipeline SDK (livepeer-runner)
Nobody has claimed this yet.
Assessment
- Difficulty
- 5/5
- Estimated time
- Over a week
- Newbie friendliness
- 25/100
- Issue type
- Feature
- Clarity
- Needs clarification
- Activity status
- Quiet
- Domain
- api, audio-video-rtc, backend
Research direction
Start with the linked pipeline-sdk.md specification and select one outstanding roadmap item rather than treating the whole epic as a single change. Read the relevant entry point or example under examples/runner/, run its existing test.sh or demo, and consider the work done when that selected item is implemented and its stated validation passes.
Written by the indexing model from the issue text.
Description
Outcome
Developers go from "I have a Python ML model" to a discoverable BYOC capability on Livepeer in under 5 minutes — surfaced in the Developer Dashboard ready for any caller to invoke.
The Pipeline SDK is the authoring surface that makes this possible: write a Python class, get a containerised, BYOC-compatible, schema-described capability.
Spec
Design lives in livepeer-specs / pipeline-sdk.md. Update the spec rather than this issue body when the design moves.
Architecture decisions (monorepo + PEP 420 namespace packages, three distributions livepeer-runner / livepeer-client / livepeer-trickle) are captured in the spec — see the Architecture and the companion Client SDK packaging section.
Roadmap
Each step yields a working SDK strictly more capable than the previous one.
- C1 —
Pipelinebase +serve()+ hello-world BYOC E2E - C2 —
setup()lifecycle + HuggingFace sentiment example - C3 — FastAPI HTTP layer (
/health,/predict,/docs,/openapi.json) - C4 — Pydantic
BaseModelfor inputs / outputs via signature introspection - C5 — Image upscale example (binary I/O via Pydantic
Base64Bytes) - C6 —
/healthstate machine matching go-livepeer'sHealthCheckwire format - C7 — SSE auto-detection from generator
predict()+ LLM chat example - C8 —
LivePipelinefor trickle transport (real-time video) — see breakdown below - C9 —
livepeer pushCLI +livepeer.yamlmanifest - C10 — Schema as Docker image label (
org.livepeer.pipeline.schema) - C11 — Agent-friendly docs (
AGENTS.md, expandedPipelinedocstring,examples/runner/_template/) - C12 — Migrate to monorepo with PEP 420 namespace packages — see client SDK packaging spec. Coordinated with #9.
- C13 — Container self-registration to orch
/capability/register— env-gated, wired intoserve()lifespan, lenient on failure (degrade/health, keep FastAPI serving). Deregister on shutdown. Retry on conn-refused / timeout / 5xx; fail fast on 400 / 404 / 405.
C8 breakdown — LivePipeline
- Step 1: skeleton (HTTP routes + ABC) —
04cc697 - Step 2: bytes-through (validate trickle wire) —
831ee44 - Step 3: frame-loop dispatch +
runner.framesnamespace —9688f67 - Step 5:
live_grayscaleexample + chroma assertion + ffplay viewer —9ef95d9series - Step 4 —
_LiveSession+ lifecycle (the only outstanding piece)-
_LiveSessionclass encapsulating per-session state - Periodic heartbeat on
events_url(gateway liveness signal) -
emit_event(payload)user-facing helper -
emit_data(payload)helper fordata_urlwhenenable_data_output=true -
on_stream_stoplifecycle hook -
Drain runner-side state on stop— measured: no leak (RSS plateaus ~170 MB after 25 sessions). State drain unnecessary. - Unified error surface — three error sources (subscribe / publish / user
process_videoraise) log distinctly today with no consistent state propagation. Add_record_error(source, exc, severity): structuredErrorEventschema (severity∈ WARN / ERROR / FATAL,source,message,timestamp,consecutive), per-source budget escalating WARN→ERROR after N consecutive failures, flippipeline._state = ERRORon terminal failures, push events viaevents_url. Quick first step (~5 LOC): flip_stateonTricklePublisherTerminalError. Full schema afterlive_transcribesurfaces real failure modes. - Verify the live-viewer demo can be brought back once heartbeat + state-drain land
- Pydantic param schema for LivePipeline — today
on_stream_start(params)andon_params_update(params)receivedict[str, Any]; users parse / validate manually. The batchPipelinealready supports typed params via signature introspection (C4). Extend the same toLivePipeline: let users typeon_stream_start(params: MyParams)and have the runner introspect, validate at the HTTP boundary, and emit a meaningful/openapi.jsonschema for/stream/start's caller-supplied params (today the schema only describes the orchestrator's protocol fields). Required for the developer dashboard / client SDK to render param controls for live capabilities. - Enrich heartbeat payload with
PipelineStatus(ai-runner pattern) — today's heartbeat is minimal{"type": "heartbeat", "timestamp"}keep-alive only. ai-runner'sreport_status_loopuses the same trickle push as both keep-alive AND status report (state, FPS, last_error, restart_count, last_params). One mechanism, dual purpose. Add FPS counters toMediaOutput/MediaPublishand swap heartbeat payload for the rich shape. Keeps/healthminimal (k8s contract).
-
Companion issues — runner / examples polish
Targeted issues spun off from this epic. Two blocking, one cosmetic for full live_transcribe fidelity (5/5 transcripts delivered to SSE).
Blocking (data-loss bugs)
- #12 — SDK-side.
_resolve_next_seqreturns-1on probe failure; combined with the publisher's+1increment this duplicate-POSTs to seg 0 on every trickle channel, dropping the first record. Observable onlive_transcribeas missingtranscript[1]. Fix: one line (return 0instead of-1) + demote the warning to debug. - Upstream: livepeer/go-livepeer#3924 — gateway's data subscriber tears down too early on
/stream/stop, dropping the finalemit_datafromon_stream_stop. Observable onlive_transcribeas missingtranscript[4](orch log showsclient disconnectedon the final POST). Fix: bounded drain loop inbyoc/trickle.go:startDataSubscribe.
Together these account for both observed transcript losses (runner emits 5 → SSE delivers 3 today). Either one fixed independently → 4/5 delivered.
Cosmetic (log noise, no functional impact)
- Upstream: livepeer/go-livepeer#3922 — spurious ERROR-level logs at every clean
/stream/stop(5 fix sites: ffmpeg subprocess output, trickle preconnect, rtmp2segment probe, orch trickle handler). Operations work; logs just look scary. Pure log-level demotion (if ctx.Err() != nil { debug }).
Examples follow-ups
-
Live viewer tool for
live_grayscale— bring back the webcam-pushed live viewer (deleted because the current PyAV decode→user→encode loop can't sustain real-time webcam load: ring buffer drains, mediamtx kicks the egress publisher). Today the example uses synthetictestsrc+ capture-to-file + replay. Bring back the webcam viewer once C8 Step 4 lands. -
Worked example covering full LivePipeline lifecycle —
live_grayscaleexercises the SDK plumbing but only overridesprocess_video.live_transcribe(Whisper STT) andlive_depth(DepthAnything V2) now exercise more of the lifecycle (setup,on_stream_start,process_audio,emit_data,emit_event,on_stream_stop). Still TODO: atest.shthat subscribes todata_urlfrom the caller side and asserts structured records arrive — needsstart_byoc_jobfrom #6. -
Exercise
on_params_updatein an example — the only LivePipeline hook with no live demo. Fires on mid-stream parameter changes (caller pushes new params to/stream/paramswithout restart). Smallest viable demo: extendlive_detectto accept{"detection_threshold": 0.5}mid-stream and update the YOLO confidence cutoff in-place, with atest.shstep that pushes a new threshold mid-run and asserts the emitted records reflect it. Alternatively document the hook in the SDK README and defer the example until a real use case demands it. -
Migrate
/stream/paramsand/stream/stoptocontrol_urlsubscribe — per the spec, the long-term shape is one HTTP endpoint (/stream/start) plus everything else over the trickle plane. Today BYOC already publishes params + keepalives tocontrol_url(byoc/trickle.go:539) but the orchestrator HTTP-forwards each message via/stream/params(byoc/stream_orchestrator.go:421). Migration must be coordinated upstream: orchestrator drops the HTTP-forwarding step + runner adds trickle subscribe in lockstep. Blocked on: trickle control-channel size / changeover bug (see "Future protocol work" below). -
Production-grade live transcribe example —
live_transcribeis intentionally the minimal lifecycle demo with explicit 3 s chunking +vad_filter=True; first-transcript latency is ~3 s and word boundaries can split mid-window. For users who actually want production live transcription, add a separate example usingwhisper_streaming'sOnlineASRProcessor(LocalAgreement-2 → ~1 s latency, cleaner boundaries) — sameLivePipelinelifecycle, differentprocess_audiobody. Same folder shape aslive_transcribe, marketed as "production transcribe". Possibly also covers VAD-driven segmentation and emit_data partial-vs-final transcript distinction. -
Recover from orchestrator capability drop — gateway sometimes drops the orchestrator from its capability pool after stream failures (
Retrying stream with a different orchestrator err=unknown swap reason→no orchestrators available, ending stream). Once dropped, every subsequent/process/stream/starteither 400s or kills mid-flight, untilregister_capabilityis re-run manually. Investigate (a) re-register watchdog, (b) healthcheck-driven re-register hook, or (c) push a fix upstream in go-livepeer's gateway swap-orch logic. -
Switch examples to
-network offchainonce go-livepeer #3906 lands. Current compose files run with-network arbitrum-one-mainnet -ethUrl https://arb1.arbitrum.io/rpc -ethPassword secret-passwordand rely onpricePerUnit=0so no real on-chain payment occurs — but the gateway still polls Arbitrum for orchestrator stake lookups (db_discovery.go), and the public RPC throttles with429 Too Many Requestslines all over the gateway log. Tracked upstream as livepeer/go-livepeer#3905. When the PR merges, drop-network,-ethUrl,-ethPasswordfrom each example'sdocker-compose.yml(5 files) and run with bare-network offchain. Eliminates the 429 noise entirely. -
Assert grayscale, not just bytes-received —
live_grayscale/test.shnow extracts U / V plane averages viaffprobe signalstatsand asserts ≈128 (chroma-zero = grayscale). -
End state: retire
examples/runner/, replace with unit tests, move worked examples to a separate repo — once the SDK stabilizes, delete the in-treeexamples/runner/folder. The lifecycle/coverage value those examples currently provide (setup, on_stream_start, process_video / process_audio, emit_data, emit_event, on_stream_stop, error paths) gets reified as proper unit tests insidelivepeer-python-gateway. The worked examples themselves (live_grayscale, live_transcribe, live_depth, replicate_flux, …) move to a standalonelivepeer/pipeline-examplesrepo that depends on the publishedlivepeer.runnerpackage as a normal pip dependency. Aligns examples with how external developers actually consume the SDK, decouples example evolution from SDK release cadence, and keeps this repo focused on the runner itself.
Performance & future improvements
Captured while building examples that surfaced specific optimization
opportunities. Not roadmap-blocking; revisit when concrete use cases
demand them.
-
Parallel
process_video/process_audioexecution — the frame
loop today dispatches both hooks sequentially in a single async for-loop,
so heavy inference in one stalls the other. Refactor into two queues fed
from one decoder, drained in separate tasks. Not needed for
live_detector
live_transcribetoday; triggered
by a real pipeline that needs it (e.g. a Moondream2-class video model
running alongside whisper). Contract change ("frames may interleave across
hooks") so deserves explicit design before flipping. -
live_describe— VLM-driven video understanding (GPU) — natural-
language scene description via Moondream2
(~1.6 GB, ~300 ms / inference on GPU). SameLivePipelinelifecycle as
live_detect, swaps YOLO for a vision-
language model that emits descriptions instead of bounding boxes. Mirrors
live_depth's GPU pattern. Compelling
for "real video understanding" positioning; not strictly needed since
live_detectalready demonstrates multi-modal LivePipeline. -
Concurrent inference patterns documented in SDK README — three-tier
pattern users adopt as inference cost grows: (1)asyncio.to_threadfor
offloading individual inference calls so the main loop stays responsive,
(2)asyncio.create_taskfor fire-and-forget windowed work
(transcribe → emit when done, decouples inference latency from frame
cadence), (3) separate process / IPC for GPU-isolated heavy models. Should
land alongside the parallelprocess_*refactor above so users understand
which knob to reach for.
Runner SDK code-quality improvements
Findings from a focused code review of src/livepeer_gateway/runner/. The
Bugs items are real correctness issues worth fixing before C13 (auto-
registration) lands so we don't bake them into a lifecycle path.
Bugs
- Async-generator pipelines silently broken (
serve.py:225) —inspect.isgeneratorfunction()returnsFalseforasync defgenerators. A user writingasync def run(self, ...): yield ...falls into the sync-call branch;StreamingResponsecan't iterate the resulting async-gen object. Fix: detectinspect.isasyncgenfunction()and use an async SSE formatter, or reject async generators with a clear error. - Sync
run()blocks the event loop (serve.py:84) —pipeline.run(...)is invoked directly inside an async FastAPI handler. Any CPU/IO-boundrun()(sentiment, every HF pipeline) stalls/healthand concurrent requests. Fix:await asyncio.to_thread(pipeline.run, ...)whenrunis sync; keep the direct call only forasync defor generators. -
/stream/startconcurrent-session race (serve.py:131) — two concurrent calls both pass the_session is Nonecheck, both construct_LiveSession, second overwrites first → orphaned tasks + publishers. Wrap session creation with anasyncio.Lockon the pipeline. -
result.frameAttributeError on raw PyAV return (live_pipeline.py:353) — a user who returns a rawav.VideoFrame(natural after PyAV work) hitsAttributeErrorbecause the SDK expects the wrapper. Either accept both (getattr(result, "frame", result)) or document+enforce the wrapper. -
on_params_updatesemantics undocumented — delta or full replacement?live_tintreads as delta;serve.py:209does full replace. Pick one and document. Compounding:session.paramsis mutated BEFORE the hook runs, so a raising hook leaves partial state. Roll back on exception.
Code quality
- Private cross-module imports (
serve.py:13-22) —_LiveSession,_run_frame_loop,_has_user_processingimported across module boundaries with underscore prefix. They're not private anymore; either drop the underscore or move the/stream/*handler factory intolive_pipeline.py. - Duplicate introspection (
live_pipeline.py:288-299) —_emit_flagsand_has_user_processingrecompute the sameprocess_*hook overrides. Collapse into one function returning(emit_video, emit_audio, has_user_processing). -
**kwargssilently swallowed (serve.py:38-56) — a user writingdef run(self, **kwargs)(allowed by the ABC) gets an empty input model and silently loses request body. Either rejectVAR_KEYWORD/VAR_POSITIONALwith a clear error, or treat**kwargsasextra="allow"on the generated model. - OpenAPI schema misses non-
BaseModelreturns (serve.py:230-231) —list[Foo]orFoo | Noneis silently dropped from the schema. Widen detection viapydantic.TypeAdapter, or document the limitation. -
/healthreads private_stateacross two classes (serve.py:105) —pipeline._stateis a private attribute onPipelineandLivePipeline(two unrelated classes that happen to share the name). Lift into a shared base orProtocol, or expose a publicstateproperty.
Naming / readability (bundle with C8 perf refactor)
- Rename
MediaOutput/MediaPublish→media_in/media_out— current names describe verbs (output, publish) but trickle direction is the opposite, which is a foot-gun every time someone reads the code. Cross-cutting refactor; bundle with the parallelprocess_video/process_audiowork in Performance & future improvements since both touch the frame loop. - Unify
/stream/*response shapes —{"status": "started", "gateway_request_id": ...}vs{"status": "ok"}vs{"status": "stopped"}across handlers. Pick one Pydantic response model.
Framework adapters (deferred — build on demand)
Migration paths for users from existing ML frameworks. Each ships as its own pip package with its own foreign dep, isolated from core SDK. Build only when a real migration ask shows up.
-
livepeer-runner-cog— wrapscog.BasePredictor -
livepeer-runner-fal— wrapsfal.App -
livepeer-runner-modal— wraps Modal@app.function -
livepeer-runner-bentoml— wraps@bentoml.service -
livepeer-runner-confyscript
Future protocol work (cross-team, gated on C9 + upstream go-livepeer)
-
Fix trickle control-channel size / segment-changeover bug (upstream) —
control_urlparams updates fail silently or get truncated when payload is more than small JSON (~1 MB practical ceiling observed). Hunch is segment-changeover behavior during large writes — possiblyFirstByteTimeout, pipe buffering on segment boundaries, or chunk-write semantics across the rollover. Workaround in byoc/stream_gateway.go:1007-1009 switched stop / params to HTTP POST after the bug bit on base64-binary payloads. Blocking dependency for the "migrate to control_url subscribe" follow-up. Upstream go-livepeer change. -
Capability identity via OCI digest — replace free-form capability names with content-hashed references like
byoc/<repo>@sha256:<digest>. Aligns BYOC with Replicate's reproducibility model. SDK side:livepeer pushcaptures the digest at publish time and bakes it into the manifest. Upstream side: orchestrator registration + gateway routing +OrchestratorInfocarry the digest. -
Cosign / Sigstore signing of capability digests — optional layer on top of digest pinning. Publisher signs the digest, gateway verifies signature against publisher's key.
-
Name:version aliases over digest-pinned wire — Replicate-style mutable names (
byoc/text-reverser:v2) that resolve to a digest at lookup time. Wire protocol always pins the digest; aliases are a UX layer.
Related
- Spec: pipeline-sdk.md
- Spec follow-ups:
livepeer-specs/_followups.md - Companion epic: #9 (Client SDK — request-side)
- Draft PR: #7
- Caller-side BYOC SDK PR: #6
- Upstream BYOC offchain: livepeer/go-livepeer#3905, livepeer/go-livepeer#3906
- Prior art: livepeer/ai-runner — pioneered the Pydantic-class I/O pattern
- Monorepo + namespace package precedents: Apache Airflow providers, Google Cloud SDK, Azure SDK for Python
uvworkspace docs: https://docs.astral.sh/uv/concepts/projects/workspaces/- PEP 420 (namespace packages): https://peps.python.org/pep-0420/
- Dominant language
- Python
- Stars
- 1
- Forks
- 7
- PR merge metrics
- No merged PRs in 30d
Getting set up
This project ships no dev container, Dockerfile or contributing guide, so setting up is up to you: start from its README, and see our first-contribution guide for the general steps.
First steps
- Read the whole issue, then the project's contributing guide.
- Comment on the issue to say you are picking it up — it saves two people doing the same work.
- Fork the repository and make your change on a branch.
- Open a pull request that references the issue number.
More from livepeer/livepeer-python-gateway
-
Difficulty 2/5 1-3 hours Newbie friendliness 78/100
livepeer/livepeer-python-gateway#64 · 1 comment ·
-
Improvement
Difficulty 1/5 Under an hour Newbie friendliness 68/100
livepeer/livepeer-python-gateway#27 · 1 comment ·
-
Difficulty 2/5 1-3 hours Newbie friendliness 76/100
-
Difficulty 1/5 Under an hour Newbie friendliness 70/100
livepeer/livepeer-python-gateway#12 · 1 comment ·
-
Difficulty 5/5 Over a week Newbie friendliness 35/100
All issues in livepeer/livepeer-python-gateway
Similar issues
-
New InternshipOpennew_internship
Difficulty 1/5 Under an hour Newbie friendliness 70/100
-
[BUG] Reports tab: "Unban" button tooltip shows raw `{{ip}}` placeholder instead of the IP addressOpenbug javascript ui
Difficulty 2/5 1-3 hours Newbie friendliness 68/100
bunkerity/bunkerweb#4001 · 1 comment ·
Maintainers usually reply within 1 day
-
bug
Difficulty 1/5 Under an hour Newbie friendliness 92/100
PedestrianDynamics/pyFDS-Evac#476 ·
Maintainers usually reply within 1 day
-
Difficulty 2/5 1-3 hours Newbie friendliness 72/100
google/differential-privacy#516 ·
-
Difficulty 2/5 1-3 hours Newbie friendliness 82/100
adobe-fonts/source-serif#153 ·