Add stream helper utilities and IVF frame iterator (#1082)

* Add stream helper utilities and IVF frame iterator

* fix lint

* some cosmetics

* fix lint

* changes

* improve test

* improve types and test

* add todo for better bitrate calculation
This commit is contained in:
Harisreedhar
2026-05-11 16:36:23 +02:00
committed by henryruhs
parent cc8bfc1af4
commit a2aedc8814
3 changed files with 138 additions and 1 deletions
+67 -1
View File
@@ -1,11 +1,18 @@
import asyncio
from typing import Tuple
import math
import os
import subprocess
from typing import Iterator, Optional, Tuple, cast
from aiortc import MediaStreamTrack, QueuedVideoStreamTrack, RTCPeerConnection, RTCRtpSender
from aiortc.mediastreams import MediaStreamError
from av import VideoFrame
from starlette.datastructures import Headers
from starlette.types import Scope
from facefusion.common_helper import is_linux, is_macos
from facefusion.streamer import process_vision_frame
from facefusion.types import Resolution, StreamBuffer, WebSocketStreamMode
def process_stream_frame(target_stream_frame : VideoFrame) -> VideoFrame:
@@ -42,3 +49,62 @@ def on_video_track(rtc_connection : RTCPeerConnection, output_track : QueuedVide
if target_track.kind == 'video':
asyncio.create_task(process_and_enqueue(target_track, output_track))
def calculate_bitrate(resolution : Resolution) -> int: # TODO : improve the bitrate calculation
pixel_total = resolution[0] * resolution[1]
bitrate_factor = 3500 / math.sqrt(1920 * 1080)
return max(400, round(math.sqrt(pixel_total) * bitrate_factor))
def calculate_buffer_size(resolution : Resolution) -> int:
return calculate_bitrate(resolution) * 2
def get_websocket_stream_mode(scope : Scope) -> Optional[WebSocketStreamMode]:
protocol_header = Headers(scope = scope).get('Sec-WebSocket-Protocol')
if protocol_header:
for protocol in protocol_header.split(','):
websocket_stream_mode = protocol.strip()
if websocket_stream_mode in [ 'image', 'video' ]:
return cast(WebSocketStreamMode, websocket_stream_mode)
return None
def read_pipe_buffer(pipe_handle : int, size : int) -> Optional[bytes]:
byte_buffer = bytearray()
frame_data = os.read(pipe_handle, size - len(byte_buffer))
while frame_data:
byte_buffer += frame_data
if len(byte_buffer) == size:
return bytes(byte_buffer)
frame_data = os.read(pipe_handle, size - len(byte_buffer))
return None
def forward_stream_frame(process : subprocess.Popen[bytes]) -> Iterator[StreamBuffer]:
pipe_handle = process.stdout.fileno()
if is_linux() or is_macos():
os.set_blocking(pipe_handle, True)
header = read_pipe_buffer(pipe_handle, 32)
if header:
frame_header = read_pipe_buffer(pipe_handle, 12)
while frame_header:
frame_size = int.from_bytes(frame_header[0:4], 'little')
frame_data = read_pipe_buffer(pipe_handle, frame_size)
if frame_data:
yield frame_data
frame_header = read_pipe_buffer(pipe_handle, 12)
+2
View File
@@ -264,6 +264,8 @@ BenchmarkCycleSet = TypedDict('BenchmarkCycleSet',
WebcamMode = Literal['inline', 'udp', 'v4l2']
StreamMode = Literal['udp', 'v4l2']
WebSocketStreamMode = Literal['image', 'video']
StreamBuffer : TypeAlias = bytes
RtcOfferSet = TypedDict('RtcOfferSet',
{