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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
38 changes: 25 additions & 13 deletions benchmarks/media.py
Original file line number Diff line number Diff line change
Expand Up @@ -204,7 +204,9 @@ async def run(self) -> VideoResult:
# enough for the size, so the encoder doesn't drop frames for bitrate
bitrate = int(self.width * self.height * self.fps * 0.2)
remote = await _remote_track(caller, callee, generator.track, max_bitrate=bitrate)
processor = webrtc.MediaStreamTrackProcessor(remote, max_buffer_size=self.max_buffer_size)
processor = webrtc.MediaStreamTrackProcessor(
webrtc.MediaStreamTrackProcessorInit(remote, max_buffer_size=self.max_buffer_size)
)
sampling = asyncio.ensure_future(run.sample_rss(self.rss_every)) if self.rss_every else None
lag = await run.phases.measure(
run.result, processor, write=lambda: run.write(generator), read=lambda: run.read(processor)
Expand Down Expand Up @@ -243,13 +245,18 @@ async def write(self, generator: webrtc.VideoTrackGenerator) -> None:
if self.phases.measuring.is_set():
self.result.sent += 1
await writer.write(
webrtc.VideoFrame(frame, format='I420', coded_width=width, coded_height=height, timestamp=number)
webrtc.VideoFrame(
frame,
webrtc.VideoFrameBufferInit(
format='I420', coded_width=width, coded_height=height, timestamp=number
),
)
)
number += 1
await asyncio.sleep(max(0.0, start + number / fps - loop.time()))

async def read(self, processor: webrtc.MediaStreamTrackProcessor) -> None:
header = {'rect': {'x': 0, 'y': 0, 'width': BITS * BLOCK, 'height': BLOCK}}
header = webrtc.VideoFrameCopyToOptions(rect=webrtc.DOMRectInit(x=0, y=0, width=BITS * BLOCK, height=BLOCK))
loop = asyncio.get_running_loop()
async for frame in processor.readable:
now = loop.time()
Expand Down Expand Up @@ -310,7 +317,7 @@ async def audio_loopback(channels: int = 2, seconds: float = 60, warmup: float =
async with _connection() as (caller, callee):
generator = webrtc.MediaStreamTrackGenerator('audio')
remote = await _remote_track(caller, callee, generator)
processor = webrtc.MediaStreamTrackProcessor(remote, max_buffer_size=50)
processor = webrtc.MediaStreamTrackProcessor(webrtc.MediaStreamTrackProcessorInit(remote, max_buffer_size=50))

async def write() -> None:
writer = generator.writable.get_writer()
Expand All @@ -319,12 +326,14 @@ async def write() -> None:
while not phases.done.is_set():
await writer.write(
webrtc.AudioData(
format='s16',
sample_rate=48000,
number_of_frames=480,
number_of_channels=channels,
timestamp=written * 10_000,
data=chunk,
webrtc.AudioDataInit(
format='s16',
sample_rate=48000,
number_of_frames=480,
number_of_channels=channels,
timestamp=written * 10_000,
data=chunk,
)
)
)
written += 1
Expand Down Expand Up @@ -370,10 +379,11 @@ async def copy_costs(
for width, height in sizes:
chroma = (width // 2) * (height // 2)
frame = webrtc.VideoFrame(
bytes(width * height + 2 * chroma), format='I420', coded_width=width, coded_height=height, timestamp=0
bytes(width * height + 2 * chroma),
webrtc.VideoFrameBufferInit(format='I420', coded_width=width, coded_height=height, timestamp=0),
)
for format in formats:
options = {'format': format}
options = webrtc.VideoFrameCopyToOptions(format=format)
destination = bytearray(frame.allocation_size(options))
runs = 0
start = time.perf_counter()
Expand All @@ -391,6 +401,8 @@ def construct_cost(width: int, height: int, budget: float = 1.0) -> float:
runs = 0
start = time.perf_counter()
while time.perf_counter() - start < budget:
webrtc.VideoFrame(data, format='I420', coded_width=width, coded_height=height, timestamp=0).close()
webrtc.VideoFrame(
data, webrtc.VideoFrameBufferInit(format='I420', coded_width=width, coded_height=height, timestamp=0)
).close()
runs += 1
return (time.perf_counter() - start) / runs * 1000
16 changes: 9 additions & 7 deletions docs/source/media.rst
Original file line number Diff line number Diff line change
Expand Up @@ -22,9 +22,10 @@ Receiving

@pc.on('track')
async def on_track(event):
async for frame in webrtc.MediaStreamTrackProcessor(event.track).readable:
rgba = bytearray(frame.allocation_size({'format': 'RGBA'}))
await frame.copy_to(rgba, {'format': 'RGBA'})
rgba_options = webrtc.VideoFrameCopyToOptions(format='RGBA')
async for frame in webrtc.MediaStreamTrackProcessor(webrtc.MediaStreamTrackProcessorInit(event.track)).readable:
rgba = bytearray(frame.allocation_size(rgba_options))
await frame.copy_to(rgba, rgba_options)
frame.close()

Sending
Expand All @@ -35,20 +36,21 @@ Sending
generator = webrtc.VideoTrackGenerator()
pc.add_track(generator.track)
writer = generator.writable.get_writer()
await writer.write(webrtc.VideoFrame(i420, format='I420', coded_width=640, coded_height=480, timestamp=0))
init = webrtc.VideoFrameBufferInit(format='I420', coded_width=640, coded_height=480, timestamp=0)
await writer.write(webrtc.VideoFrame(i420, init))

microphone = webrtc.MediaStreamTrackGenerator('audio')
pc.add_track(microphone)
await microphone.writable.get_writer().write(
webrtc.AudioData(format='s16', sample_rate=48000, number_of_frames=480, number_of_channels=1,
timestamp=0, data=pcm)
webrtc.AudioData(webrtc.AudioDataInit(format='s16', sample_rate=48000, number_of_frames=480,
number_of_channels=1, timestamp=0, data=pcm))
)

Transforming
------------

.. code-block:: python

processor = webrtc.MediaStreamTrackProcessor(track)
processor = webrtc.MediaStreamTrackProcessor(webrtc.MediaStreamTrackProcessorInit(track))
generator = webrtc.VideoTrackGenerator()
await processor.readable.pipe_through(webrtc.TransformStream(transformer)).pipe_to(generator.writable)
7 changes: 7 additions & 0 deletions docs/source/webrtc.models.dictionary.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
webrtc.models.dictionary
========================

.. automodule:: webrtc.models.dictionary
:members:
:undoc-members:
:show-inheritance:
3 changes: 2 additions & 1 deletion docs/source/webrtc.models.rst
Original file line number Diff line number Diff line change
Expand Up @@ -14,15 +14,16 @@ Submodules

webrtc.models.audio_data
webrtc.models.blob
webrtc.models.dictionary
webrtc.models.events
webrtc.models.media_track_constraints
webrtc.models.rtc_certificate
webrtc.models.rtc_configuration
webrtc.models.rtc_ice_candidate
webrtc.models.rtc_rtp_transceiver_init
webrtc.models.rtc_session_description
webrtc.models.rtc_session_description_init
webrtc.models.rtc_stats
webrtc.models.rtp_parameters
webrtc.models.rtp_source
webrtc.models.rtp_transceiver_init
webrtc.models.video_frame
7 changes: 7 additions & 0 deletions docs/source/webrtc.models.rtc_rtp_transceiver_init.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
webrtc.models.rtc\_rtp\_transceiver\_init
========================================

.. automodule:: webrtc.models.rtc_rtp_transceiver_init
:members:
:undoc-members:
:show-inheritance:
7 changes: 0 additions & 7 deletions docs/source/webrtc.models.rtp_transceiver_init.rst

This file was deleted.

21 changes: 11 additions & 10 deletions examples/echo.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,32 +24,33 @@

async def grayscale(frame: webrtc.VideoFrame, controller: webrtc.TransformStreamDefaultController) -> None:
"""Transforms an I420 frame: U and V at 128 leave only the luma."""
data = bytearray(frame.allocation_size({'format': 'I420'}))
await frame.copy_to(data, {'format': 'I420'})
i420 = webrtc.VideoFrameCopyToOptions(format='I420')
data = bytearray(frame.allocation_size(i420))
await frame.copy_to(data, i420)
luma = frame.coded_width * frame.coded_height
data[luma:] = b'\x80' * (len(data) - luma)
controller.enqueue(
webrtc.VideoFrame(
data,
format='I420',
coded_width=frame.coded_width,
coded_height=frame.coded_height,
timestamp=frame.timestamp,
webrtc.VideoFrameBufferInit(
format='I420', coded_width=frame.coded_width, coded_height=frame.coded_height, timestamp=frame.timestamp
),
)
)
frame.close()


async def watch(track: webrtc.MediaStreamTrack) -> None:
"""Reads the echoed frames for a while, then prints whether the last one is gray."""
reader = webrtc.MediaStreamTrackProcessor(track).readable.get_reader()
reader = webrtc.MediaStreamTrackProcessor(webrtc.MediaStreamTrackProcessorInit(track)).readable.get_reader()
loop = asyncio.get_running_loop()
end = loop.time() + SECONDS
frames = 0
while loop.time() < end:
frame = (await reader.read()).value
rgba = bytearray(frame.allocation_size({'format': 'RGBA'}))
await frame.copy_to(rgba, {'format': 'RGBA'})
options = webrtc.VideoFrameCopyToOptions(format='RGBA')
rgba = bytearray(frame.allocation_size(options))
await frame.copy_to(rgba, options)
frame.close()
frames += 1
red, green, blue = rgba[0:3]
Expand Down Expand Up @@ -91,7 +92,7 @@ async def main() -> None:

@echo.on('track')
def on_echo_track(event: webrtc.RTCTrackEvent) -> None:
readable = webrtc.MediaStreamTrackProcessor(event.track).readable
readable = webrtc.MediaStreamTrackProcessor(webrtc.MediaStreamTrackProcessorInit(event.track)).readable
pipe = readable.pipe_through(webrtc.TransformStream({'transform': grayscale})).pipe_to(generator.writable)
pipes.append(asyncio.ensure_future(pipe))

Expand Down
17 changes: 11 additions & 6 deletions examples/janus_streaming.py
Original file line number Diff line number Diff line change
Expand Up @@ -129,10 +129,13 @@ def fullscreen() -> Iterator[None]:
async def watch(track: webrtc.MediaStreamTrack) -> None:
"""Draws the frames of the video track until it ends."""
# a buffer of one frame drops the frames the terminal is too slow for
async for frame in webrtc.MediaStreamTrackProcessor(track, max_buffer_size=1).readable:
async for frame in webrtc.MediaStreamTrackProcessor(
webrtc.MediaStreamTrackProcessorInit(track, max_buffer_size=1)
).readable:
with frame:
rgbx = bytearray(frame.allocation_size({'format': 'RGBX'}))
await frame.copy_to(rgbx, {'format': 'RGBX'})
options = webrtc.VideoFrameCopyToOptions(format='RGBX')
rgbx = bytearray(frame.allocation_size(options))
await frame.copy_to(rgbx, options)
size = frame.visible_rect
draw(rgbx, int(size.width), int(size.height))

Expand All @@ -141,9 +144,11 @@ async def listen(track: webrtc.MediaStreamTrack) -> None:
"""Plays the audio track until it ends."""
speakers = None
try:
async for data in webrtc.MediaStreamTrackProcessor(track, max_buffer_size=50).readable:
async for data in webrtc.MediaStreamTrackProcessor(
webrtc.MediaStreamTrackProcessorInit(track, max_buffer_size=50)
).readable:
with data:
options = {'plane_index': 0, 'format': 's16'}
options = webrtc.AudioDataCopyToOptions(plane_index=0, format='s16')
samples = bytearray(data.allocation_size(options))
data.copy_to(samples, options)
speakers = speakers or Speakers(int(data.sample_rate), data.number_of_channels)
Expand All @@ -157,7 +162,7 @@ async def answer(pc: webrtc.RTCPeerConnection, offer: dict[str, str]) -> dict[st
"""Answers with all the ICE candidates in the SDP, since there is no trickling."""
gathered = asyncio.Event()
pc.on('icegatheringstatechange', lambda _: pc.ice_gathering_state == 'complete' and gathered.set())
await pc.set_remote_description(offer)
await pc.set_remote_description(webrtc.RTCSessionDescriptionInit.from_json(offer))
await pc.set_local_description(await pc.create_answer())
with contextlib.suppress(asyncio.TimeoutError):
await asyncio.wait_for(gathered.wait(), 5)
Expand Down
25 changes: 15 additions & 10 deletions examples/openai_live.py
Original file line number Diff line number Diff line change
Expand Up @@ -238,7 +238,7 @@ async def run(self, hang_up: asyncio.Event) -> None:

console.info(f'Creating a {args.model} session...')
answer = await self._create_session(self.pc.local_description.sdp)
await self.pc.set_remote_description({'type': 'answer', 'sdp': answer})
await self.pc.set_remote_description(webrtc.RTCSessionDescriptionInit('answer', answer))

self.microphone.start()
self.tasks.append(asyncio.ensure_future(self._send_microphone(generator.writable.get_writer())))
Expand Down Expand Up @@ -312,23 +312,28 @@ async def _send_microphone(self, writer: webrtc.WritableStreamDefaultWriter) ->
)
guarded = not self.args.barge_in and time.monotonic() - self.assistant_spoke_at < ECHO_TAIL
data = webrtc.AudioData(
format='s16',
sample_rate=SAMPLE_RATE,
number_of_frames=FRAME,
number_of_channels=1,
timestamp=timestamp,
data=silence if guarded else chunk,
webrtc.AudioDataInit(
format='s16',
sample_rate=SAMPLE_RATE,
number_of_frames=FRAME,
number_of_channels=1,
timestamp=timestamp,
data=silence if guarded else chunk,
)
)
await writer.write(data)
timestamp += 10_000

async def _play(self, track: webrtc.MediaStreamTrack) -> None:
"""Plays the assistant's audio."""
heard = False
async for data in webrtc.MediaStreamTrackProcessor(track, max_buffer_size=50).readable:
async for data in webrtc.MediaStreamTrackProcessor(
webrtc.MediaStreamTrackProcessorInit(track, max_buffer_size=50)
).readable:
with data:
samples = bytearray(data.allocation_size({'plane_index': 0, 'format': 's16'}))
data.copy_to(samples, {'plane_index': 0, 'format': 's16'})
options = webrtc.AudioDataCopyToOptions(plane_index=0, format='s16')
samples = bytearray(data.allocation_size(options))
data.copy_to(samples, options)
rate, channels = int(data.sample_rate), data.number_of_channels
if peak(samples) > VOICE_LEVEL:
self.assistant_spoke_at = time.monotonic()
Expand Down
9 changes: 6 additions & 3 deletions examples/recorder.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,10 +30,13 @@ async def record(track: webrtc.MediaStreamTrack, file: BinaryIO) -> None:
"""Writes the frames of a track to a file until the track ends."""
frames = 0
with file:
async for media in webrtc.MediaStreamTrackProcessor(track, max_buffer_size=30).readable:
async for media in webrtc.MediaStreamTrackProcessor(
webrtc.MediaStreamTrackProcessorInit(track, max_buffer_size=30)
).readable:
if track.kind == 'audio':
data = bytearray(media.allocation_size({'plane_index': 0}))
media.copy_to(data, {'plane_index': 0})
options = webrtc.AudioDataCopyToOptions(plane_index=0)
data = bytearray(media.allocation_size(options))
media.copy_to(data, options)
else:
data = bytearray(media.allocation_size())
await media.copy_to(data)
Expand Down
14 changes: 8 additions & 6 deletions examples/telegram_group_calls.py
Original file line number Diff line number Diff line change
Expand Up @@ -133,12 +133,14 @@ async def send_audio_data(generator: webrtc.MediaStreamTrackGenerator, file: Bin
frames = len(data) // 4
await writer.write(
webrtc.AudioData(
format='s16',
sample_rate=48000,
number_of_frames=frames,
number_of_channels=2,
timestamp=chunks * 10_000,
data=data[: frames * 4],
webrtc.AudioDataInit(
format='s16',
sample_rate=48000,
number_of_frames=frames,
number_of_channels=2,
timestamp=chunks * 10_000,
data=data[: frames * 4],
)
)
)
chunks += 1
Expand Down
Loading
Loading