From e117706b3a044ebd9f9be715e7dd25d8cc7f1d31 Mon Sep 17 00:00:00 2001 From: AmitMY Date: Wed, 30 Sep 2026 11:43:01 +0200 Subject: [PATCH 1/5] feat(rtc): expose pre-encoded video publishing in Python --- livekit-rtc/README.md | 45 +++++ livekit-rtc/livekit/rtc/__init__.py | 21 ++- livekit-rtc/livekit/rtc/encoded_video.py | 127 ++++++++++++++ livekit-rtc/livekit/rtc/track.py | 5 +- tests/rtc/test_encoded_video.py | 204 +++++++++++++++++++++++ 5 files changed, 400 insertions(+), 2 deletions(-) create mode 100644 livekit-rtc/livekit/rtc/encoded_video.py create mode 100644 tests/rtc/test_encoded_video.py diff --git a/livekit-rtc/README.md b/livekit-rtc/README.md index 516f9b6c..91494d75 100644 --- a/livekit-rtc/README.md +++ b/livekit-rtc/README.md @@ -4,3 +4,48 @@ Python SDK to integrate LiveKit's real-time video, audio, and data capabilities See https://docs.livekit.io/ for more information. +## Publishing pre-encoded video + +Use `EncodedVideoSource` when an upstream encoder already provides compressed +frames. The source uses the same native passthrough encoder as the Rust SDK; +Python does not decode or re-encode the video. + +```python +from livekit import rtc + +source = rtc.EncodedVideoSource(width=640, height=360) +track = rtc.LocalVideoTrack.create_video_track("camera", source) +await room.local_participant.publish_track( + track, + rtc.TrackPublishOptions( + video_codec=rtc.VideoCodec.VP9, + video_encoder=rtc.VideoEncoderBackend.ENCODER_BACKEND_PRE_ENCODED, + simulcast=False, + ), +) +# Supply each complete encoded frame at its playback time. +source.capture_frame(rtc.EncodedVideoFrame( + data=encoded_vp9_frame, + width=640, + height=360, + codec=rtc.VideoCodec.VP9, + frame_type=rtc.EncodedFrameType.ENCODED_FRAME_KEY, + timestamp_us=capture_timestamp_us, +)) +feedback = source.take_feedback() +# Forward feedback.keyframe_requested and feedback.rate_control to your encoder. +# At shutdown, unpublish the track and release the source: +await room.local_participant.unpublish_track(track.sid) +await source.aclose() +``` + +Supported codecs are H.264, H.265, VP8, VP9 and AV1, subject to receiver support. +Submit complete access units (Annex B for H.264/H.265), not WebM/MP4 container +bytes or RTP packets. A container must first be demuxed; demuxing extracts +compressed frames and does not reconstruct pixels. The codec must match the +publication and stay fixed, and simulcast must be disabled. + +The caller controls pacing and handles keyframe and bitrate feedback. In +particular, a prerecorded file cannot generate a new keyframe on demand for a +late subscriber or after packet loss. WebM's auxiliary alpha data is not part of +the VP9 color access unit and is not transported by this API. diff --git a/livekit-rtc/livekit/rtc/__init__.py b/livekit-rtc/livekit/rtc/__init__.py index d5b9bcd6..4be4c124 100644 --- a/livekit-rtc/livekit/rtc/__init__.py +++ b/livekit-rtc/livekit/rtc/__init__.py @@ -35,6 +35,7 @@ SimulateScenarioKind, TrackPublishOptions, VideoEncoding, + VideoEncoderBackend, ) from ._proto.track_pb2 import ( FrameMetadataFeature, @@ -43,7 +44,13 @@ TrackSource, ParticipantTrackPermission, ) -from ._proto.video_frame_pb2 import FrameMetadata, VideoBufferType, VideoCodec, VideoRotation +from ._proto.video_frame_pb2 import ( + EncodedFrameType, + FrameMetadata, + VideoBufferType, + VideoCodec, + VideoRotation, +) from ._proto.track_publication_pb2 import VideoQuality from .audio_frame import AudioFrame from .audio_source import AudioSource @@ -95,6 +102,12 @@ from .video_frame import ( VideoFrame, ) +from .encoded_video import ( + EncodedVideoFrame, + EncodedVideoSource, + EncodedVideoSourceFeedback, + EncodedRateControl, +) from .video_source import VideoSource from .video_stream import VideoFrameEvent, VideoStream from .audio_resampler import AudioResampler, AudioResamplerQuality @@ -138,6 +151,12 @@ from .frame_processor import FrameProcessor __all__ = [ + "EncodedFrameType", + "EncodedVideoFrame", + "EncodedVideoSource", + "EncodedVideoSourceFeedback", + "EncodedRateControl", + "VideoEncoderBackend", "ConnectionQuality", "ConnectionState", "DataPacketKind", diff --git a/livekit-rtc/livekit/rtc/encoded_video.py b/livekit-rtc/livekit/rtc/encoded_video.py new file mode 100644 index 00000000..85a2268b --- /dev/null +++ b/livekit-rtc/livekit/rtc/encoded_video.py @@ -0,0 +1,127 @@ +# Copyright 2026 LiveKit, Inc. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +from __future__ import annotations + +from dataclasses import dataclass + +from ._ffi_client import FfiClient, FfiHandle +from ._proto import ffi_pb2 as proto_ffi +from ._proto import video_frame_pb2 as proto_video +from ._utils import get_address + + +@dataclass(frozen=True) +class EncodedVideoFrame: + """One complete encoded access unit, without a container or RTP headers. + + ``data`` must contain a single frame in ``codec`` (Annex B for H.264/H.265). + ``timestamp_us`` is its capture time in microseconds, increasing along the + stream. The codec must match the publication and remain fixed for its lifetime. + """ + + data: bytes + width: int + height: int + codec: proto_video.VideoCodec.ValueType + frame_type: proto_video.EncodedFrameType.ValueType + timestamp_us: int + metadata: proto_video.FrameMetadata | None = None + + def __post_init__(self) -> None: + if not isinstance(self.data, bytes): + raise TypeError("encoded frame data must be bytes") + if not self.data: + raise ValueError("encoded frame data must not be empty") + if self.width <= 0 or self.height <= 0: + raise ValueError("encoded frame dimensions must be positive") + + +@dataclass(frozen=True) +class EncodedRateControl: + """The latest bitrate and frame-rate targets requested by WebRTC.""" + + target_bitrate_bps: int + framerate_fps: float + + +@dataclass(frozen=True) +class EncodedVideoSourceFeedback: + """Encoder feedback consumed since the previous call to ``take_feedback``.""" + + keyframe_requested: bool + rate_control: EncodedRateControl | None + + +class EncodedVideoSource: + """Publish pre-encoded video without decoding or re-encoding it. + + Create a ``LocalVideoTrack`` with this source and publish it with a matching + ``video_codec``, ``video_encoder=ENCODER_BACKEND_PRE_ENCODED`` and + ``simulcast=False``. Feed complete access units at the intended playback rate; + this source does not pace or buffer a file for playback. + + Poll ``take_feedback`` and forward keyframe and rate-control requests to the + upstream encoder. A demuxer alone cannot produce a new keyframe or adapt the + encoded bitrate. Call ``capture_frame`` from a single producer. + """ + + def __init__(self, width: int, height: int) -> None: + """Create a source with the initial encoded frame dimensions.""" + if width <= 0 or height <= 0: + raise ValueError("encoded source dimensions must be positive") + req = proto_ffi.FfiRequest() + req.new_video_source.type = proto_video.VideoSourceType.VIDEO_SOURCE_ENCODED + req.new_video_source.resolution.width = width + req.new_video_source.resolution.height = height + resp = FfiClient.instance.request(req) + self._ffi_handle = FfiHandle(resp.new_video_source.source.handle.id) + + def capture_frame(self, frame: EncodedVideoFrame) -> bool: + """Submit a frame; return whether the native source accepted it. + + The native implementation copies the payload before this call returns. + Acceptance does not guarantee delivery to subscribers. + """ + req = proto_ffi.FfiRequest() + capture = req.capture_encoded_video_frame + capture.source_handle = self._ffi_handle.handle + capture.buffer.data_ptr = get_address(frame.data) + capture.buffer.data_len = len(frame.data) + capture.codec = frame.codec + capture.frame_type = frame.frame_type + capture.width = frame.width + capture.height = frame.height + capture.timestamp_us = frame.timestamp_us + if frame.metadata is not None: + capture.metadata.CopyFrom(frame.metadata) + response: proto_ffi.FfiResponse = FfiClient.instance.request(req) + return response.capture_encoded_video_frame.accepted + + def take_feedback(self) -> EncodedVideoSourceFeedback: + """Consume pending feedback, including requests from late subscribers.""" + req = proto_ffi.FfiRequest() + req.take_encoded_video_source_feedback.source_handle = self._ffi_handle.handle + feedback = FfiClient.instance.request(req).take_encoded_video_source_feedback + rate_control = None + if feedback.HasField("rate_control"): + rate_control = EncodedRateControl( + target_bitrate_bps=feedback.rate_control.target_bitrate_bps, + framerate_fps=feedback.rate_control.framerate_fps, + ) + return EncodedVideoSourceFeedback(feedback.keyframe_requested, rate_control) + + async def aclose(self) -> None: + """Release the native source handle.""" + self._ffi_handle.dispose() diff --git a/livekit-rtc/livekit/rtc/track.py b/livekit-rtc/livekit/rtc/track.py index 684e2493..c8a2fc1d 100644 --- a/livekit-rtc/livekit/rtc/track.py +++ b/livekit-rtc/livekit/rtc/track.py @@ -26,6 +26,7 @@ from .audio_stream import AudioStream from .room import Room from .video_source import VideoSource + from .encoded_video import EncodedVideoSource from .platform_audio import PlatformAudioSource @@ -213,7 +214,9 @@ def __init__(self, info: proto_track.OwnedTrack): super().__init__(info) @staticmethod - def create_video_track(name: str, source: "VideoSource") -> "LocalVideoTrack": + def create_video_track( + name: str, source: "Union[VideoSource, EncodedVideoSource]" + ) -> "LocalVideoTrack": req = proto_ffi.FfiRequest() req.create_video_track.name = name req.create_video_track.source_handle = source._ffi_handle.handle diff --git a/tests/rtc/test_encoded_video.py b/tests/rtc/test_encoded_video.py new file mode 100644 index 00000000..0e1c3963 --- /dev/null +++ b/tests/rtc/test_encoded_video.py @@ -0,0 +1,204 @@ +from __future__ import annotations + +import ctypes +from dataclasses import replace +from typing import Any + +import pytest + +from livekit import rtc +from livekit.rtc import encoded_video +from livekit.rtc._proto import ffi_pb2 as proto_ffi + + +# A 16x16 red VP9 keyframe generated with libvpx-vp9, yuv420p, realtime. +VP9_KEYFRAME = bytes.fromhex( + "824983420000f000f60038241c184a00003060000010bffff4a3dffffff9c97fffffffbd840000" +) + + +def frame(**kwargs: Any) -> rtc.EncodedVideoFrame: + return replace( + rtc.EncodedVideoFrame( + data=VP9_KEYFRAME, + width=16, + height=16, + codec=rtc.VideoCodec.VP9, + frame_type=rtc.EncodedFrameType.ENCODED_FRAME_KEY, + timestamp_us=123456, + ), + **kwargs, + ) + + +@pytest.mark.parametrize("width,height", [(0, 16), (16, 0), (-1, 16)]) +def test_invalid_dimensions(width: int, height: int) -> None: + with pytest.raises(ValueError, match="dimensions"): + frame(width=width, height=height) + with pytest.raises(ValueError, match="dimensions"): + rtc.EncodedVideoSource(width, height) + + +def test_empty_or_mutable_payload() -> None: + with pytest.raises(ValueError, match="empty"): + frame(data=b"") + with pytest.raises(TypeError, match="bytes"): + frame(data=bytearray(VP9_KEYFRAME)) + + +@pytest.mark.parametrize("accepted", [True, False]) +@pytest.mark.parametrize("metadata", [None, rtc.FrameMetadata(frame_id=7, user_timestamp=99)]) +async def test_capture_ffi_contract( + monkeypatch: pytest.MonkeyPatch, accepted: bool, metadata: rtc.FrameMetadata | None +) -> None: + requests = [] + + def request(self: Any, req: proto_ffi.FfiRequest) -> proto_ffi.FfiResponse: + requests.append(req) + response = proto_ffi.FfiResponse() + if req.HasField("new_video_source"): + assert req.new_video_source.type == encoded_video.proto_video.VIDEO_SOURCE_ENCODED + assert req.new_video_source.resolution.width == 16 + assert req.new_video_source.resolution.height == 16 + response.new_video_source.source.handle.id = 0 + else: + capture = req.capture_encoded_video_frame + # Read foreign memory during the synchronous call, while it must be valid. + assert ( + ctypes.string_at(capture.buffer.data_ptr, capture.buffer.data_len) == VP9_KEYFRAME + ) + assert capture.codec == rtc.VideoCodec.VP9 + assert capture.frame_type == rtc.EncodedFrameType.ENCODED_FRAME_KEY + assert (capture.width, capture.height, capture.timestamp_us) == (16, 16, 123456) + assert capture.HasField("metadata") == (metadata is not None) + if metadata is not None: + assert capture.metadata == metadata + response.capture_encoded_video_frame.accepted = accepted + return response + + monkeypatch.setattr(encoded_video.FfiClient, "request", request) + source = rtc.EncodedVideoSource(16, 16) + try: + assert source.capture_frame(frame(metadata=metadata)) is accepted + assert len(requests) == 2 + finally: + await source.aclose() + + +@pytest.mark.parametrize("has_rate_control", [False, True]) +async def test_feedback(monkeypatch: pytest.MonkeyPatch, has_rate_control: bool) -> None: + source = rtc.EncodedVideoSource(16, 16) + calls = 0 + + def request(self: Any, req: proto_ffi.FfiRequest) -> proto_ffi.FfiResponse: + nonlocal calls + assert req.take_encoded_video_source_feedback.source_handle == source._ffi_handle.handle + result = proto_ffi.FfiResponse() + feedback = result.take_encoded_video_source_feedback + feedback.keyframe_requested = calls == 0 + if has_rate_control and calls == 0: + feedback.rate_control.target_bitrate_bps = 123000 + feedback.rate_control.framerate_fps = 29.97 + calls += 1 + return result + + monkeypatch.setattr(encoded_video.FfiClient, "request", request) + try: + feedback = source.take_feedback() + assert feedback.keyframe_requested + assert feedback.rate_control == ( + rtc.EncodedRateControl(123000, 29.97) if has_rate_control else None + ) + assert source.take_feedback() == rtc.EncodedVideoSourceFeedback(False, None) + finally: + await source.aclose() + + +async def test_native_source_lifecycle() -> None: + # Exercise the shipped native FFI, not just serialization mocks. + for _ in range(10): + source = rtc.EncodedVideoSource(16, 16) + try: + track = rtc.LocalVideoTrack.create_video_track("encoded", source) + assert track.kind == rtc.TrackKind.KIND_VIDEO + assert isinstance(source.capture_frame(frame()), bool) + assert source.take_feedback() == rtc.EncodedVideoSourceFeedback(False, None) + finally: + await source.aclose() + await source.aclose() # Closing is idempotent. + assert source._ffi_handle.disposed + del track + + +async def test_publish_preencoded_video() -> None: + """A receiver decodes the access unit published through the real native FFI.""" + import asyncio + import contextlib + import os + import time + import uuid + + from livekit import api + + if not all(os.getenv(key) for key in ("LIVEKIT_URL", "LIVEKIT_API_KEY", "LIVEKIT_API_SECRET")): + pytest.skip("LiveKit server credentials are required") + + room_name = f"encoded-video-{uuid.uuid4().hex}" + publisher, subscriber = rtc.Room(), rtc.Room() + source = rtc.EncodedVideoSource(16, 16) + subscribed: asyncio.Future[rtc.Track] = asyncio.get_running_loop().create_future() + video_stream = None + producer = None + + @subscriber.on("track_subscribed") + def on_track( + track: rtc.Track, + publication: rtc.RemoteTrackPublication, + participant: rtc.RemoteParticipant, + ) -> None: + if not subscribed.done(): + subscribed.set_result(track) + + def token(identity: str) -> str: + return ( + api.AccessToken() + .with_identity(identity) + .with_grants(api.VideoGrants(room_join=True, room=room_name)) + .to_jwt() + ) + + try: + await subscriber.connect(os.environ["LIVEKIT_URL"], token("subscriber")) + await publisher.connect(os.environ["LIVEKIT_URL"], token("publisher")) + track = rtc.LocalVideoTrack.create_video_track("preencoded", source) + await publisher.local_participant.publish_track( + track, + rtc.TrackPublishOptions( + video_codec=rtc.VideoCodec.VP9, + video_encoder=rtc.VideoEncoderBackend.ENCODER_BACKEND_PRE_ENCODED, + simulcast=False, + ), + ) + remote = await asyncio.wait_for(subscribed, 10) + video_stream = rtc.VideoStream(remote, format=rtc.VideoBufferType.RGBA) + + async def send() -> None: + while True: + source.capture_frame(frame(timestamp_us=time.monotonic_ns() // 1000)) + await asyncio.sleep(0.04) + + producer = asyncio.create_task(send()) + received = await asyncio.wait_for(video_stream.__anext__(), 10) + assert (received.frame.width, received.frame.height) == (16, 16) + red, green, blue, alpha = received.frame.data[:4] + assert red > 200 and green < 50 and blue < 50 and alpha == 255 + finally: + if producer is not None: + producer.cancel() + with contextlib.suppress(asyncio.CancelledError): + await producer + if video_stream is not None: + await video_stream.aclose() + await publisher.disconnect() + await subscriber.disconnect() + await source.aclose() From 0c31793b9680e1cf0e9193f090830888d7c479e8 Mon Sep 17 00:00:00 2001 From: AmitMY Date: Wed, 30 Sep 2026 11:47:53 +0200 Subject: [PATCH 2/5] fix(ci): validate generated stubs on fork pull requests --- .github/workflows/build-rtc.yml | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/.github/workflows/build-rtc.yml b/.github/workflows/build-rtc.yml index 14841ddd..1310faab 100644 --- a/.github/workflows/build-rtc.yml +++ b/.github/workflows/build-rtc.yml @@ -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 @@ -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/"]' From 3b73b73db66b936584f3d566dcf875ce31f95876 Mon Sep 17 00:00:00 2001 From: AmitMY Date: Wed, 30 Sep 2026 11:55:32 +0200 Subject: [PATCH 3/5] feat(examples): stream demuxed WebM and prototype single-track alpha --- .github/workflows/tests.yml | 6 +- examples/preencoded-webm/README.md | 119 ++++++++ examples/preencoded-webm/publish.py | 139 ++++++++++ examples/preencoded-webm/requirements.txt | 3 + examples/preencoded-webm/webm.py | 103 +++++++ tests/rtc/test_webm_example.py | 313 ++++++++++++++++++++++ 6 files changed, 680 insertions(+), 3 deletions(-) create mode 100644 examples/preencoded-webm/README.md create mode 100644 examples/preencoded-webm/publish.py create mode 100644 examples/preencoded-webm/requirements.txt create mode 100644 examples/preencoded-webm/webm.py create mode 100644 tests/rtc/test_webm_example.py diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 5a3aaca8..749db67c 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -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' @@ -109,7 +109,7 @@ 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' @@ -117,7 +117,7 @@ jobs: 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) diff --git a/examples/preencoded-webm/README.md b/examples/preencoded-webm/README.md new file mode 100644 index 00000000..a4f338cb --- /dev/null +++ b/examples/preencoded-webm/README.md @@ -0,0 +1,119 @@ +# 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. + +## Experimental single-track alpha transport + +**This is a backend wire-format prototype, not interoperable transparent WebRTC.** +Opaque publishing works with existing VP8/VP9 receivers. Alpha mode must only be +used with matching test receivers in an isolated room; an ordinary browser/mobile +LiveKit video track does not recognize the additional data. No browser or mobile +playback adapter is implemented here. + +By default, the demuxer rejects alpha WebM instead of silently dropping its alpha. +To opt into the experimental VP9 envelope: + +```sh +.venv/bin/python examples/preencoded-webm/publish.py transparent.webm \ + --room alpha-experiment --experimental-alpha +``` + +The prototype extracts the VP9 color access unit and the separately compressed +alpha access unit from WebM `BlockAdditional` with `BlockAddID=1`. It publishes both +inside **one encoded-frame payload on one media track**, sharing one frame type, +timestamp, and RTP frame boundary. There is no second track or data channel. +The native passthrough encoder forwards the bytes, and normal RTP packetization +can split them across multiple packets. + +The provisional payload is: + +```text +[color access unit][alpha access unit] +[color length: u32 big endian][alpha length: u32 big endian]["LKWA"][version: u8 = 1] +``` + +`LKWA` and version 1 are local proposal values, not an assigned LiveKit format. +Lengths exclude the 13-byte footer. Both components must be nonempty and the +entire payload must be at most 16 MiB in this example. A receiver validates the +footer and sizes after RTP reassembly, then extracts the original compressed +components **before** handing anything to an ordinary VP9 decoder. The included +`unpack_alpha` function specifies that byte contract; it is not a video renderer. +Sparse/reused alpha blocks and dimension changes are not supported. A WebM +keyframe must contain independently decodable color and alpha components. + +This deliberately does not overload LiveKit's `FrameMetadata.user_data`: the +native LKTS trailer has a one-byte length and a maximum total size of 255 bytes, +while alpha access units can be many kilobytes. The media envelope precedes any +native metadata trailer. Combining it with frame metadata or E2EE needs separate +interoperability testing; this example does not enable either. + +### Required before this could become a supported feature + +- Agree on a versioned transport profile with LiveKit maintainers. The example's + track name identifies a prototype; it does **not** negotiate receiver support. +- Add capability signaling and enforce it for every subscriber, including late + joins. Unsupported receivers must be rejected or explicitly offered an opaque + fallback. Today's server does not strip this custom alpha envelope for them. +- Implement a receiver path for each target platform. A browser adapter could + investigate remuxing compressed color and alpha into WebM for native playback; + browser/MSE alpha support and latency must be tested. Native mobile SDKs need + their own supported decode/playback path. No cross-platform guarantee follows + from successful byte transport. +- Validate loss/reordering/retransmission, congestion feedback, keyframe recovery, + codec profiles, resource limits, metadata, encryption, and reconnection. + +These are interoperability requirements, not additional Python encoder settings. +Keep the alpha mode experimental until they are resolved. + +## Validation + +```sh +uv run --with 'av>=16.1' pytest tests/rtc/test_encoded_video.py tests/rtc/test_webm_example.py +``` + +Tests generate synthetic alpha WebM, verify unchanged compressed payloads and +key/delta timestamps, read the first frame before a pipe reaches EOF, reject +malformed/oversized envelopes, and check cancellation joins the active demux read. +The base API's server-backed test publishes VP9 and verifies subscriber pixels. + +A separate local transport probe used Python FFI 0.12.80, LiveKit server 1.13.7, +and a Go SDK 2.18.1 subscriber to reassemble RTP without decoding. A 37,630-byte +compound frame (27,578 color bytes, 10,039 alpha bytes, 13-byte footer) crossed the +server in 33 RTP packets with identical SHA-256 before and after. This demonstrates +transport through the tested stack; it does not validate an alpha player or +compatibility with every LiveKit deployment. diff --git a/examples/preencoded-webm/publish.py b/examples/preencoded-webm/publish.py new file mode 100644 index 00000000..8f1963e2 --- /dev/null +++ b/examples/preencoded-webm/publish.py @@ -0,0 +1,139 @@ +"""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, experimental_alpha=args.experimental_alpha) + 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="experimental-webm-alpha-v1" if args.experimental_alpha else "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") + parser.add_argument( + "--experimental-alpha", + action="store_true", + help="use the unnegotiated alpha prototype; matching test receivers only", + ) + logging.basicConfig(level=logging.INFO) + asyncio.run(main(parser.parse_args())) diff --git a/examples/preencoded-webm/requirements.txt b/examples/preencoded-webm/requirements.txt new file mode 100644 index 00000000..b13ee1d4 --- /dev/null +++ b/examples/preencoded-webm/requirements.txt @@ -0,0 +1,3 @@ +# Requires Python 3.10+. +# Install the local SDK from the repository root first: uv sync --dev +av>=16.1 diff --git a/examples/preencoded-webm/webm.py b/examples/preencoded-webm/webm.py new file mode 100644 index 00000000..a82ced71 --- /dev/null +++ b/examples/preencoded-webm/webm.py @@ -0,0 +1,103 @@ +"""Demux WebM without decoding; prototype a single-track VP9 alpha envelope. + +The alpha envelope is experimental, not a negotiated LiveKit or WebRTC format. +It requires a matching receiver to strip it before invoking a video decoder. +""" + +from __future__ import annotations + +import struct +from typing import Iterator + +import av +from livekit import rtc + +# Provisional wire format for isolated interoperability experiments only: +# [color access unit][alpha access unit][color_len:u32][alpha_len:u32][LKWA][v1] +# Both lengths are big endian. The footer is inside the encoded frame, before +# any native LiveKit metadata trailer. Do not reuse the 255-byte LKTS metadata +# envelope: alpha access units can span many RTP packets. +_FOOTER = struct.Struct(">II4sB") +_MAX_FRAME_BYTES = 16 * 1024 * 1024 + + +def pack_alpha(color: bytes, alpha: bytes) -> bytes: + """Package two compressed VP9 access units without inspecting their pixels.""" + if not color or not alpha: + raise ValueError("color and alpha access units must not be empty") + if len(color) + len(alpha) + _FOOTER.size > _MAX_FRAME_BYTES: + raise ValueError("alpha envelope exceeds the 16 MiB example limit") + return color + alpha + _FOOTER.pack(len(color), len(alpha), b"LKWA", 1) + + +def unpack_alpha(data: bytes) -> tuple[bytes, bytes]: + """Validate the experimental envelope and return its compressed components. + + Included to specify and test the wire contract, not as a playback adapter. + """ + if not _FOOTER.size < len(data) <= _MAX_FRAME_BYTES: + raise ValueError("invalid alpha envelope size") + color_size, alpha_size, magic, version = _FOOTER.unpack_from(data, len(data) - _FOOTER.size) + if magic != b"LKWA" or version != 1: + raise ValueError("unknown alpha envelope") + if not color_size or not alpha_size or color_size + alpha_size != len(data) - _FOOTER.size: + raise ValueError("invalid alpha envelope lengths") + return data[:color_size], data[color_size : -_FOOTER.size] + + +def demux( + container: av.container.InputContainer, *, experimental_alpha: bool = False +) -> Iterator[rtc.EncodedVideoFrame]: + """Read timestamped VP8/VP9 packets, including from non-seekable live WebM. + + Alpha is rejected by default rather than silently discarded. The opt-in + prototype requires VP9 and an alpha access unit in every packet; sparse alpha + and changing dimensions are outside this example's scope. + """ + 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") + has_alpha = stream.metadata.get("alpha_mode") == "1" + if has_alpha and not experimental_alpha: + raise ValueError( + "WebM has alpha; use the experimental alpha receiver/profile or an opaque source" + ) + if experimental_alpha and (not has_alpha or codec != rtc.VideoCodec.VP9): + raise ValueError("experimental alpha requires VP9 WebM with AlphaMode=1") + + 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) + additional = bytes(packet.get_sidedata("matroska_block_additional")) + if additional: + if len(additional) <= 8 or int.from_bytes(additional[:8], "big") != 1: + raise ValueError("unsupported WebM BlockAdditional mapping") + if not experimental_alpha: + raise ValueError("WebM alpha cannot be discarded implicitly") + data = pack_alpha(data, additional[8:]) + elif experimental_alpha: + raise ValueError("missing alpha access unit in WebM packet") + 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, + ) diff --git a/tests/rtc/test_webm_example.py b/tests/rtc/test_webm_example.py new file mode 100644 index 00000000..5482c227 --- /dev/null +++ b/tests/rtc/test_webm_example.py @@ -0,0 +1,313 @@ +"""Run with `uv run --with 'av>=16.1' pytest tests/rtc/test_webm_example.py`.""" + +from __future__ import annotations + +import asyncio +import importlib.util +import io +import os +import struct +import threading +from fractions import Fraction +from pathlib import Path + +import numpy as np +import pytest + +av = pytest.importorskip( + "av", minversion="16.1", reason="optional WebM example dependency (Python 3.10+)" +) +from livekit import rtc # noqa: E402 + +_EXAMPLE = Path(__file__).resolve().parents[2] / "examples" / "preencoded-webm" +_spec = importlib.util.spec_from_file_location("webm", _EXAMPLE / "webm.py") +webm = importlib.util.module_from_spec(_spec) +_spec.loader.exec_module(webm) + + +@pytest.fixture(scope="module") +def alpha_webm() -> bytes: + """Synthetic media only; alpha is intentionally larger than one RTP packet.""" + buf = io.BytesIO() + with av.open(buf, "w", format="webm", options={"live": "1", "cluster_time_limit": "1"}) as c: + stream = c.add_stream("libvpx-vp9", rate=25) + stream.width = stream.height = 160 + stream.pix_fmt = "yuva420p" + stream.options = { + "deadline": "realtime", + "cpu-used": "8", + "lag-in-frames": "0", + "auto-alt-ref": "0", + } + pixels = np.random.default_rng(4).integers(0, 256, (160, 160, 4), dtype=np.uint8) + for i in range(12): + frame = av.VideoFrame.from_ndarray(pixels, format="rgba") + frame.pts = i + frame.time_base = Fraction(1, 25) + c.mux(stream.encode(frame)) + c.mux(stream.encode(None)) + return buf.getvalue() + + +def test_demux_preserves_compressed_color_and_alpha(alpha_webm: bytes) -> None: + with av.open(io.BytesIO(alpha_webm)) as c: + packets = [p for p in c.demux(video=0) if p.size] + with av.open(io.BytesIO(alpha_webm)) as c: + # No decode call: payloads must be byte-for-byte identical to the demuxer. + frames = list(webm.demux(c, experimental_alpha=True)) + assert len(frames) == len(packets) == 12 + assert [f.timestamp_us for f in frames] == list(range(0, 480000, 40000)) + for frame, packet in zip(frames, packets): + color, alpha = webm.unpack_alpha(frame.data) + assert color == bytes(packet) + assert alpha == bytes(packet.get_sidedata("matroska_block_additional"))[8:] + assert frame.frame_type == ( + rtc.EncodedFrameType.ENCODED_FRAME_KEY + if packet.is_keyframe + else rtc.EncodedFrameType.ENCODED_FRAME_DELTA + ) + assert len(webm.unpack_alpha(frames[0].data)[1]) > 1200 + + +def test_alpha_requires_explicit_opt_in(alpha_webm: bytes) -> None: + with av.open(io.BytesIO(alpha_webm)) as c: + with pytest.raises(ValueError, match="alpha"): + next(webm.demux(c)) + + +@pytest.mark.parametrize("codec", ["libvpx", "libvpx-vp9"]) +def test_opaque_webm(codec: str) -> None: + buf = io.BytesIO() + with av.open(buf, "w", format="webm") as c: + stream = c.add_stream(codec, rate=25) + stream.width = stream.height = 16 + stream.pix_fmt = "yuv420p" + stream.options = {"deadline": "realtime", "lag-in-frames": "0"} + c.mux( + stream.encode( + av.VideoFrame.from_ndarray(np.zeros((16, 16, 3), dtype=np.uint8), format="rgb24") + ) + ) + c.mux(stream.encode(None)) + with av.open(io.BytesIO(buf.getvalue())) as c: + packets = [bytes(p) for p in c.demux(video=0) if p.size] + with av.open(io.BytesIO(buf.getvalue())) as c: + assert [f.data for f in webm.demux(c)] == packets + with av.open(io.BytesIO(buf.getvalue())) as c: + with pytest.raises(ValueError, match="AlphaMode"): + next(webm.demux(c, experimental_alpha=True)) + + +def test_demux_before_pipe_eof(alpha_webm: bytes) -> None: + read_fd, write_fd = os.pipe() + decoded_first = threading.Event() + writer_closed = threading.Event() + + def write() -> None: + try: + with os.fdopen(write_fd, "wb") as pipe: + for start in range(0, len(alpha_webm), 4096): + pipe.write(alpha_webm[start : start + 4096]) + pipe.flush() + decoded_first.wait(5) + finally: + writer_closed.set() + + writer = threading.Thread(target=write) + writer.start() + try: + with ( + os.fdopen(read_fd, "rb") as pipe, + av.open( + pipe, + format="webm", + buffer_size=4096, + options={"probesize": "4096", "analyzeduration": "0"}, + ) as c, + ): + frames = webm.demux(c, experimental_alpha=True) + first = next(frames) + assert not writer_closed.is_set() + decoded_first.set() + assert first.frame_type == rtc.EncodedFrameType.ENCODED_FRAME_KEY + assert len(list(frames)) == 11 + finally: + decoded_first.set() + writer.join(timeout=5) + assert not writer.is_alive() + + +@pytest.mark.parametrize( + "data", + [ + b"", + b"short", + b"caa" + struct.pack(">II4sB", 1, 2, b"LKWA", 2), + b"caa" + struct.pack(">II4sB", 1, 3, b"LKWA", 1), + b"caa" + struct.pack(">II4sB", 0, 3, b"LKWA", 1), + b"caa" + struct.pack(">II4sB", 1, 2, b"NOPE", 1), + ], +) +def test_reject_bad_envelopes(data: bytes) -> None: + with pytest.raises(ValueError): + webm.unpack_alpha(data) + + +def test_bound_envelope_size(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr(webm, "_MAX_FRAME_BYTES", 32) + with pytest.raises(ValueError, match="limit"): + webm.pack_alpha(b"c" * 16, b"a" * 16) + with pytest.raises(ValueError): + webm.pack_alpha(b"", b"a") + with pytest.raises(ValueError): + webm.unpack_alpha(b"x" * 33) + + +async def test_cancel_waits_for_reader(monkeypatch: pytest.MonkeyPatch) -> None: + import sys + + monkeypatch.setitem(sys.modules, "webm", webm) + spec = importlib.util.spec_from_file_location("webm_publisher", _EXAMPLE / "publish.py") + publisher = importlib.util.module_from_spec(spec) + spec.loader.exec_module(publisher) + started, release, finished = threading.Event(), threading.Event(), threading.Event() + + def frames(): + started.set() + assert release.wait(5) + finished.set() + yield 1 + + task = asyncio.create_task(publisher.read_next(iter(frames()))) + try: + assert await asyncio.to_thread(started.wait, 5) + task.cancel() + await asyncio.sleep(0) + assert not task.done() + release.set() + with pytest.raises(asyncio.CancelledError): + await task + assert finished.is_set() + finally: + release.set() + + +def test_http_stream_before_response_finishes(alpha_webm: bytes) -> None: + from http.server import BaseHTTPRequestHandler, HTTPServer + + first_frame, response_finished = threading.Event(), threading.Event() + + class Handler(BaseHTTPRequestHandler): + def do_GET(self) -> None: + self.send_response(200) + self.send_header("Content-Type", "video/webm") + self.end_headers() + self.wfile.write(alpha_webm) + self.wfile.flush() + first_frame.wait(5) + response_finished.set() + + def log_message(self, *args) -> None: + pass + + server = HTTPServer(("127.0.0.1", 0), Handler) + worker = threading.Thread(target=server.handle_request) + worker.start() + try: + with av.open( + f"http://127.0.0.1:{server.server_port}/live.webm", + format="webm", + timeout=(5.0, 5.0), + options={"probesize": "4096", "analyzeduration": "0"}, + ) as c: + frame = next(webm.demux(c, experimental_alpha=True)) + assert not response_finished.is_set() + assert webm.unpack_alpha(frame.data)[1] + finally: + first_frame.set() + worker.join(timeout=5) + server.server_close() + assert not worker.is_alive() + + +@pytest.mark.parametrize("broken_input", [False, True]) +async def test_publish_primes_keyframe_and_cleans_up( + monkeypatch: pytest.MonkeyPatch, broken_input: bool +) -> None: + import sys + from types import SimpleNamespace + + monkeypatch.setitem(sys.modules, "webm", webm) + spec = importlib.util.spec_from_file_location("webm_publisher", _EXAMPLE / "publish.py") + publisher = importlib.util.module_from_spec(spec) + spec.loader.exec_module(publisher) + sent = [] + closed = [] + unpublished = [] + + class Source: + def __init__(self, *args): + pass + + def capture_frame(self, frame): + sent.append(frame) + return True + + def take_feedback(self): + return rtc.EncodedVideoSourceFeedback( + False, rtc.EncodedRateControl(100000, 25) if len(sent) > 1 else None + ) + + async def aclose(self): + closed.append(True) + + async def subscribed(): + pass + + async def publish_track(*args): + return SimpleNamespace(sid="track", wait_for_subscription=subscribed) + + async def unpublish_track(sid): + unpublished.append(sid) + + monkeypatch.setattr(publisher.rtc, "EncodedVideoSource", Source) + monkeypatch.setattr(publisher.rtc.LocalVideoTrack, "create_video_track", lambda *args: None) + room = SimpleNamespace( + local_participant=SimpleNamespace( + publish_track=publish_track, unpublish_track=unpublish_track + ) + ) + + def frames(): + yield rtc.EncodedVideoFrame( + b"key", 16, 16, rtc.VideoCodec.VP9, rtc.EncodedFrameType.ENCODED_FRAME_KEY, 0 + ) + if broken_input: + raise ValueError("broken stream") + yield rtc.EncodedVideoFrame( + b"delta", 16, 16, rtc.VideoCodec.VP9, rtc.EncodedFrameType.ENCODED_FRAME_DELTA, 40000 + ) + + if broken_input: + with pytest.raises(ValueError, match="broken stream"): + await publisher.publish(room, frames(), name="test") + else: + task = asyncio.create_task(publisher.publish(room, frames(), name="test")) + try: + + async def wait_for_frames(): + while len(sent) < 3: + await asyncio.sleep(0.01) + + await asyncio.wait_for(wait_for_frames(), 5) + assert not task.done() # EOF must not immediately tear down the track. + finally: + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + assert [f.data for f in sent] == ( + [b"key", b"key"] if broken_input else [b"key", b"key", b"delta"] + ) + assert sent[1].timestamp_us > sent[0].timestamp_us + assert closed == [True] + assert unpublished == ["track"] From d57c4659d6cd898e9a41b410614c53ec3ee985c5 Mon Sep 17 00:00:00 2001 From: AmitMY Date: Wed, 30 Sep 2026 12:01:39 +0200 Subject: [PATCH 4/5] test(webm): prevent seeking anonymous pipes on Windows --- tests/rtc/test_webm_example.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/tests/rtc/test_webm_example.py b/tests/rtc/test_webm_example.py index 5482c227..ffa7239a 100644 --- a/tests/rtc/test_webm_example.py +++ b/tests/rtc/test_webm_example.py @@ -99,6 +99,11 @@ def test_opaque_webm(codec: str) -> None: def test_demux_before_pipe_eof(alpha_webm: bytes) -> None: + class PipeReader(io.BufferedReader): + def seekable(self) -> bool: + # Windows can report an anonymous pipe as seekable; FFmpeg must not seek it. + return False + read_fd, write_fd = os.pipe() decoded_first = threading.Event() writer_closed = threading.Event() @@ -117,7 +122,7 @@ def write() -> None: writer.start() try: with ( - os.fdopen(read_fd, "rb") as pipe, + PipeReader(io.FileIO(read_fd, "rb")) as pipe, av.open( pipe, format="webm", From 6ef36d8bb02c6be08f914b00b8c951f52a1a9e72 Mon Sep 17 00:00:00 2001 From: AmitMY Date: Wed, 30 Sep 2026 12:20:08 +0200 Subject: [PATCH 5/5] feat(examples): preserve transparent WebM using existing byte streams --- examples/preencoded-webm/README.md | 88 ++------ examples/preencoded-webm/publish.py | 13 +- examples/preencoded-webm/webm.py | 68 +------ examples/webm-byte-stream/README.md | 60 ++++++ examples/webm-byte-stream/stream_webm.py | 75 +++++++ tests/rtc/test_webm_example.py | 243 ++++++++++++++++++----- 6 files changed, 353 insertions(+), 194 deletions(-) create mode 100644 examples/webm-byte-stream/README.md create mode 100644 examples/webm-byte-stream/stream_webm.py diff --git a/examples/preencoded-webm/README.md b/examples/preencoded-webm/README.md index a4f338cb..11a8dd2e 100644 --- a/examples/preencoded-webm/README.md +++ b/examples/preencoded-webm/README.md @@ -37,68 +37,23 @@ 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. -## Experimental single-track alpha transport +## Transparency and container-preserving transport -**This is a backend wire-format prototype, not interoperable transparent WebRTC.** -Opaque publishing works with existing VP8/VP9 receivers. Alpha mode must only be -used with matching test receivers in an isolated room; an ordinary browser/mobile -LiveKit video track does not recognize the additional data. No browser or mobile -playback adapter is implemented here. +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. -By default, the demuxer rejects alpha WebM instead of silently dropping its alpha. -To opt into the experimental VP9 envelope: +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. -```sh -.venv/bin/python examples/preencoded-webm/publish.py transparent.webm \ - --room alpha-experiment --experimental-alpha -``` - -The prototype extracts the VP9 color access unit and the separately compressed -alpha access unit from WebM `BlockAdditional` with `BlockAddID=1`. It publishes both -inside **one encoded-frame payload on one media track**, sharing one frame type, -timestamp, and RTP frame boundary. There is no second track or data channel. -The native passthrough encoder forwards the bytes, and normal RTP packetization -can split them across multiple packets. - -The provisional payload is: - -```text -[color access unit][alpha access unit] -[color length: u32 big endian][alpha length: u32 big endian]["LKWA"][version: u8 = 1] -``` - -`LKWA` and version 1 are local proposal values, not an assigned LiveKit format. -Lengths exclude the 13-byte footer. Both components must be nonempty and the -entire payload must be at most 16 MiB in this example. A receiver validates the -footer and sizes after RTP reassembly, then extracts the original compressed -components **before** handing anything to an ordinary VP9 decoder. The included -`unpack_alpha` function specifies that byte contract; it is not a video renderer. -Sparse/reused alpha blocks and dimension changes are not supported. A WebM -keyframe must contain independently decodable color and alpha components. - -This deliberately does not overload LiveKit's `FrameMetadata.user_data`: the -native LKTS trailer has a one-byte length and a maximum total size of 255 bytes, -while alpha access units can be many kilobytes. The media envelope precedes any -native metadata trailer. Combining it with frame metadata or E2EE needs separate -interoperability testing; this example does not enable either. - -### Required before this could become a supported feature - -- Agree on a versioned transport profile with LiveKit maintainers. The example's - track name identifies a prototype; it does **not** negotiate receiver support. -- Add capability signaling and enforce it for every subscriber, including late - joins. Unsupported receivers must be rejected or explicitly offered an opaque - fallback. Today's server does not strip this custom alpha envelope for them. -- Implement a receiver path for each target platform. A browser adapter could - investigate remuxing compressed color and alpha into WebM for native playback; - browser/MSE alpha support and latency must be tested. Native mobile SDKs need - their own supported decode/playback path. No cross-platform guarantee follows - from successful byte transport. -- Validate loss/reordering/retransmission, congestion feedback, keyframe recovery, - codec profiles, resource limits, metadata, encryption, and reconnection. - -These are interoperability requirements, not additional Python encoder settings. -Keep the alpha mode experimental until they are resolved. +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 @@ -106,14 +61,7 @@ Keep the alpha mode experimental until they are resolved. uv run --with 'av>=16.1' pytest tests/rtc/test_encoded_video.py tests/rtc/test_webm_example.py ``` -Tests generate synthetic alpha WebM, verify unchanged compressed payloads and -key/delta timestamps, read the first frame before a pipe reaches EOF, reject -malformed/oversized envelopes, and check cancellation joins the active demux read. -The base API's server-backed test publishes VP9 and verifies subscriber pixels. - -A separate local transport probe used Python FFI 0.12.80, LiveKit server 1.13.7, -and a Go SDK 2.18.1 subscriber to reassemble RTP without decoding. A 37,630-byte -compound frame (27,578 color bytes, 10,039 alpha bytes, 13-byte footer) crossed the -server in 33 RTP packets with identical SHA-256 before and after. This demonstrates -transport through the tested stack; it does not validate an alpha player or -compatibility with every LiveKit deployment. +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. diff --git a/examples/preencoded-webm/publish.py b/examples/preencoded-webm/publish.py index 8f1963e2..9aaa6162 100644 --- a/examples/preencoded-webm/publish.py +++ b/examples/preencoded-webm/publish.py @@ -107,7 +107,7 @@ async def main(args: argparse.Namespace) -> None: buffer_size=4096, options={"probesize": "4096", "analyzeduration": "0"}, ) as container: - frames = demux(container, experimental_alpha=args.experimental_alpha) + frames = demux(container) room = rtc.Room() token = ( api.AccessToken() @@ -117,11 +117,7 @@ async def main(args: argparse.Namespace) -> None: ) try: await room.connect(os.environ["LIVEKIT_URL"], token) - await publish( - room, - frames, - name="experimental-webm-alpha-v1" if args.experimental_alpha else "webm", - ) + await publish(room, frames, name="webm") finally: await room.disconnect() @@ -130,10 +126,5 @@ async def main(args: argparse.Namespace) -> None: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("input", help="WebM file or streaming HTTP(S) URL") parser.add_argument("--room", default="preencoded-webm") - parser.add_argument( - "--experimental-alpha", - action="store_true", - help="use the unnegotiated alpha prototype; matching test receivers only", - ) logging.basicConfig(level=logging.INFO) asyncio.run(main(parser.parse_args())) diff --git a/examples/preencoded-webm/webm.py b/examples/preencoded-webm/webm.py index a82ced71..90753669 100644 --- a/examples/preencoded-webm/webm.py +++ b/examples/preencoded-webm/webm.py @@ -1,58 +1,18 @@ -"""Demux WebM without decoding; prototype a single-track VP9 alpha envelope. - -The alpha envelope is experimental, not a negotiated LiveKit or WebRTC format. -It requires a matching receiver to strip it before invoking a video decoder. -""" +"""Demux opaque VP8/VP9 WebM for pre-encoded media-track publishing.""" from __future__ import annotations -import struct from typing import Iterator import av from livekit import rtc -# Provisional wire format for isolated interoperability experiments only: -# [color access unit][alpha access unit][color_len:u32][alpha_len:u32][LKWA][v1] -# Both lengths are big endian. The footer is inside the encoded frame, before -# any native LiveKit metadata trailer. Do not reuse the 255-byte LKTS metadata -# envelope: alpha access units can span many RTP packets. -_FOOTER = struct.Struct(">II4sB") -_MAX_FRAME_BYTES = 16 * 1024 * 1024 - - -def pack_alpha(color: bytes, alpha: bytes) -> bytes: - """Package two compressed VP9 access units without inspecting their pixels.""" - if not color or not alpha: - raise ValueError("color and alpha access units must not be empty") - if len(color) + len(alpha) + _FOOTER.size > _MAX_FRAME_BYTES: - raise ValueError("alpha envelope exceeds the 16 MiB example limit") - return color + alpha + _FOOTER.pack(len(color), len(alpha), b"LKWA", 1) - - -def unpack_alpha(data: bytes) -> tuple[bytes, bytes]: - """Validate the experimental envelope and return its compressed components. - Included to specify and test the wire contract, not as a playback adapter. - """ - if not _FOOTER.size < len(data) <= _MAX_FRAME_BYTES: - raise ValueError("invalid alpha envelope size") - color_size, alpha_size, magic, version = _FOOTER.unpack_from(data, len(data) - _FOOTER.size) - if magic != b"LKWA" or version != 1: - raise ValueError("unknown alpha envelope") - if not color_size or not alpha_size or color_size + alpha_size != len(data) - _FOOTER.size: - raise ValueError("invalid alpha envelope lengths") - return data[:color_size], data[color_size : -_FOOTER.size] - - -def demux( - container: av.container.InputContainer, *, experimental_alpha: bool = False -) -> Iterator[rtc.EncodedVideoFrame]: +def demux(container: av.container.InputContainer) -> Iterator[rtc.EncodedVideoFrame]: """Read timestamped VP8/VP9 packets, including from non-seekable live WebM. - Alpha is rejected by default rather than silently discarded. The opt-in - prototype requires VP9 and an alpha access unit in every packet; sparse alpha - and changing dimensions are outside this example's scope. + 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") @@ -61,13 +21,8 @@ def demux( codec = codecs.get(stream.codec_context.name) if codec is None: raise ValueError("this example supports VP8/VP9 WebM only") - has_alpha = stream.metadata.get("alpha_mode") == "1" - if has_alpha and not experimental_alpha: - raise ValueError( - "WebM has alpha; use the experimental alpha receiver/profile or an opaque source" - ) - if experimental_alpha and (not has_alpha or codec != rtc.VideoCodec.VP9): - raise ValueError("experimental alpha requires VP9 WebM with AlphaMode=1") + 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): @@ -80,15 +35,8 @@ def demux( raise ValueError("WebM presentation timestamps must increase") last_timestamp = timestamp data = bytes(packet) - additional = bytes(packet.get_sidedata("matroska_block_additional")) - if additional: - if len(additional) <= 8 or int.from_bytes(additional[:8], "big") != 1: - raise ValueError("unsupported WebM BlockAdditional mapping") - if not experimental_alpha: - raise ValueError("WebM alpha cannot be discarded implicitly") - data = pack_alpha(data, additional[8:]) - elif experimental_alpha: - raise ValueError("missing alpha access unit in WebM 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, diff --git a/examples/webm-byte-stream/README.md b/examples/webm-byte-stream/README.md new file mode 100644 index 00000000..4b09a68d --- /dev/null +++ b/examples/webm-byte-stream/README.md @@ -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. diff --git a/examples/webm-byte-stream/stream_webm.py b/examples/webm-byte-stream/stream_webm.py new file mode 100644 index 00000000..c726a767 --- /dev/null +++ b/examples/webm-byte-stream/stream_webm.py @@ -0,0 +1,75 @@ +"""Forward intact WebM from a file or HTTP response over a LiveKit byte stream.""" + +from __future__ import annotations + +import argparse +import asyncio +import os +from collections.abc import AsyncIterator + +import aiofiles +import aiohttp +from livekit import api, rtc + +_CHUNK_SIZE = 15_000 + + +async def read_chunks(source: str) -> AsyncIterator[bytes]: + if source.startswith(("http://", "https://")): + timeout = aiohttp.ClientTimeout(total=None, sock_connect=5, sock_read=5) + async with aiohttp.ClientSession(timeout=timeout) as session: + async with session.get(source) as response: + response.raise_for_status() + async for chunk in response.content.iter_chunked(_CHUNK_SIZE): + yield chunk + else: + async with aiofiles.open(source, "rb") as file: + while chunk := await file.read(_CHUNK_SIZE): + yield chunk + + +async def send_webm( + participant: rtc.LocalParticipant, chunks: AsyncIterator[bytes], *, recipient: str +) -> None: + writer = await participant.stream_bytes( + name="video.webm", + topic="webm", + mime_type="video/webm", + destination_identities=[recipient], + ) + reason = "WebM source interrupted" + try: + async for chunk in chunks: + await writer.write(chunk) + reason = "" + finally: + await writer.aclose(reason=reason) + + +async def main(args: argparse.Namespace) -> None: + room = rtc.Room() + token = ( + api.AccessToken() + .with_identity("webm-byte-publisher") + .with_grants(api.VideoGrants(room_join=True, room=args.room)) + .to_jwt() + ) + try: + await room.connect(os.environ["LIVEKIT_URL"], token) + if args.recipient not in room.remote_participants: + raise ValueError("recipient must join and register its 'webm' handler before sending") + chunks = read_chunks(args.input) + try: + await send_webm(room.local_participant, chunks, recipient=args.recipient) + finally: + await chunks.aclose() + finally: + await room.disconnect() + + +if __name__ == "__main__": + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("input", help="valid streaming WebM file or HTTP(S) URL") + parser.add_argument("--room", default="webm-stream") + parser.add_argument("--recipient", required=True, help="identity of the prepared receiver") + asyncio.run(main(parser.parse_args())) diff --git a/tests/rtc/test_webm_example.py b/tests/rtc/test_webm_example.py index ffa7239a..9f483502 100644 --- a/tests/rtc/test_webm_example.py +++ b/tests/rtc/test_webm_example.py @@ -6,7 +6,6 @@ import importlib.util import io import os -import struct import threading from fractions import Fraction from pathlib import Path @@ -25,14 +24,13 @@ _spec.loader.exec_module(webm) -@pytest.fixture(scope="module") -def alpha_webm() -> bytes: - """Synthetic media only; alpha is intentionally larger than one RTP packet.""" +def make_webm(*, alpha: bool) -> bytes: + """Generate synthetic streaming WebM without private media fixtures.""" buf = io.BytesIO() with av.open(buf, "w", format="webm", options={"live": "1", "cluster_time_limit": "1"}) as c: stream = c.add_stream("libvpx-vp9", rate=25) stream.width = stream.height = 160 - stream.pix_fmt = "yuva420p" + stream.pix_fmt = "yuva420p" if alpha else "yuv420p" stream.options = { "deadline": "realtime", "cpu-used": "8", @@ -49,29 +47,35 @@ def alpha_webm() -> bytes: return buf.getvalue() -def test_demux_preserves_compressed_color_and_alpha(alpha_webm: bytes) -> None: - with av.open(io.BytesIO(alpha_webm)) as c: +@pytest.fixture(scope="module") +def opaque_webm() -> bytes: + return make_webm(alpha=False) + + +@pytest.fixture(scope="module") +def alpha_webm() -> bytes: + return make_webm(alpha=True) + + +def test_demux_preserves_compressed_video(opaque_webm: bytes) -> None: + with av.open(io.BytesIO(opaque_webm)) as c: packets = [p for p in c.demux(video=0) if p.size] - with av.open(io.BytesIO(alpha_webm)) as c: - # No decode call: payloads must be byte-for-byte identical to the demuxer. - frames = list(webm.demux(c, experimental_alpha=True)) + with av.open(io.BytesIO(opaque_webm)) as c: + frames = list(webm.demux(c)) assert len(frames) == len(packets) == 12 assert [f.timestamp_us for f in frames] == list(range(0, 480000, 40000)) for frame, packet in zip(frames, packets): - color, alpha = webm.unpack_alpha(frame.data) - assert color == bytes(packet) - assert alpha == bytes(packet.get_sidedata("matroska_block_additional"))[8:] + assert frame.data == bytes(packet) assert frame.frame_type == ( rtc.EncodedFrameType.ENCODED_FRAME_KEY if packet.is_keyframe else rtc.EncodedFrameType.ENCODED_FRAME_DELTA ) - assert len(webm.unpack_alpha(frames[0].data)[1]) > 1200 -def test_alpha_requires_explicit_opt_in(alpha_webm: bytes) -> None: +def test_reject_alpha_webm(alpha_webm: bytes) -> None: with av.open(io.BytesIO(alpha_webm)) as c: - with pytest.raises(ValueError, match="alpha"): + with pytest.raises(ValueError, match="alpha.*preserves the WebM container"): next(webm.demux(c)) @@ -93,12 +97,9 @@ def test_opaque_webm(codec: str) -> None: packets = [bytes(p) for p in c.demux(video=0) if p.size] with av.open(io.BytesIO(buf.getvalue())) as c: assert [f.data for f in webm.demux(c)] == packets - with av.open(io.BytesIO(buf.getvalue())) as c: - with pytest.raises(ValueError, match="AlphaMode"): - next(webm.demux(c, experimental_alpha=True)) -def test_demux_before_pipe_eof(alpha_webm: bytes) -> None: +def test_demux_before_pipe_eof(opaque_webm: bytes) -> None: class PipeReader(io.BufferedReader): def seekable(self) -> bool: # Windows can report an anonymous pipe as seekable; FFmpeg must not seek it. @@ -111,8 +112,8 @@ def seekable(self) -> bool: def write() -> None: try: with os.fdopen(write_fd, "wb") as pipe: - for start in range(0, len(alpha_webm), 4096): - pipe.write(alpha_webm[start : start + 4096]) + for start in range(0, len(opaque_webm), 4096): + pipe.write(opaque_webm[start : start + 4096]) pipe.flush() decoded_first.wait(5) finally: @@ -130,7 +131,7 @@ def write() -> None: options={"probesize": "4096", "analyzeduration": "0"}, ) as c, ): - frames = webm.demux(c, experimental_alpha=True) + frames = webm.demux(c) first = next(frames) assert not writer_closed.is_set() decoded_first.set() @@ -142,32 +143,6 @@ def write() -> None: assert not writer.is_alive() -@pytest.mark.parametrize( - "data", - [ - b"", - b"short", - b"caa" + struct.pack(">II4sB", 1, 2, b"LKWA", 2), - b"caa" + struct.pack(">II4sB", 1, 3, b"LKWA", 1), - b"caa" + struct.pack(">II4sB", 0, 3, b"LKWA", 1), - b"caa" + struct.pack(">II4sB", 1, 2, b"NOPE", 1), - ], -) -def test_reject_bad_envelopes(data: bytes) -> None: - with pytest.raises(ValueError): - webm.unpack_alpha(data) - - -def test_bound_envelope_size(monkeypatch: pytest.MonkeyPatch) -> None: - monkeypatch.setattr(webm, "_MAX_FRAME_BYTES", 32) - with pytest.raises(ValueError, match="limit"): - webm.pack_alpha(b"c" * 16, b"a" * 16) - with pytest.raises(ValueError): - webm.pack_alpha(b"", b"a") - with pytest.raises(ValueError): - webm.unpack_alpha(b"x" * 33) - - async def test_cancel_waits_for_reader(monkeypatch: pytest.MonkeyPatch) -> None: import sys @@ -197,7 +172,7 @@ def frames(): release.set() -def test_http_stream_before_response_finishes(alpha_webm: bytes) -> None: +def test_http_stream_before_response_finishes(opaque_webm: bytes) -> None: from http.server import BaseHTTPRequestHandler, HTTPServer first_frame, response_finished = threading.Event(), threading.Event() @@ -207,7 +182,7 @@ def do_GET(self) -> None: self.send_response(200) self.send_header("Content-Type", "video/webm") self.end_headers() - self.wfile.write(alpha_webm) + self.wfile.write(opaque_webm) self.wfile.flush() first_frame.wait(5) response_finished.set() @@ -225,9 +200,9 @@ def log_message(self, *args) -> None: timeout=(5.0, 5.0), options={"probesize": "4096", "analyzeduration": "0"}, ) as c: - frame = next(webm.demux(c, experimental_alpha=True)) + frame = next(webm.demux(c)) assert not response_finished.is_set() - assert webm.unpack_alpha(frame.data)[1] + assert frame.data finally: first_frame.set() worker.join(timeout=5) @@ -316,3 +291,165 @@ async def wait_for_frames(): assert sent[1].timestamp_us > sent[0].timestamp_us assert closed == [True] assert unpublished == ["track"] + + +_byte_example = _EXAMPLE.parent / "webm-byte-stream" / "stream_webm.py" +_byte_spec = importlib.util.spec_from_file_location("stream_webm", _byte_example) +byte_publisher = importlib.util.module_from_spec(_byte_spec) +_byte_spec.loader.exec_module(byte_publisher) + + +@pytest.mark.parametrize("failure", [None, ValueError, asyncio.CancelledError]) +async def test_byte_publisher_preserves_chunks_and_closes(alpha_webm: bytes, failure) -> None: + from types import SimpleNamespace + + written, reasons, options = [], [], [] + + async def write(chunk): + written.append(chunk) + + async def close(**kwargs): + reasons.append(kwargs["reason"]) + + async def open_stream(**kwargs): + options.append(kwargs) + return SimpleNamespace(write=write, aclose=close) + + async def chunks(): + yield alpha_webm[:15000] + if failure: + raise failure() + yield alpha_webm[15000:] + + participant = SimpleNamespace(stream_bytes=open_stream) + if failure: + with pytest.raises(failure): + await byte_publisher.send_webm(participant, chunks(), recipient="viewer") + assert b"".join(written) == alpha_webm[:15000] + assert reasons == ["WebM source interrupted"] + else: + await byte_publisher.send_webm(participant, chunks(), recipient="viewer") + assert b"".join(written) == alpha_webm + assert reasons == [""] + assert options == [ + dict( + name="video.webm", + topic="webm", + mime_type="video/webm", + destination_identities=["viewer"], + ) + ] + + +async def test_byte_file_reader(alpha_webm: bytes, tmp_path: Path) -> None: + path = tmp_path / "alpha.webm" + path.write_bytes(alpha_webm) + chunks = [chunk async for chunk in byte_publisher.read_chunks(str(path))] + assert all(0 < len(chunk) <= 15000 for chunk in chunks) + assert b"".join(chunks) == alpha_webm + + +async def test_byte_http_reader_before_eof(alpha_webm: bytes) -> None: + from aiohttp import web + + first_received = asyncio.Event() + + async def handle(request): + response = web.StreamResponse(headers={"Content-Type": "video/webm"}) + await response.prepare(request) + await response.write(alpha_webm[:15000]) + await asyncio.wait_for(first_received.wait(), 5) + await response.write(alpha_webm[15000:]) + return response + + app = web.Application() + app.router.add_get("/video.webm", handle) + runner = web.AppRunner(app) + await runner.setup() + site = web.TCPSite(runner, "127.0.0.1", 0) + await site.start() + port = runner.addresses[0][1] + chunks = byte_publisher.read_chunks(f"http://127.0.0.1:{port}/video.webm") + try: + first = await asyncio.wait_for(chunks.__anext__(), 5) + assert first and alpha_webm.startswith(first) + first_received.set() + received = first + b"".join([chunk async for chunk in chunks]) + assert received == alpha_webm + finally: + first_received.set() + await chunks.aclose() + await runner.cleanup() + + +async def test_byte_stream_preserves_alpha_webm_before_eof(alpha_webm: bytes) -> None: + """Round-trip a complete standard WebM through the real LiveKit byte API.""" + import contextlib + import uuid + + from livekit import api + + if not all(os.getenv(key) for key in ("LIVEKIT_URL", "LIVEKIT_API_KEY", "LIVEKIT_API_SECRET")): + pytest.skip("LiveKit server credentials are required") + + room_name = f"webm-bytes-{uuid.uuid4().hex}" + sender, receiver = rtc.Room(), rtc.Room() + first_received = asyncio.Event() + completed = asyncio.get_running_loop().create_future() + consumers = [] + + async def consume(reader): + try: + assert reader.info.mime_type == "video/webm" + received = bytearray() + async for chunk in reader: + received.extend(chunk) + first_received.set() + completed.set_result(bytes(received)) + except Exception as exc: + first_received.set() + if not completed.done(): + completed.set_exception(exc) + finally: + reader.close() + + def on_stream(reader, identity): + consumers.append(asyncio.create_task(consume(reader))) + + receiver.register_byte_stream_handler("webm", on_stream) + + def token(identity): + return ( + api.AccessToken() + .with_identity(identity) + .with_grants(api.VideoGrants(room_join=True, room=room_name)) + .to_jwt() + ) + + async def chunks(): + yield alpha_webm[:15000] + # Do not produce the rest until the receiver sees a chunk. This fails + # if delivery buffers the whole WebM until EOF. + await asyncio.wait_for(first_received.wait(), 10) + yield alpha_webm[15000:] + + try: + await receiver.connect(os.environ["LIVEKIT_URL"], token("viewer")) + await sender.connect(os.environ["LIVEKIT_URL"], token("sender")) + await asyncio.wait_for( + byte_publisher.send_webm(sender.local_participant, chunks(), recipient="viewer"), 20 + ) + received = await asyncio.wait_for(completed, 10) + assert received == alpha_webm + with av.open(io.BytesIO(received)) as container: + assert container.streams.video[0].metadata["alpha_mode"] == "1" + packets = [p for p in container.demux(video=0) if p.size] + assert len(packets) == 12 + assert all(p.get_sidedata("matroska_block_additional") for p in packets) + finally: + for task in consumers: + task.cancel() + with contextlib.suppress(asyncio.CancelledError): + await task + await sender.disconnect() + await receiver.disconnect()