Skip to content

Latest commit

 

History

8 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Adaptive Stream

A video platform that ingests an upload, transcodes it into a multi-resolution HLS ladder across a horizontally scalable worker fleet, and streams it back with adaptive bitrate switching driven by real measured bandwidth.

Built as a learning vehicle for distributed-systems patterns — Kafka partitioning, idempotent consumers, at-least-once delivery, and stateless worker scaling — rather than as a video-codec exercise. ffmpeg and hls.js do the media work; the interesting parts are the pipeline around them.

Status Increments 1–3 complete and working end to end; paused
Built August 2026
Languages TypeScript (API, web) · Rust (workers) · SQL
Size ~3.5k lines of source across three services
Tests 109 (66 API · 25 worker · 18 web)
Runs on docker compose up — one command, six services

A glimpse

The player adapting from 1080p down to 240p as the network is throttled

Real playback, real time, no editing. The player is left on Auto: it starts conservatively at 240p, measures the connection and climbs to 1080p, then the network is simulated at 400 kbps — and one buffer-length later it steps back down to 240p on its own. Watch the estimated bandwidth settle to 0.40 Mbps, matching the simulated rate, and the switch log record each decision.

The lag between throttling and the visible switch is not sluggishness; it is the player draining the ~12 seconds of 1080p it had already downloaded.

  Upload ──▶ Kafka ──▶ Rust workers (× N) ──▶ HLS ladder ──▶ Adaptive playback
                        one rendition each      240p…1080p     auto or manual

Drop a video into the browser. The API streams it to disk and publishes an event. A Rust worker probes it, decides a bitrate ladder capped at the source resolution, and fans out one job per rendition. Any worker picks up any job, so five renditions transcode in parallel across the fleet. As each finishes it reports back; when the last one lands, the API assembles the master playlist and the video becomes playable.

In the player you can pin a resolution, or leave it on Auto and watch hls.js climb and drop the ladder as bandwidth changes — with an overlay showing the bandwidth estimate actually driving each decision.

Architecture

┌──────────────────────────────────────────────────────────────────────┐
│  React + TypeScript + hls.js                                         │
│  upload · library with live status · player with ABR controls        │
└──────┬────────────────────────────────────────┬──────────────────────┘
       │ REST (multipart, polling)              │ HLS playlists + segments
       ▼                                        ▼
┌──────────────────────────────────────────────────────────────────────┐
│  Node API — Express                    sole writer of pipeline state │
│  ┌───────────┬──────────────┬────────────────┬───────────────────┐   │
│  │ REST      │ HLS delivery │ Kafka producer │ Kafka consumer    │   │
│  │ routes    │ + throttle   │                │ group: api-state  │   │
│  └───────────┴──────────────┴────────────────┴───────────────────┘   │
└──────┬───────────────┬──────────────────────┬────────────────────────┘
       │ SQL           │ fs (storage seam)    │ produce / consume
       ▼               ▼                      ▼
┌────────────┐  ┌──────────────┐  ┌──────────────────────────────┐
│ PostgreSQL │  │ /data volume │  │ Kafka (KRaft) · 5 topics     │
│ source of  │  │ raw/ · hls/  │  │                              │
│   truth    │  └──────┬───────┘  └──────────────┬───────────────┘
└────────────┘         │                         │
                       ▼                         ▼
              ┌──────────────────────────────────────────────┐
              │  Rust worker × N  (tokio · rdkafka · ffmpeg)  │
              │  analyzer role  ·  transcoder role            │
              │  stateless — holds no database connection     │
              └──────────────────────────────────────────────┘

Pipeline

POST /api/videos ──▶ row ──▶ [video.uploaded]
                                    │
                          Rust worker · analyzer
                          ffprobe → bitrate ladder
                            │              │
              [transcode.jobs] × N    [video.analyzed]
                    │  │  │
              worker worker worker      ← one rendition each
              ffmpeg ffmpeg ffmpeg
                    │  │  │
                 [transcode.events]
                          │
                 Node consumer (api-state)
                 upsert job row; when every rendition is DONE
                 → master.m3u8, video READY

Tech stack, and why

Layer Choice Why this one
API Node.js + Express + TypeScript The API is almost entirely I/O — streaming uploads to disk, database writes, Kafka fan-in. Node's event loop is a good fit and the ecosystem for multipart handling and Kafka clients is mature.
Workers Rust (tokio + rdkafka) Transcoding is CPU-bound and long-running. Rust gives predictable memory under sustained load and no GC pauses mid-job, and the worker is a natural place to own a subprocess pipeline.
Queue Apache Kafka (KRaft) Needed consumer groups for work distribution, partitions for parallelism, and offset semantics for at-least-once delivery. A simpler broker would have hidden the exact concepts this project exists to learn. KRaft mode means one container, no ZooKeeper.
Database PostgreSQL + Prisma Postgres for real constraints — the UNIQUE (video_id, rendition) index is the pipeline's idempotency anchor, not a decoration. Prisma for typed access and a migration story that stays in the repo.
Frontend React + TypeScript + Vite + hls.js hls.js is the reference client-side HLS implementation and already contains a well-tuned ABR algorithm. Writing my own would have been re-implementing a solved problem instead of learning the pipeline.
Media ffmpeg / ffprobe (CLI) Deliberately off the shelf. Codec internals were explicitly not a goal; the CLI is invoked as a subprocess so the CPU burn never blocks the async runtime.
Orchestration Docker Compose Six services with health gates and a shared volume, reproducible with one command, and --scale worker=N demonstrates horizontal scaling without a Kubernetes detour.

The polyglot split is the point

Node handles what is I/O-bound; Rust handles what is CPU-bound. The boundary between them is a message queue rather than a function call, which is what makes the two halves independently scalable — and what forced every interesting correctness question in this project.


Key design decisions

1. Partition keys are chosen per topic, and they differ on purpose

The single most consequential decision in the system.

Topic Key Why
transcode.jobs videoId:rendition scatters one video's renditions across partitions
transcode.events videoId gathers one video's events onto a single partition

Keying transcode.jobs by videoId alone would send every rendition of a video to the same partition, therefore the same consumer, therefore sequential transcoding — the fleet would silently degrade to one worker per video while still looking healthy.

transcode.events goes the other way deliberately. Because all of a video's events land on one partition, one consumer processes them serially, so "was I the last rendition to finish?" is answered by partition ordering plus a guarded SQL update. Kafka's ordering guarantee is used as the mutex — no distributed lock, no leader election, no advisory lock in Postgres.

Partition count is also the parallelism ceiling: six partitions on transcode.jobs means a seventh worker replica sits idle with no assignment.

2. Workers hold no database connection

Workers consume Kafka, read and write a shared volume, and emit events. Every state transition is written by the API's api-state consumer group.

This keeps one schema owner and one migration story rather than duplicating the data model across two languages, and it means scaling to twenty workers is not a connection-pool conversation.

3. Self-describing messages instead of enforced ordering

The analyzer produces job messages directly rather than routing them back through the API. That removes a network round trip, but transcode.jobs and transcode.events are separate topics with no ordering between them — a fast 240p rendition can finish before the API has processed video.analyzed.

Two changes make that safe:

  • Job rows are upserted, so a completion can create the row it needs.
  • Every job and event carries totalRenditions, so completion detection never depends on having seen the analysis first.

Out-of-order tolerance replaced enforced ordering. It is also the more realistic pattern: at any real throughput you cannot assume the control plane keeps pace with the data plane.

4. At-least-once delivery, and idempotency everywhere it lands

Offsets are committed only after work completes, so redelivery is normal rather than exceptional. Every consumer is built for it:

  • UNIQUE (video_id, rendition) makes duplicate job rows impossible.
  • Output paths are deterministic, so a re-run overwrites rather than duplicates.
  • Job status is ranked (QUEUED < RUNNING < DONE = FAILED) so a late progress event cannot drag a finished rendition backwards.
  • Renditions are written to a staging directory and rename()d into place — atomic on POSIX, so a crashed worker never leaves a half-written playlist where a player could fetch it.

5. "Unprocessable" and "failed to process" are different failures

A consumer that swallows every error commits the offset and silently loses data. A consumer that retries every error blocks its partition forever on a poison pill. So the two are classified:

  • Schema or JSON failure → the message will never become valid; log and skip.
  • Anything else (database outage, a bug) → rethrow, leave the offset uncommitted, let it redeliver.

This one was learned the hard way — see Key learnings.

6. Analysis is a cheap gate in front of expensive work

ffprobe runs in seconds and reads only container headers. Rejecting a corrupt file there costs milliseconds; discovering the same problem after dispatching five ffmpeg jobs costs CPU-minutes. It also decides the ladder, capped at the source resolution, so a 480p upload never produces an upscaled "1080p" that is pure waste and a dishonest quality label.

Keeping ANALYZING and TRANSCODING as distinct states matters for the same reason: an analysis failure means your file is broken (terminal, re-upload), while a transcode failure means our worker had a problem (retryable). One merged PROCESSING state throws away the distinction a user actually needs.

7. A storage seam from day one

Both the API and the workers reach storage through a narrow interface with a disk implementation today and an S3 implementation planned. Object keys are identical strings in both — hls/{id}/720p/index.m3u8 is a filesystem path now and an object key later — so the database never learns which backend it is on.

The honest limitation this creates: because media lives on a shared Docker volume, workers scale horizontally only within a single host. Removing that boundary is exactly what the object-storage phase is for. It is a deliberate trade, not an oversight.

8. A reconciler instead of a transactional outbox

Inserting the video row and producing video.uploaded are not atomic; a crash between them strands a video with no work queued. The textbook fix is a transactional outbox.

That was deliberately deferred. Instead, a reconciler re-produces anything left sitting in UPLOADED past a threshold — and because every consumer is already idempotent, a duplicate costs nothing. The outbox is a later upgrade rather than a prerequisite, and knowing which problem it solves is more valuable than having built it reflexively.

9. Single-step upload with three cleanup paths

POST /api/videos takes the file and creates the record in one call. The UUID is generated in the API before any bytes move, so the upload streams straight to its final path with no temp-then-rename.

Orphaned bytes are handled at three levels: an abort handler for a client that disconnects mid-upload, the error handler for a failed request, and a periodic sweeper for the crash window between "file written" and "row committed."

originalFileName is stored as metadata and reduced to a basename; on-disk names derive from the videoId and a whitelisted extension, so an upload named ../../etc/passwd.mp4 still lands inside the media root.

10. Keyframe alignment, decided long before playback

Every rendition is forced to cut at identical timestamps (-force_key_frames expr:gte(t,n_forced*4) with a fixed GOP). Without it, a player cannot switch renditions cleanly — and the failure stays completely invisible until someone tries to change quality mid-playback, long after the transcoding code looked finished.

Verified by segment count: a 90-second source produces exactly 23 segments in every rendition.

11. Making adaptation demonstrable, not just functional

Local bandwidth is effectively infinite, so ABR pins to the top rendition and never switches — the headline feature would look broken in a demo. Two things fix that:

  • A dev-only server-side throttle paces segment responses to a chosen bitrate. Because hls.js derives its estimate from real download times, this drives the genuine ABR path rather than faking a switch. It is gated behind an env flag, because an endpoint that lets any caller hold a connection open at 50 kbps is a denial-of-service primitive.
  • A capped buffer. maxBufferLength alone is only a target; hls.js grows toward maxMaxBufferLength (600s by default), so a short video downloads entirely up front and bandwidth is never re-evaluated again.

Observed on a 90s 1080p source: auto started at 240p, climbed to 1080p, then dropped back to 240p within one segment of throttling to 400 kbps.

12. A shared contract, enforced by both languages

contracts/ holds one JSON fixture per event. The Node suite parses them with zod and the Rust suite parses them with serde. A schema change that breaks the other language fails a test rather than production — which is the main risk a polyglot pipeline introduces.


Setup

Prerequisites

Only Docker is required to run the stack. For working on the code directly:

Tool Version Needed for
Docker + Compose v2 everything
Node.js 24+ API and web development, running their tests
Rust 1.85+ (edition 2024) worker development
cmake any building rdkafka, which compiles librdkafka from source

Run it

git clone <repo-url> adaptive-stream
cd adaptive-stream
docker compose up -d --build

First build takes a few minutes — the worker image compiles librdkafka from source. Then:

Service URL
Web UI http://localhost:5173
API http://localhost:3000
Health http://localhost:3000/health
Postgres localhost:5432 (adaptive / adaptive)
Kafka localhost:29092 (external listener)

Upload a video through the UI and watch the status move UPLOADED → TRANSCODING → READY, with per-rendition progress. When it turns green, hit Play.

Scale the workers

docker compose up -d --scale worker=5

Each worker transcodes one rendition at a time, so throughput comes from replicas. Confirm the partitions actually spread:

docker compose exec kafka /opt/kafka/bin/kafka-consumer-groups.sh \
  --bootstrap-server kafka:9092 --describe --group worker-transcoder

Test videos

Generated by the worker image's own ffmpeg, so nothing needs installing:

./scripts/make-fixtures.sh

Use long-1080p.mp4 for the adaptive-streaming demo — a short clip has too few segments for ABR to visibly react.

Re-recording the demo GIF

The GIF at the top is captured from the running stack by driving headless Chrome against the real UI — the level switches in it come from hls.js reacting to genuinely throttled downloads, not from a script that pretends to switch.

cd scripts/capture && npm install
OUT_DIR=/tmp/abr-frames npm run record

It needs the stack up with ENABLE_BANDWIDTH_THROTTLE=true, a READY 90s video, and Chrome installed (CHROME_PATH overrides the default macOS location). Plain Chromium will not do: it lacks the H.264 decoders these renditions need.

Tests

cd api && npm install && npm test        # 66 — needs Docker for testcontainers
cd web && npm install && npm test        # 18
cd worker && cargo test                  # 25

The API suite starts a throwaway PostgreSQL container and applies real migrations to it, so constraints, enums and defaults behave truly rather than through a mock.

Configuration

Variable Default Purpose
DATABASE_URL — Postgres connection string
KAFKA_BROKERS localhost:9092 broker list
MEDIA_ROOT /data shared media volume
MAX_UPLOAD_BYTES 2 GiB upload limit, returns 413
SWEEPER_INTERVAL_MS 15 min orphan sweep / reconcile cadence
RECONCILE_STUCK_AFTER_MS 5 min age before an UPLOADED video is re-produced
ENABLE_BANDWIDTH_THROTTLE false dev only — enables the network simulator
WORKER_ROLES analyzer,transcoder roles this worker runs
MAX_ATTEMPTS 3 transcode retries before the DLQ
FFMPEG_PRESET veryfast encoder speed/quality trade

Layout

api/         Express + Prisma — REST, HLS delivery, pipeline state consumer
worker/      Rust — analyzer and transcoder roles
web/         React + hls.js — upload, library, adaptive player
contracts/   Event schemas + fixtures parsed by both languages' tests
docs/        Design spec
scripts/     Test-video generation and demo capture

Key learnings

Partition keys are a design decision, not a detail. The difference between a fleet that transcodes in parallel and one that quietly serialises is a single string. Worse, the broken version looks completely healthy — consumers connected, lag at zero, work getting done, just N times slower than it should be. Nothing alerts on "your parallelism silently collapsed."

Ordering guarantees can replace coordination. I expected completion detection ("was I the last rendition?") to need a lock. It doesn't: keying events by video id puts them on one partition and one consumer, which serialises them for free. Choosing a key was cheaper and more reliable than any locking scheme I would have written.

Swallowing errors in a consumer is data loss with extra steps. My first consumer caught everything and logged it, on the reasoning that a poison pill shouldn't block a partition. Then a stale Prisma client made every write fail — and because the handler swallowed the error, kafkajs committed the offsets and five completion events were permanently gone. The video sat in TRANSCODING forever with no trace of why. Distinguishing unprocessable from failed to process is not a nicety; getting it wrong turns a recoverable outage into silent corruption.

At-least-once forces you to design idempotency up front. Every consumer assumes it will see the same message twice, and that assumption shaped the schema — the unique constraint, the deterministic paths, the status ranking. Retrofitting this onto a system that assumed exactly-once would have been a rewrite.

Some bugs are invisible until the last mile. The hardcoded CODECS string was wrong for three of five renditions (240p is H.264 level 2.1, 1080p is 4.0) and nothing failed — playlists generated, files written, tests green. Keyframe alignment is the same shape: get it wrong and everything looks finished until a player tries to switch. Both were only caught by inspecting real output rather than trusting green tests.

Local development hides the problem your product solves. On localhost there is no bandwidth constraint, so adaptive bitrate has nothing to adapt to. Building the throttle was not polish — without it I could not tell working ABR from broken ABR, and neither could anyone I showed it to.

Bind mounts and baked images drift. The dev image generated its Prisma client at build time while prisma/ was mounted at runtime, so a schema change on the host left the container's client stale and every write failed with a confusing error. Anything derived from a mounted file has to be derived at boot, not at build.


Roadmap

# Increment Status
1 Upload to disk, video CRUD, React library ✅ done
2 Kafka pipeline, Rust transcoding workers, HLS ladder ✅ done
3 HLS delivery, hls.js player, manual + automatic ABR ✅ done
4 Live progress over SSE, richer failure surfacing planned
5 Object storage (S3/MinIO), removing the single-host limit planned

Known limitations

Stated plainly, because they are deliberate trades rather than surprises:

  • Single-host scaling. Media on a shared volume means workers scale within one machine. Increment 5 removes this.
  • No authentication. Anyone with a video's UUID can fetch it.
  • No test drives a real broker. Handlers are tested against real Postgres and the Rust logic is unit-tested, but the Kafka round trip is verified only by manual end-to-end runs.
  • The retry/DLQ path is untested. It is implemented and reviewed, but no fault-injection test exercises it.
  • Polling, not push. The library page polls every two seconds; SSE is Increment 4.
  • No transactional outbox. Covered by the reconciler, as described above.

About

An event driven adaptive streaming application for streaming videos (multiple resolutions) based on client's bandwidth.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages