moq-rs 0.5.1


pip install moq-rs

  Latest version

Released: Sep 26, 2026


Meta
Author: Luke Curley
Requires Python: >=3.10

Classifiers

Development Status
  • 3 - Alpha

Intended Audience
  • Developers

Programming Language
  • Python :: 3

Topic
  • Multimedia :: Sound/Audio
  • Multimedia :: Video

moq

Python bindings for Media over QUIC: real-time pub/sub with built-in caching, fan-out, and prioritization, on top of QUIC.

Installed as moq-rs (the moq name is taken on PyPI), imported as moq.

It wraps the auto-generated moq-ffi UniFFI bindings with a Pythonic API: no Moq prefixes, async iterators, context managers, and simplified connection setup. At session setup it negotiates either the moq-lite or moq-transport wire protocol.

Installation

pip install moq-rs

# or with uv
uv add moq-rs

This pulls in the moq-ffi native bindings automatically. moq-rs is pure Python and is versioned independently of moq-ffi; it floats to the latest compatible moq-ffi patch.

Quick Start

Subscribe to a stream

import asyncio
import moq


async def main():
    async with moq.connect("https://cdn.moq.dev/anon") as client:
        async for announcement in client.announced():
            # A route covers a prefix and carries no broadcast, so resolve the path.
            broadcast = await client.request_broadcast(announcement.prefix)
            catalog = await broadcast.catalog()

            for name, track in catalog.audio.items():
                frames = await broadcast.subscribe_media(name, track)
                async with frames:
                    async for frame in frames:
                        print(f"Got frame: {len(frame.payload)} bytes, ts={frame.timestamp_us}")


asyncio.run(main())

Publish a stream

import asyncio
import moq


async def main():
    async with moq.Client("https://cdn.moq.dev/anon") as client:
        broadcast = client.create_broadcast("my-stream")

        # Publish an Opus audio track (init bytes from your encoder)
        audio = broadcast.publish_audio(moq.AudioFormat.OPUS, opus_init_bytes)

        # Write frames
        # Audio has no keyframes, so `cut` is what gives it group boundaries.
        audio.write_frame(payload, timestamp_us=0)
        audio.cut()
        audio.write_frame(payload, timestamp_us=20000)
        audio.cut()

        broadcast.announce()

        # Clean up
        audio.finish()
        broadcast.close()


asyncio.run(main())

Host a server

import asyncio
import moq


async def main():
    async with moq.Server("127.0.0.1:4443", tls_generate=["localhost"]) as server:
        broadcast = server.create_broadcast("hello")
        track = broadcast.publish_track("events")
        broadcast.announce()
        print(f"listening on https://{server.local_addr}")

        sessions = []
        async for request in server:
            safe_path = request.path.split("?", 1)[0]
            safe_url = request.url.split("?", 1)[0] if request.url else None
            print(f"  + {request.transport} {safe_path} from {safe_url}")
            sessions.append(await request.accept())


asyncio.run(main())

Reject a request instead of accepting it with await request.reject(403).

Advanced: Manual origin wiring

For full control over the origin topology:

import moq

origin = moq.OriginProducer()
client = moq.Client(
    "https://cdn.moq.dev/anon",
    publish=origin,
    subscribe=origin,
)

API

Connection

  • connect(url, *, tls_verify=True, tls_roots=None, tls_system_roots=None, tls_fingerprints=None, tls_cert=None, tls_key=None, bind=None, max_streams=None, reconnect=True, backoff=None, publish=None, subscribe=None). Shorthand for Client(...); use as async with moq.connect(url) as client:.
  • Client(url, *, tls_verify=True, tls_roots=None, tls_system_roots=None, tls_fingerprints=None, tls_cert=None, tls_key=None, bind=None, max_streams=None, reconnect=True, backoff=None, publish=None, subscribe=None). Async context manager for connecting to a relay.
    • tls_roots. PEM root certificate file path(s) to trust instead of the system roots.
    • tls_system_roots. Whether to trust platform roots in addition to custom roots.
    • tls_fingerprints. Hex SHA-256 fingerprint(s) to pin the peer's certificate to, the native equivalent of serverCertificateHashes. Accepts the values a server reports via cert_fingerprints(), so you can trust a self-signed certificate without tls_verify=False.
    • tls_cert, tls_key. Paired PEM certificate chain and private key paths for mTLS.
    • max_streams. Raise the peer's inbound stream cap.
    • reconnect, backoff. Redial with a Backoff when the transport drops; reconnect=False dials once.
    • .session. The established Session (or None before connecting / after exit).
  • Server(bind="[::]:443", *, tls_cert=(), tls_key=(), tls_generate=(), publish=None, subscribe=None). Async context manager + async iterator of incoming Requests.
    • .local_addr. The bound address (useful when binding to port 0).
    • .cert_fingerprints(). SHA-256 fingerprints of the configured TLS certificates, for serverCertificateHashes browser cert pinning.
    • .create_broadcast(path) → BroadcastProducer. Create an unannounced broadcast, invisible to everyone; announce() makes it discoverable and reachable; close() ends it.
  • Request. An incoming session, yielded by async for request in server.
    • .url, .path, .query, .transport. The query-free path is uniform across transports; the root or missing path is "". The encoded query may contain credentials.
    • .set_publish(origin), .set_consume(origin). Per-request overrides, captured at accept(). Raise if the request is already answered, cancelled, or currently accepting.
    • await .accept() → Session. Complete the handshake (hold the result to keep the connection alive).
    • await .reject(code). Reject with an application error code; 401 and 403 map to unauthorized.
    • .cancel(). Cancel an in-flight accept()/reject() call.
  • Session. An established connection. Holding it keeps the connection alive; it is also an async with context manager that shuts down on exit.
    • await .closed(). Wait until the session closes.
    • .cancel(code), .shutdown(). Close with an error code, or gracefully (code 0).
    • .publish() → OriginProducer, .consume() → OriginConsumer. The wired origin sides.
    • .stats() → ConnectionStats. Snapshot RTT, bandwidth estimates, and byte/packet counters.
    • await .status() → ConnectionStatus, .epoch(). Watch reconnects; the epoch counts connections, 1 on the first.
    • .bandwidth() → Bandwidth. Divide the send estimate between encoders and app-owned tracks.

Publishing

  • BroadcastProducer(). Create a broadcast to publish tracks into.
    • .dynamic() → BroadcastDynamic
    • .publish_audio(format, init, *, label=None, track=None) → MediaProducer. init is required: an OpusHead or AudioSpecificConfig resolves the whole rendition. track names the track; otherwise a unique name is derived from the format.
    • .publish_video(format, init=b"", *, label=None, hint=None, track=None) → MediaProducer. init may be empty for a format that resolves in band; a VideoHint pins catalog fields the stream can't reveal (bitrate) or publishes the catalog before the first keyframe. track names the track as in publish_audio.
    • .encode_video(input, output, *, bandwidth=None) → VideoProducer. Encode raw VideoFrames inside the binding; .write(frame) each one.
    • .encode_audio(name, input, output, *, bandwidth=None) → AudioProducer. Encode raw PCM AudioFrames; the codec is output.codec, e.g. AudioCodec.opus(), with output.frame_duration_us setting the Opus frame length.
    • .close() ends the broadcast for good; a second call is a no-op. .finish() is its deprecated alias.
  • BroadcastDynamic. Async source of tracks requested by subscribers.
    • await .requested_track() → TrackRequest. Call .accept() on it for a TrackProducer, or .abort(code) to reject.
    • Async iterator yielding TrackRequest
  • MediaProducer. Write frames to a track.
    • .write_frame(payload, timestamp_us=0)
    • .cut() / .seek(sequence) draw a group boundary (audio has none of its own)
    • .finish()
  • TrackProducer / GroupProducer. Write raw payloads with no codec parsing.
    • .write_frame(payload, timestamp_us=0) writes a payload with a presentation timestamp in microseconds.
    • .create_group(sequence) creates a sparse or replayed group at an explicit sequence.
    • .finish() ends at the live edge; the handle remains so .abort(error_code) can still run.
    • .finish_at(final_sequence) declares the first group that will never be produced while leaving lower groups writable.
    • .abort(error_code) terminates the track or group with an application error.
    • .append_datagram(payload, timestamp_us=0) -> sequence (TrackProducer) sends a best-effort datagram. Payloads are capped at 1200 bytes and there is no stream fallback.

Subscribing

  • BroadcastConsumer. Subscribe to tracks within a broadcast.
    • await .subscribe_catalog() → CatalogConsumer
    • await .subscribe_track(name, subscription=None) → TrackConsumer
    • await .subscribe_media(name, track, subscription=None) → MediaConsumer. track is the catalog record (e.g. catalog.video[name]); its container tells the decoder how to parse the bitstream.
    • await .catalog() → Catalog (convenience)
  • CatalogConsumer. Async iterator of Catalog.
  • MediaConsumer. Async iterator of MediaFrame.
  • TrackConsumer. Async iterator of raw groups, in sequence order.
    • await .next_group() → GroupConsumer | None. Sequence order; what the default iteration yields.
    • await .recv_group() → GroupConsumer | None. Arrival order, which may be out of sequence. Prefer it when latency matters more than order.
    • .groups_as_arrived(). Async iterator over recv_group().
    • .read_frame() -> Frame | None returns the first timestamped frame of the next group. Empty groups are skipped; None is track EOF.
    • await .recv_datagram() -> Datagram | None for best-effort raw track datagrams.
    • .info() → TrackInfo
    • .update(subscription). Change delivery priority, staleness, or group range after subscribing.
  • GroupConsumer. Async iterator of timestamped Frames.
    • .read_frame() -> Frame | None returns a timestamped raw frame.

Every handle whose cleanup is cancel() is an async context manager, so exiting async with releases it: the consumers (CatalogConsumer, MediaConsumer, MediaGroupConsumer, TrackConsumer, AudioConsumer, GroupConsumer, JsonSnapshotConsumer, JsonStreamConsumer, AnnounceConsumer, AnnouncedBroadcast) and the dynamic sources (OriginDynamic, BroadcastDynamic, TrackDynamic).

Origin (advanced)

  • OriginProducer(*, cache_capacity_bytes=None). Manage broadcast announcements. Set cache_capacity_bytes to bound cached groups under this origin.
    • .consume() → OriginConsumer
    • .dynamic(prefix, route=Route()) → OriginDynamic
    • .create_broadcast(path) → BroadcastProducer
  • OriginDynamic. Async source of broadcasts requested by consumers.
    • await .requested_broadcast() → BroadcastRequest. Call .accept(broadcast) to serve it, or .reject(code) to fail the requester.
    • Async iterator yielding BroadcastRequest
  • OriginConsumer. Discover broadcasts.
    • .announced(prefix, filter=None) → AnnounceConsumer (async iterator); filter is a pattern relative to the literal prefix, while each update's .prefix stays origin-relative and .captures reports wildcard matches
    • .announced_broadcast(path) → AnnouncedBroadcast (awaitable, waits until something serves the path)
    • .request_broadcast(path) → BroadcastConsumer (awaitable; announced now or a dynamic fallback, else raises)

Types

  • Catalog. .audio: dict[str, Audio], .video: dict[str, Video], .display, .rotation, .flip.
  • Frame. .payload: bytes, .timestamp_us: int. The unit of every write and every raw read.
  • MediaFrame. .payload: bytes, .timestamp_us: int, .keyframe: bool. Returned by media subscriptions. keyframe marks a group start or video keyframe; for audio it is true only at a group start.
  • Datagram. .sequence: int, .timestamp_us: int, .payload: bytes. Delivered only on datagram-capable transports and lite-05 or newer moq-lite.
  • Audio. .codec, .sample_rate, .channel_count, .bitrate, .description.
  • Video. .codec, .coded: Dimensions, .display_aspect, .bitrate, .stalled, .framerate, .description. A true .stalled recommends temporarily avoiding the rendition without making it unavailable.
  • Subscription. Subscriber delivery preferences: priority, staleness, and optional group range.
  • TrackInfo. Publisher track properties: priority, cache window, and timescale.
  • Dimensions. .width: int, .height: int.
  • Container. The catalog container enum, carried on each Video/Audio record.

Logging and errors

  • log_level(level="info"). Initialize logging for the underlying Rust layer ("error", "warn", "info", "debug", "trace"). Call once per process.
  • Error. The exception raised by all operations. Catch a specific case via its variants, e.g. except moq.Error.AlreadyResponded:, except moq.Error.Cancelled:, or except moq.Error.Busy: when a setter races an in-flight connect/listen/accept.
  • is_shutdown(err). True for Cancelled and Closed, which arise from graceful shutdown rather than an actual failure. Use it to break out of an async for without treating the expected end-of-stream error as a problem.
  • is_auth(err). True for Unauthorized (HTTP 401) and Forbidden (HTTP 403), and for a protocol Unauthorized session close. Retrying without new credentials won't help, so surface these rather than reconnect.
  • protocol_error(err). The structured protocol failure (session or stream scope, verbatim wire code, known kind) when the peer sent one.

See Also

  • moq-ffi. The raw UniFFI bindings this package wraps. Use it directly only if you need the unwrapped Moq-prefixed API.
  • MoQ project. Full monorepo with Rust server, TypeScript browser lib, and more.

Wheel compatibility matrix

Platform Python 3
any

Files in release

Extras: None
Dependencies:
moq-ffi (~=0.4.3)