Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions .github/workflows/build-rtc.yml
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ jobs:
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
submodules: true
repository: ${{ github.event.pull_request.head.repo.full_name }}
ref: ${{ github.event.pull_request.head.ref }}

- uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7
Expand All @@ -48,7 +49,13 @@ jobs:
- name: generate python stubs
run: ./generate_proto.sh

# Fork PR tokens cannot push generated changes. Verify them read-only.
- name: Check generated stubs on forks
if: github.event.pull_request.head.repo.full_name != github.repository
run: git diff --exit-code -- livekit/rtc/_proto

- name: Add changes
if: github.event.pull_request.head.repo.full_name == github.repository
uses: EndBug/add-and-commit@cc9c08ba6c8df3b93a8f2db63e89b98368ae2ae8 # v11
with:
add: '["livekit-rtc/"]'
Expand Down
6 changes: 3 additions & 3 deletions .github/workflows/tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -100,7 +100,7 @@ jobs:
uv venv .test-venv
source .test-venv/bin/activate
uv pip install rtc-wheel/*.whl ./livekit-api ./livekit-protocol
uv pip install pytest pytest-asyncio numpy matplotlib
uv pip install pytest pytest-asyncio numpy matplotlib "av>=16.1; python_version >= '3.10'"

- name: Create venv and install dependencies (macOS)
if: runner.os == 'macOS'
Expand All @@ -109,15 +109,15 @@ jobs:
source .test-venv/bin/activate
uv pip install "${{ steps.select-wheel-macos.outputs.wheel }}"
uv pip install ./livekit-api ./livekit-protocol
uv pip install pytest pytest-asyncio numpy matplotlib
uv pip install pytest pytest-asyncio numpy matplotlib "av>=16.1; python_version >= '3.10'"

- name: Create venv and install dependencies (Windows)
if: runner.os == 'Windows'
run: |
uv venv .test-venv
$wheel = (Get-ChildItem rtc-wheel\*.whl)[0].FullName
uv pip install --python .test-venv $wheel .\livekit-api .\livekit-protocol
uv pip install --python .test-venv pytest pytest-asyncio numpy matplotlib
uv pip install --python .test-venv pytest pytest-asyncio numpy matplotlib "av>=16.1; python_version >= '3.10'"
shell: pwsh

- name: Run tests (Unix)
Expand Down
67 changes: 67 additions & 0 deletions examples/preencoded-webm/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
# Pre-encoded WebM streaming

Publish VP8/VP9 WebM from a file or streaming HTTP(S) response using the Python
pre-encoded video source. PyAV only demuxes packets: there is no `decode`, raw
pixel conversion, or encoder in the publishing path. Audio is outside this example.

The WebM example requires Python 3.10+ and PyAV 16.1+ for packet side-data access.
The underlying pre-encoded SDK API also supports Python 3.9.

From the repository root (with the bundled FFI installed):

```sh
uv sync --dev
uv pip install -r examples/preencoded-webm/requirements.txt
export LIVEKIT_URL=ws://localhost:7880
export LIVEKIT_API_KEY=devkey
export LIVEKIT_API_SECRET=secret
.venv/bin/python examples/preencoded-webm/publish.py input.webm --room demo
# The input can also be an HTTP(S) URL producing a live WebM response.
```

Use WebM with a keyframe at the start, increasing timestamps, fixed dimensions,
and no B-frame reordering. The example supports VP8 and VP9; receiver codec/profile
support still matters. Reads run off the asyncio loop, one packet at a time, with
network timeouts and bounded probing so the first frame does not wait for EOF.
Presentation timestamps drive pacing, including variable frame-rate input. The
publisher waits up to 30 seconds for a subscriber and briefly repeats the initial
keyframe until the native encoder reports its first rate target (at most five
seconds). This avoids consuming the only keyframe of a short clip before encoder
initialization; it is not a general keyframe-recovery mechanism. At EOF the track
stays published until Ctrl-C, since capture acceptance is not a delivery or drain
acknowledgement.

The demuxer cannot satisfy a request for a new keyframe or lower bitrate; feedback
is logged and must be connected to an upstream encoder in a production integration.
Use short GOPs and an appropriate bitrate for a file demonstration. Late joiners
and packet loss may require waiting until the next keyframe. This example does not
provide adaptive encoding or a loss-recovery policy.

## Transparency and container-preserving transport

This example publishes opaque video as a standard VP8/VP9 media track. It removes
WebM container framing, so it rejects alpha and BlockAdditional side data instead
of silently losing them. It does not define a custom alpha payload or transport.

To preserve transparent WebM, carry its original bytes, including the WebM header,
track metadata, Clusters, and BlockAdditional elements. LiveKit's existing
[byte-stream API](https://docs.livekit.io/transport/data/byte-streams/) can carry
such bytes incrementally, but is a data stream rather than a video media track.
An ordinary VP9 RTP media track does not accept a WebM container as its payload.

A future client adapter could feed an intact WebM byte stream into native media
playback without application-side decoding or re-encoding. Streaming playback,
alpha support, and latency still need validation on each target browser/mobile
platform. The adjacent [WebM byte-stream example](../webm-byte-stream/) forwards the
container intact. No client playback adapter is included.

## Validation

```sh
uv run --with 'av>=16.1' pytest tests/rtc/test_encoded_video.py tests/rtc/test_webm_example.py
```

Tests generate synthetic WebM, verify unchanged compressed video bytes and
key/delta timestamps, reject alpha, read the first frame before a pipe or HTTP
response reaches EOF, and check startup and cancellation cleanup. The base API's
server-backed test publishes VP9 and verifies subscriber pixels.
130 changes: 130 additions & 0 deletions examples/preencoded-webm/publish.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
"""Publish compressed WebM video from a file or a streaming HTTP response."""

from __future__ import annotations

import argparse
import asyncio
import logging
import os
import time
from dataclasses import replace
from typing import Iterator, TypeVar

import av
from livekit import api, rtc

from webm import demux

T = TypeVar("T")
logger = logging.getLogger(__name__)


async def read_next(frames: Iterator[T]) -> T | None:
"""Keep blocking HTTP/demux reads off the loop and join them on cancellation."""
pending = asyncio.create_task(asyncio.to_thread(next, frames, None))
try:
return await asyncio.shield(pending)
except asyncio.CancelledError:
# av.open's read timeout bounds this wait. The container must not be
# closed while a worker is using it.
try:
await pending
finally:
raise


async def publish(room: rtc.Room, frames: Iterator[rtc.EncodedVideoFrame], *, name: str) -> None:
"""Pace and publish one compressed video stream; always release its source."""
frame = await read_next(frames)
if frame is None:
raise ValueError("WebM contains no video frames")
if frame.frame_type != rtc.EncodedFrameType.ENCODED_FRAME_KEY:
raise ValueError("WebM must start at a keyframe")
source = rtc.EncodedVideoSource(frame.width, frame.height)
publication = None
try:
track = rtc.LocalVideoTrack.create_video_track(name, source)
publication = await room.local_participant.publish_track(
track,
rtc.TrackPublishOptions(
video_codec=frame.codec,
video_encoder=rtc.VideoEncoderBackend.ENCODER_BACKEND_PRE_ENCODED,
simulcast=False,
),
)
# Do not consume a short file before the first subscriber is ready.
await asyncio.wait_for(publication.wait_for_subscription(), timeout=30)
origin_pts = frame.timestamp_us
origin_clock = time.monotonic_ns() // 1000
warned_feedback = False
started = False
startup_deadline = time.monotonic() + 5
while frame is not None:
timestamp = origin_clock + frame.timestamp_us - origin_pts
await asyncio.sleep(max(0, (timestamp - time.monotonic_ns() // 1000) / 1_000_000))
if not source.capture_frame(replace(frame, timestamp_us=timestamp)):
raise RuntimeError("native source rejected the encoded frame")
feedback = source.take_feedback()
if not started:
# Track subscription can precede native encoder initialization.
# Keep the initial keyframe until it has processed a frame and
# reported its rate target; otherwise a short file can lose its
# only keyframe before the encoder is ready.
if feedback.rate_control is None:
if time.monotonic() >= startup_deadline:
raise TimeoutError("pre-encoded encoder did not start")
await asyncio.sleep(0.04)
origin_clock = time.monotonic_ns() // 1000
continue
started = True
if not warned_feedback and (feedback.keyframe_requested or feedback.rate_control):
logger.warning(
"Encoder feedback received: %s; this demux-only example cannot "
"change the upstream encoder. Use short GOPs and a suitable bitrate.",
feedback,
)
warned_feedback = True
frame = await read_next(frames)
# There is no encoded-video drain acknowledgement. Keep the track alive
# until cancellation rather than guessing when the final frame arrived.
logger.info("WebM ended; press Ctrl-C to unpublish and disconnect")
await asyncio.Event().wait()
finally:
try:
if publication is not None:
await room.local_participant.unpublish_track(publication.sid)
finally:
await source.aclose()


async def main(args: argparse.Namespace) -> None:
# Open before connecting; URL open/read timeouts prevent an abandoned network
# source from indefinitely holding a worker during shutdown.
with av.open(
args.input,
format="webm",
timeout=(5.0, 5.0),
buffer_size=4096,
options={"probesize": "4096", "analyzeduration": "0"},
) as container:
frames = demux(container)
room = rtc.Room()
token = (
api.AccessToken()
.with_identity("webm-publisher")
.with_grants(api.VideoGrants(room_join=True, room=args.room))
.to_jwt()
)
try:
await room.connect(os.environ["LIVEKIT_URL"], token)
await publish(room, frames, name="webm")
finally:
await room.disconnect()


if __name__ == "__main__":
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("input", help="WebM file or streaming HTTP(S) URL")
parser.add_argument("--room", default="preencoded-webm")
logging.basicConfig(level=logging.INFO)
asyncio.run(main(parser.parse_args()))
3 changes: 3 additions & 0 deletions examples/preencoded-webm/requirements.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
# Requires Python 3.10+.
# Install the local SDK from the repository root first: uv sync --dev
av>=16.1
51 changes: 51 additions & 0 deletions examples/preencoded-webm/webm.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
"""Demux opaque VP8/VP9 WebM for pre-encoded media-track publishing."""

from __future__ import annotations

from typing import Iterator

import av
from livekit import rtc


def demux(container: av.container.InputContainer) -> Iterator[rtc.EncodedVideoFrame]:
"""Read timestamped VP8/VP9 packets, including from non-seekable live WebM.

Media-track publishing does not preserve the WebM container. Reject alpha
rather than silently dropping data that ordinary VP8/VP9 RTP cannot represent.
"""
if not container.streams.video:
raise ValueError("WebM has no video stream")
stream = container.streams.video[0]
codecs = {"vp8": rtc.VideoCodec.VP8, "vp9": rtc.VideoCodec.VP9}
codec = codecs.get(stream.codec_context.name)
if codec is None:
raise ValueError("this example supports VP8/VP9 WebM only")
if stream.metadata.get("alpha_mode") == "1":
raise ValueError("WebM alpha requires a transport that preserves the WebM container")

last_timestamp = None
for packet in container.demux(stream):
if not packet.size:
continue # PyAV's flush sentinel, not a video frame.
if packet.pts is None or packet.time_base is None:
raise ValueError("WebM packet is missing its presentation timestamp")
timestamp = int(packet.pts * packet.time_base * 1_000_000)
if last_timestamp is not None and timestamp <= last_timestamp:
raise ValueError("WebM presentation timestamps must increase")
last_timestamp = timestamp
data = bytes(packet)
if packet.get_sidedata("matroska_block_additional"):
raise ValueError("WebM BlockAdditional cannot be discarded by media-track publishing")
yield rtc.EncodedVideoFrame(
data=data,
width=stream.width,
height=stream.height,
codec=codec,
frame_type=(
rtc.EncodedFrameType.ENCODED_FRAME_KEY
if packet.is_keyframe
else rtc.EncodedFrameType.ENCODED_FRAME_DELTA
),
timestamp_us=timestamp,
)
60 changes: 60 additions & 0 deletions examples/webm-byte-stream/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
# Stream intact WebM through LiveKit

Forward the original bytes of a WebM file or streaming HTTP response through
LiveKit's existing byte-stream API, with MIME type `video/webm`. WebM headers,
timestamps, compressed color, and alpha `BlockAdditional` data remain unchanged.
There is no demuxing, decoding, re-encoding, remuxing, or custom media format.
The publisher does not parse or validate the input; supply valid streaming WebM.

From the repository root (Python 3.10+, with the bundled FFI installed):

```sh
uv sync --dev
export LIVEKIT_URL=ws://localhost:7880
export LIVEKIT_API_KEY=devkey
export LIVEKIT_API_SECRET=secret
uv run python examples/webm-byte-stream/stream_webm.py input.webm \
--room demo --recipient viewer
# input.webm can also be a streaming HTTP(S) URL.
```

The receiver must already be connected and have a `webm` byte-stream handler
registered before sending. The example checks participant presence, which does
not prove that its handler is ready. Applications need their own readiness flow.
Streams do not replay their beginning to participants who join midway; restart
from a suitable WebM initialization segment and keyframe for a new viewer.

Reads and awaited writes are incremental, using chunks of at most 15 KB. Chunk
boundaries need not align with WebM elements: concatenating the received chunks
reconstructs the original byte stream. HTTP reads have a five-second inactivity
timeout; a failed or cancelled transfer closes the writer with an error reason.
A file is sent as quickly as the transport allows; this is not a frame scheduler.
The receiver/player uses WebM timestamps for playout.

## Why a byte stream

LiveKit's pre-encoded media source takes codec access units. The standard
[VP9 RTP payload format](https://www.rfc-editor.org/rfc/rfc9628.html#section-4)
describes VP9 frame transport, not WebM container or alpha `BlockAdditional`
transport. Sending WebM as if it were a VP9 access unit would require a custom
receiver and transport convention. Changing the MIME type alone cannot add that
support to an ordinary WebRTC video track.

LiveKit's [byte streams](https://docs.livekit.io/transport/data/byte-streams/)
already support incremental arbitrary bytes. This example uses that existing
transport rather than inventing an RTP payload. It is a data stream, not a
published video media track: ordinary track attachment, media congestion feedback,
and automatic video subscription do not apply. Byte streams use reliable delivery;
loss recovery can delay later bytes and increase latency. LiveKit's lossy data
tracks are not interchangeable: dropping arbitrary WebM chunks corrupts the stream.

## Future playback

A client handler must feed incoming bytes into a streaming media player. On the
web, the intended path is Media Source Extensions with a compatible WebM byte
stream; see the [WebM byte-stream specification](https://www.w3.org/TR/mse-byte-stream-format-webm/).
This preserves the possibility of native alpha composition without an
application-side decode/re-encode step. An arbitrary WebM file is not necessarily
MSE-compatible, and file playback support does not prove streaming alpha support.
Browser/mobile playback, alpha support, buffering, and latency need platform tests.
No client playback adapter is implemented in this example.
Loading
Loading