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
2 changes: 1 addition & 1 deletion .github/scripts/fuzz.sh
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ if [ ! -x "$BUILD/venv/bin/python" ]; then
# built from source with this Clang, so its sanitizer runtime is the one the extension is instrumented for
CLANG_BIN="$(command -v clang)" LIBFUZZER_LIB="$(clang -print-runtime-dir)/libclang_rt.fuzzer_no_main.a" \
"$BUILD/venv/bin/pip" install -q --no-binary atheris atheris \
cmake ninja "pybind11>=3.0" pytest
cmake ninja "pybind11>=3.0" "typing_extensions>=4.10" pytest
fi
# shellcheck disable=SC1091
source "$BUILD/venv/bin/activate"
Expand Down
2 changes: 1 addition & 1 deletion .github/scripts/sanitizers-macos.sh
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ fi

if [ ! -x "$BUILD/venv/bin/python" ]; then
uv venv -q "$BUILD/venv" --python "$PYTHON"
uv pip install -q --python "$BUILD/venv/bin/python" cmake ninja "pybind11>=3.0" pytest pytest-asyncio pytest-timeout
uv pip install -q --python "$BUILD/venv/bin/python" cmake ninja "pybind11>=3.0" "typing_extensions>=4.10" pytest pytest-asyncio pytest-timeout
fi
export PATH="$BUILD/venv/bin:$PATH"

Expand Down
2 changes: 1 addition & 1 deletion .github/scripts/sanitizers.sh
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ dnf install -y -q clang lld compiler-rt llvm
"$PYTHON" -m venv "$BUILD/venv"
# shellcheck disable=SC1091
source "$BUILD/venv/bin/activate"
python -m pip install -q cmake ninja "pybind11>=3.0" pytest pytest-asyncio pytest-timeout
python -m pip install -q cmake ninja "pybind11>=3.0" "typing_extensions>=4.10" pytest pytest-asyncio pytest-timeout

CC=clang CXX=clang++ cmake -S "$SRC" -B "$BUILD" -G Ninja \
-DCMAKE_BUILD_TYPE=RelWithDebInfo \
Expand Down
3 changes: 3 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,9 @@ jobs:
- uses: astral-sh/setup-uv@v7
- run: uvx ruff check
- run: uvx ruff format --check
# the imports of the checked code, without building the extension: its stub is in stubs/
- run: uv sync --group dev --group wpt --no-install-project
- run: uvx pyrefly check
- run: make format-check
# clang-tidy reads the libwebrtc headers, from the cache the wheels use
- uses: actions/cache@v6
Expand Down
7 changes: 6 additions & 1 deletion Makefile
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
.PHONY: dev test asan tsan fuzz lint format format-check tidy stub wheels doc clean
.PHONY: dev test asan tsan fuzz lint typecheck format format-check tidy stub wheels doc clean

# pinned to the clang-tidy of .github/scripts/tidy.sh
CLANG_FORMAT := uvx clang-format==22.1.8
Expand Down Expand Up @@ -30,6 +30,9 @@ lint: format-check
uvx ruff check
uvx ruff format --check

typecheck:
uvx pyrefly check

format:
uvx ruff check --fix
uvx ruff format
Expand All @@ -45,6 +48,8 @@ tidy:
stub:
uv run --no-sync pybind11-stubgen wrtc -o build/stubs
cp build/stubs/wrtc.pyi stubs/wrtc/__init__.pyi
# collections.abc.Buffer is 3.12+
perl -pi -e 's/collections\.abc\.Buffer/typing_extensions.Buffer/g; s/^import typing$$/import typing\nimport typing_extensions/' stubs/wrtc/__init__.pyi

# wheels for the current platform, exactly as CI builds them
wheels:
Expand Down
2 changes: 1 addition & 1 deletion benchmarks/__main__.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ def _commit() -> str:
try:
commit = subprocess.check_output([git, 'rev-parse', '--short', 'HEAD'], text=True).strip()
dirty = subprocess.check_output([git, 'status', '--porcelain'], text=True).strip()
return commit + (' with uncommitted changes' if dirty else '')
return commit + (' with uncommitted changes' if dirty != '' else '')
except (OSError, subprocess.CalledProcessError):
return 'unknown'

Expand Down
18 changes: 10 additions & 8 deletions benchmarks/measure.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@

def percentile(values: Sequence[float], fraction: float) -> float:
"""The value below which the fraction of the values are, or NaN without values."""
if not values:
if len(values) == 0:
return float('nan')
ordered = sorted(values)
return ordered[min(len(ordered) - 1, int(fraction * len(ordered)))]
Expand All @@ -46,7 +46,7 @@ def __init__(self, interval: float = 0.005) -> None:
self.interval = interval
# compact: a 5 ms probe collects 12000 a minute
self.lags = array.array('d')
self._task: asyncio.Task | None = None
self._task: asyncio.Task[None] | None = None

async def _probe(self) -> None:
loop = asyncio.get_running_loop()
Expand All @@ -62,7 +62,7 @@ def __enter__(self) -> Self:
def __exit__(
self, exc_type: type[BaseException] | None, exc: BaseException | None, traceback: TracebackType | None
) -> None:
if self._task:
if self._task is not None:
self._task.cancel()

@property
Expand Down Expand Up @@ -103,27 +103,29 @@ def stop(self) -> Usage:
@property
def cpu_percent(self) -> float:
"""Of one core: the process uses several threads (encoders, decoders, network)."""
return self.cpu / self.wall * 100 if self.wall else float('nan')
return self.cpu / self.wall * 100 if self.wall != 0 else float('nan')


def slope_mb_per_minute(samples: Sequence[tuple[float, int]]) -> float:
"""The trend of (seconds, bytes) samples, by least squares: NaN for less than two."""
if not samples:
if len(samples) == 0:
return float('nan')
xs = [t for t, _ in samples]
ys = [b / 1e6 for _, b in samples]
mean_x, mean_y = statistics.fmean(xs), statistics.fmean(ys)
denominator = sum((x - mean_x) ** 2 for x in xs)
if not denominator:
if denominator == 0:
return float('nan')
return sum((x - mean_x) * (y - mean_y) for x, y in zip(xs, ys)) / denominator * 60


def machine() -> str:
"""The CPU, the OS and Python of the machine."""
cpu = platform.processor() or platform.machine()
cpu = platform.processor()
if cpu == '':
cpu = platform.machine()
sysctl = shutil.which('sysctl') if sys.platform == 'darwin' else None
if sysctl:
if sysctl is not None:
with contextlib.suppress(OSError, subprocess.CalledProcessError):
cpu = subprocess.check_output([sysctl, '-n', 'machdep.cpu.brand_string'], text=True).strip()
return (
Expand Down
56 changes: 42 additions & 14 deletions benchmarks/media.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@
from tests.helpers import connect, rss_bytes, wait_for_event

if TYPE_CHECKING:
from collections.abc import AsyncIterator, Awaitable, Iterable, Sequence
from collections.abc import AsyncGenerator, Awaitable, Iterable, Sequence

# the frame number, drawn as bits in blocks of luma at the top left of each frame
BITS = 16
Expand Down Expand Up @@ -85,7 +85,7 @@ def received_sizes(self) -> str:

def _frames(width: int, height: int, count: int = 10) -> list[bytearray]:
"""I420 frames of a gradient with a bar moving across, so the encoder has motion to encode."""
frames = []
frames: list[bytearray] = []
chroma = (width // 2) * (height // 2)
row = bytes((x * 255 // width) for x in range(width))
base = bytearray(row * height + bytes([128]) * (2 * chroma))
Expand All @@ -101,14 +101,14 @@ def _frames(width: int, height: int, count: int = 10) -> list[bytearray]:

def _draw_number(frame: bytearray, width: int, number: int) -> None:
for bit in range(BITS):
value = LUMA_ONE if number >> bit & 1 else LUMA_ZERO
value = LUMA_ONE if number >> bit & 1 != 0 else LUMA_ZERO
block = bytes([value]) * BLOCK
for y in range(BLOCK):
start = y * width + bit * BLOCK
frame[start : start + BLOCK] = block


def _read_number(luma: bytes, stride: int) -> int:
def _read_number(luma: bytes | bytearray, stride: int) -> int:
number = 0
for bit in range(BITS):
# the center of the block, away from the blur of compression at its edges
Expand All @@ -118,8 +118,24 @@ def _read_number(luma: bytes, stride: int) -> int:
return number


def _video(media: object) -> webrtc.VideoFrame:
"""The media a processor of a video track reads, as a video frame."""
if not isinstance(media, webrtc.VideoFrame):
msg = f'expected a video frame, not {media!r}'
raise TypeError(msg)
return media


def _audio(media: object) -> webrtc.AudioData:
"""The media a processor of an audio track reads, as audio data."""
if not isinstance(media, webrtc.AudioData):
msg = f'expected audio data, not {media!r}'
raise TypeError(msg)
return media


@contextlib.asynccontextmanager
async def _connection() -> AsyncIterator[tuple[webrtc.RTCPeerConnection, webrtc.RTCPeerConnection]]:
async def _connection() -> AsyncGenerator[tuple[webrtc.RTCPeerConnection, webrtc.RTCPeerConnection], None]:
caller, callee = webrtc.RTCPeerConnection(), webrtc.RTCPeerConnection()
try:
yield caller, callee
Expand All @@ -138,12 +154,16 @@ async def _remote_track(
sender = caller.add_track(track)
track_event = wait_for_event(callee, 'track', 30)
await connect(caller, callee, 30)
if max_bitrate:
if max_bitrate is not None and max_bitrate != 0:
parameters = sender.get_parameters()
for encoding in parameters.encodings:
encoding.max_bitrate = max_bitrate
await sender.set_parameters(parameters)
return (await track_event).track
event = await track_event
if not isinstance(event, webrtc.RTCTrackEvent):
msg = f'expected a track event, not {event!r}'
raise TypeError(msg)
return event.track


@dataclass
Expand Down Expand Up @@ -207,11 +227,15 @@ async def run(self) -> VideoResult:
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
sampling = (
asyncio.ensure_future(run.sample_rss(self.rss_every))
if self.rss_every is not None and self.rss_every != 0
else None
)
lag = await run.phases.measure(
run.result, processor, write=lambda: run.write(generator), read=lambda: run.read(processor)
)
if sampling:
if sampling is not None:
await sampling
run.result.lag_p95_ms, run.result.lag_max_ms = lag.p95_ms, lag.max_ms
generator.track.stop()
Expand Down Expand Up @@ -258,15 +282,16 @@ async def write(self, generator: webrtc.VideoTrackGenerator) -> None:
async def read(self, processor: webrtc.MediaStreamTrackProcessor) -> None:
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:
async for media in processor.readable:
now = loop.time()
frame = _video(media)
luma = bytearray(frame.allocation_size(header))
await frame.copy_to(luma, header)
size = (frame.coded_width, frame.coded_height)
frame.close()
if self.phases.measuring.is_set():
self._count(size, _read_number(luma, BITS * BLOCK), now)
if self.options.consumer_delay:
if self.options.consumer_delay != 0:
await asyncio.sleep(self.options.consumer_delay)
if self.phases.done.is_set():
break
Expand Down Expand Up @@ -342,7 +367,8 @@ async def write() -> None:
await asyncio.sleep(max(0.0, start + written / 100 - loop.time()))

async def read() -> None:
async for audio in processor.readable:
async for media in processor.readable:
audio = _audio(media)
if phases.measuring.is_set():
result.received += 1
result.received_frames += audio.number_of_frames
Expand Down Expand Up @@ -372,10 +398,12 @@ def megapixels_per_second(self) -> float:


async def copy_costs(
sizes: Iterable[tuple[int, int]], formats: Sequence[str] = ('I420', 'RGBA', 'BGRA'), budget: float = 1.0
sizes: Iterable[tuple[int, int]],
formats: Sequence[webrtc.VideoPixelFormatValue] = ('I420', 'RGBA', 'BGRA'),
budget: float = 1.0,
) -> list[CopyResult]:
"""How long VideoFrame.copy_to takes: a copy of the planes, or a conversion to RGB."""
results = []
results: list[CopyResult] = []
for width, height in sizes:
chroma = (width // 2) * (height // 2)
frame = webrtc.VideoFrame(
Expand Down
21 changes: 11 additions & 10 deletions cmake/libcxx/update.py
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,7 @@ def git_fetch(url: str, ref: str, repo: str) -> str:
def resolve_llvm(deps: str, name: str) -> str:
"""Chromium mirrors each llvm-project dir as its own repo: finds the llvm-project commit with the same tree."""
match = re.search(rf"'{re.escape(LLVM_DIRS[name][0])}':\s*'([^'@]+)@([0-9a-f]+)'", deps)
if not match:
if match is None:
msg = f'{LLVM_DIRS[name][0]} is not in the WebRTC DEPS'
raise SystemExit(msg)
url, revision = match.groups()
Expand All @@ -158,25 +158,26 @@ def resolve_llvm(deps: str, name: str) -> str:
commits: list[Commit] = json.loads(gh(f'{REPO}/commits?path={name}&since={since}&until={until}&per_page=100'))
for commit in commits:
if tree_sha(commit['sha'], name) == tree:
return str(commit['sha'])
return commit['sha']
msg = f'no llvm-project commit has the {name} tree of {url}@{revision}'
raise SystemExit(msg)


def runtime_sources(tag: str) -> list[str]:
"""The libc++ and libc++abi sources of Chromium's Linux build, as llvm-project paths."""
sources = []
sources: list[str] = []
for gn, llvm in (('libc%2B%2B', 'libcxx'), ('libc%2B%2Babi', 'libcxxabi')):
for line in chromium(f'buildtools/third_party/{gn}/BUILD.gn', tag).decode().splitlines():
match = re.search(r'"//third_party/libc\+\+(?:abi)?/src/src/([^"]+\.cpp)"', line)
if match and not line.lstrip().startswith('#') and 'win32' not in match[1]:
if match is not None and not line.lstrip().startswith('#') and 'win32' not in match[1]:
sources.append(f'{llvm}/src/{match[1]}')
return sorted(set(sources))


def includes(path: str, data: bytes) -> set[str]:
"""The llvm-project paths a file may include: next to it, or from libcxx/src and llvm-libc."""
found = set()
found: set[str] = set()
include: bytes # one group, so findall gives its bytes
for include in INCLUDE.findall(data):
name = include.decode()
found.add(posixpath.normpath(posixpath.join(posixpath.dirname(path), name)))
Expand All @@ -193,13 +194,13 @@ def pin_runtime(commits: dict[str, str], tag: str) -> None:
pinned: dict[str, bytes] = {}
queue = list(sources)
with ThreadPoolExecutor(8) as pool:
while queue:
while len(queue) > 0:
pinned.update(zip(queue, pool.map(blob, [tree[p] for p in queue])))
found = set().union(*(includes(path, pinned[path]) for path in queue))
found = set[str]().union(*(includes(path, pinned[path]) for path in queue))
queue = sorted(c for c in found if c in tree and c not in pinned)
# CMake compiles every pinned libcxx/libcxxabi .cpp, so only the sources may be among them
stray = [p for p in pinned if p.endswith('.cpp') and p not in sources and not p.startswith('libc/')]
if stray:
if len(stray) > 0:
msg = f'sources include other .cpp files: {stray}'
raise SystemExit(msg)
lines = [f'{hashlib.sha256(data).hexdigest()} {path}\n' for path, data in sorted(pinned.items())]
Expand All @@ -222,7 +223,7 @@ def pin_headers(commit: str, headers: Path) -> None:
remote = {e['path']: e['sha'] for e in tree}
local = {str(p.relative_to(headers)): blob_sha(p.read_bytes()) for p in headers.rglob('*') if p.is_file()}
mismatches = [path for path, sha in local.items() if remote.get(path) != sha]
if mismatches:
if len(mismatches) > 0:
msg = f'the prebuilt headers differ from llvm-project@{commit}: {mismatches[:10]}'
raise SystemExit(msg)

Expand All @@ -240,7 +241,7 @@ def sha256(entry: TreeEntry) -> str:

def pin_config(tag: str) -> None:
"""Pins Chromium's build-generated libc++ config."""
config = []
config: list[str] = []
for name in CHROMIUM_CONFIG:
data = chromium(f'buildtools/third_party/libc%2B%2B/{name}', tag)
config.append(f'{hashlib.sha256(data).hexdigest()} {name}\n')
Expand Down
14 changes: 11 additions & 3 deletions examples/echo.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,14 +40,22 @@ async def grayscale(frame: webrtc.VideoFrame, controller: webrtc.TransformStream
frame.close()


def video_frame(media: object) -> webrtc.VideoFrame:
"""The media a processor of a video track reads, as a video frame."""
if not isinstance(media, webrtc.VideoFrame):
msg = f'expected a video frame, not {media!r}'
raise TypeError(msg)
return media


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(webrtc.MediaStreamTrackProcessorInit(track)).readable.get_reader()
loop = asyncio.get_running_loop()
end = loop.time() + SECONDS
frames = 0
frames, rgba = 0, bytearray()
while loop.time() < end:
frame = (await reader.read()).value
frame = video_frame((await reader.read()).value)
options = webrtc.VideoFrameCopyToOptions(format='RGBA')
rgba = bytearray(frame.allocation_size(options))
await frame.copy_to(rgba, options)
Expand All @@ -64,7 +72,7 @@ def trickle(caller: webrtc.RTCPeerConnection, callee: webrtc.RTCPeerConnection)
async def on_candidate(
event: webrtc.RTCPeerConnectionIceEvent, other: webrtc.RTCPeerConnection = other
) -> None:
if event.candidate:
if event.candidate is not None:
await other.add_ice_candidate(event.candidate)

pc.on('icecandidate', on_candidate)
Expand Down
Loading
Loading