Python
moq-rs on PyPI (the moq name was taken), imported as moq. It wraps the generated moq-ffi bindings in asyncio: async context managers for sessions, async iterators for announcements, groups, and frames, and no Moq prefixes. Python 3.10+, with wheels for Linux x86_64/aarch64, macOS arm64, and Windows x64.
pip install moq-rs # or: uv add moq-rsimport asyncio, moq
async def main():
async with moq.Client("https://cdn.moq.dev/anon") as client:
# The filter is relative to the literal prefix; prefixes stay origin-relative.
async for event in client.announced("live/", filter="*/camera"):
if not isinstance(event, moq.AnnounceEventStart):
continue # AnnounceEventUpdate, AnnounceEventEnd, or AnnounceEventLive
announcement = event.announce
print(announcement.captures) # what * matched, or None for a partial overlap
broadcast = await client.request_broadcast(announcement.prefix)
catalog = await broadcast.catalog()
name, track = next(iter(catalog.audio.items()))
async for frame in await broadcast.subscribe_media(name, track):
print(frame.timestamp_us, len(frame.payload))
asyncio.run(main())import asyncio, moq
async def main():
# opus_init_bytes, payload, pts, and rgba come from your encoder or capture source.
async with moq.Client("https://cdn.moq.dev/anon") as client:
broadcast = client.create_broadcast("my-stream.hang")
# Already-encoded frames: the catalog is filled from the bitstream
audio = broadcast.publish_audio(moq.AudioFormat.OPUS, opus_init_bytes)
audio.write_frame(payload, timestamp_us=0)
audio.cut() # audio has no keyframes, so this is what gives it groups
# Or raw pixels, encoded inside the binding (VideoToolbox, Media Foundation, NVENC, openh264)
video = broadcast.encode_video(
moq.VideoEncoderInput(format=moq.VideoPixelFormat.RGBA, width=1280, height=720, framerate=30),
moq.VideoEncoderOutput(
codec=moq.VideoCodec.H264,
track="camera",
kind=moq.VideoEncoderKind.AUTO(), # pyright: ignore[reportArgumentType]
),
)
video.write(moq.VideoFrame(timestamp_us=pts, data=rgba))
# Raw bytes and JSON
events = broadcast.publish_track("events")
events.write_frame(b'{"cmd": "ready"}', 0)
status = broadcast.publish_json_snapshot("status", compression=True)
status.update({"state": "live", "viewers": 42})
broadcast.announce()
audio.finish()
video.finish()
events.finish()
status.finish()
broadcast.close()
asyncio.run(main())For already-encoded live output, call audio.flush(timestamp_us) after each audio.write_frame with the same broadcast-clock PTS. It samples the transport handoff for catalog jitter. File, pipe, and network imports should omit flush; raw-pixel and PCM encoders inside the binding measure their own output.
Call audio.discontinuity() when the source seeks, pauses, or changes its time base. It publishes a timeline marker and restarts handoff measurement without lowering advertised jitter. Resume with timestamps that continue forward on the broadcast media clock; this does not permit timestamp rewinds. On a track from publish_video or publish_video_on_track, resume with a keyframe: a delta frame before it fails.
The three advertising operations, as the other bindings spell them: client.create_broadcast(path) (or OriginProducer.create_broadcast) returns an unannounced producer, invisible to everyone; broadcast.announce(route) / broadcast.unannounce() own that exact-path advertisement, and broadcast.close() ends the broadcast for good (a second call is a no-op); origin.dynamic(prefix, route) claims prefix and every path beneath it ("" for everything). Hold the returned handle while the claim should stay advertised, and reject the requests you will not serve. A route is a capability, not an inventory. announced(prefix, filter=...) combines a literal root with an optional relative pattern and yields AnnounceEvents: AnnounceEventStart, AnnounceEventUpdate, or AnnounceEventEnd carrying an Announce as .announce, whose .prefix stays relative to the origin and whose .captures reports what the pattern wildcards matched, or AnnounceEventLive once every route live at subscribe time has been delivered. Break on AnnounceEventLive to list what is live and stop. Paths with a .-prefixed segment below the prefix are hidden unless hidden=True.
Sessions reconnect with backoff when the transport drops and re-announce local broadcasts. session.epoch() counts the connections, 1 on the first, pairing with session.status() to log each reconnect; moq.Backoff tunes the pacing (timeout_us=0 retries forever); and moq.connect(..., max_streams=...) raises the peer's inbound stream cap.
The WebSocket fallback races QUIC after a 200 ms head start. Pass websocket_enabled=False to moq.connect for a QUIC-only relay, or a websocket_delay timedelta to change the head start.
Everything in the shared feature list is here: moq.Server with per-request accept/reject, fetch_group and fetch_media_group, dynamic() handlers for on-demand tracks and dynamic(prefix) for broadcasts, append_datagram/recv_datagram, set_catalog_section, and a producer's demand(), a TrackDemand whose used()/unused() let capture idle when nobody is subscribed. request.set_publish/set_consume raise if the request is already answered, cancelled, or currently accepting. session.bandwidth() divides the connection's send estimate; pass it to encode_video / encode_audio or reserve a share for an app-owned track. moq.is_auth(err) and moq.is_shutdown(err) classify errors. moq.protocol_error(err) is the structured protocol failure (scope, verbatim code, kind) when the peer sent one. str(err) is the Rust error message. Catch moq.Error.Busy when a setter races an in-flight connect, listen, or accept. Each server request reports a moq.Transport enum, including QUIC, Iroh, WebSocket, TCP, and Unix sockets.
encode_audio encodes raw PCM inside the binding. Its codec is an object, moq.AudioCodec.opus() or moq.AudioCodec.aac(), and AudioEncoderOutput.frame_duration_us sets the Opus frame length: 2500, 5000, 10000, 20000 (the default), 40000, or 60000. 0 takes the codec's own frame, which AAC needs. AAC-LC encodes through the platform's encoder, so a host without one refuses it.
Audio channels also names the speaker layout, by the WAVE convention: 1 is mono, 2 stereo, 3 2.1, 4 quad, 5 5.0, 6 5.1, 7 6.1, and 8 7.1, interleaved front left, front right, center, LFE, back, then side. Decoding remixes to the count you ask for; past 8 channels the samples pass through but can't be remixed.
Each frame from decode_video owns its decoded picture until it is released, including after the consumer is cancelled. frame.pixels(format) converts it on demand: VideoPixelFormat.I420, or VideoPixelFormat.RGBA for four bytes a pixel. Drop frames promptly, since held frames hold decoder buffers. resize is best effort: only NVDEC has a built-in scaler, so read each frame's own width() and height() rather than assuming it took. VideoDecoderOutput(surface=True) keeps the decoder's surface for frame.surface() instead of downloading it. Only macOS has one, so decode_video fails as unsupported elsewhere.
Connection stats
session.stats() returns a ConnectionStats snapshot. Each field is None when the transport backend does not report it (native QUIC reports all of them; browser WebTransport reports few or none) or before it is available, which is not the same as zero.
| Field | Unit | Meaning |
|---|---|---|
rtt_us | microseconds | Smoothed round-trip time. |
estimated_send_rate_bps | bits per second | Send bandwidth from the congestion controller. |
estimated_recv_rate_bps | bits per second | Receive bandwidth from MoQ PROBE. |
bytes_sent | bytes | Total sent, including retransmissions and overhead. |
bytes_received | bytes | Total received, including duplicates and overhead. |
bytes_lost | bytes | Total lost, detected via retransmission or acknowledgement. |
packets_sent | datagrams | Total datagrams sent. |
packets_received | datagrams | Total datagrams received. |
packets_lost | datagrams | Total datagrams detected as lost. |
- API reference: moq-rs.readthedocs.io
- Source and examples:
py/moq-rs - Raw bindings:
moq-ffion PyPI, for the unwrapped API
Raw track publisher metadata has an optional maximum age. Omitting it imposes no publisher age limit; zero keeps the live edge. Local cache limits still apply, and media imports explicitly retain 30 seconds. See publisher retention.
Await session.shutdown() to drain finished tracks before disconnecting. It raises if delivery has not completed within one second. cancel(code) stays immediate. Session and client async context managers await shutdown on a clean exit and cancel on an error, so the body's exception survives; finish or abort live tracks first. IETF media streams are not drained yet.