From e117706b3a044ebd9f9be715e7dd25d8cc7f1d31 Mon Sep 17 00:00:00 2001 From: AmitMY Date: Wed, 30 Sep 2026 11:43:01 +0200 Subject: [PATCH 1/2] 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/2] 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/"]'