Go
moq.dev/moq: context.Context cancellation, error returns, and Go 1.23 range-over-func iterators for live streams. The native core arrives as a prebuilt static library through the moq.dev/moq-ffi module, so go get is all it takes (CGO_ENABLED=1, the default on Unix). Targets: linux/amd64, linux/arm64, darwin/arm64 (macOS 12.3+), windows/amd64.
go get moq.dev/moq@latestimport "fmt"
import "log"
import "moq.dev/moq"
// Subscribe. The iterator is live, so run it in its own goroutine.
client, err := moq.Dial(ctx, "https://relay.example.com", moq.WithTLSRoots("ca.pem"))
if err != nil {
log.Fatal(err)
}
defer client.Close()
filter := "*/camera"
announced, err := client.Announced(moq.AnnounceOptions{Prefix: "live/", Filter: &filter})
if err != nil {
log.Fatal(err)
}
for event, err := range announced.All(ctx) {
if err != nil {
if moq.IsShutdown(err) { break }
log.Fatal(err)
}
ann, ok := event.(moq.AnnounceEventStart)
if !ok {
continue // AnnounceEventUpdate, AnnounceEventEnd, or AnnounceEventLive
}
// Prefix stays origin-relative; Captures reports what each wildcard matched.
fmt.Printf("captures: %v\n", ann.Announce.Captures)
broadcast, err := client.RequestBroadcast(ctx, ann.Announce.Prefix)
if err != nil {
log.Fatal(err)
}
catalog, err := broadcast.Catalog(ctx)
if err != nil {
log.Fatal(err)
}
fmt.Printf("%+v\n", catalog)
}// Publish encoded frames, or raw pixels with the codec inside the binding.
// opusInit, packet, pts, and rgba come from your encoder or capture source.
broadcast, _ := client.CreateBroadcast("my-stream.hang")
audio, _ := broadcast.PublishAudio(moq.AudioFormatOpus, opusInit)
_ = audio.WriteFrame(moq.Frame{Payload: packet, TimestampUs: 20_000})
track := "camera"
video, _ := broadcast.EncodeVideo(
moq.VideoEncoderInput{Format: moq.VideoPixelFormatRgba, Width: 1280, Height: 720, Framerate: 30},
moq.VideoEncoderOutput{Codec: moq.VideoCodecH264, Track: &track, Kind: moq.AutoEncoder()},
nil,
)
_ = video.Write(moq.VideoFrame{TimestampUs: pts, Data: rgba})
_ = broadcast.Announce(moq.Route{})
_ = audio.Finish()
_ = video.Finish()
broadcast.Close() // keep the producer reachable while publishing, then close explicitlyFor locally encoded media, call MediaProducer.Flush(timestampUs) after WriteFrame with the same broadcast-clock PTS. It measures catalog jitter at the transport handoff. File, pipe, and network imports should omit Flush; built-in encoders observe their own output.
Call media.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 video track, resume with a keyframe: a delta frame before it fails.
The three advertising operations: client.CreateBroadcast(path) (or origin.CreateBroadcast) 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 OriginDynamic while the claim should stay advertised, and reject the requests you will not serve. A route is a capability, not an inventory. Announced(options) combines a literal prefix with an optional relative pattern and yields an AnnounceEvent: AnnounceEventStart, AnnounceEventUpdate, or AnnounceEventEnd carrying an Announce, whose Prefix stays relative to the origin and whose Captures reports the wildcard matches, 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.
An OriginProducer from moq.NewOriginProducer has no Close: its origin ends when the garbage collector reaches the last owner (each producer, published broadcast, and OriginDynamic), and every consumer made from it then fails with moq.ErrClosed. Keep an owner reachable (a field on a long-lived struct, or runtime.KeepAlive) for as long as the origin should serve.
Every call that can block takes a context.Context first. Cancelling it returns ctx.Err() promptly and tears the in-flight native work down, so a per-call deadline bounds resource use rather than just your wait. What it tears down depends on the call: a one-shot (SubscribeTrack, FetchGroup, RequestBroadcast, Server.Accept, ...) aborts alone and leaves its object usable, while a stream read (Next, RecvGroup, ReadFrame, and the iter.Seq2 iterators over them) cancels the stream it reads, which is what a range loop over a cancelled context wants.
Sessions reconnect with backoff when the transport drops and re-announce local broadcasts, so a worker rides out a relay restart. Session().Epoch() counts the connections, 1 on the first, pairing with Session().Status(ctx) to log each reconnect by number; moq.WithBackoff tunes the pacing, with moq.RetryForever as the timeout; and moq.WithQUICMaxStreams raises the peer's inbound stream cap for a subscriber to many tracks.
The WebSocket fallback races QUIC after a 200 ms head start. moq.WithWebSocketEnabled(false) turns it off for a QUIC-only relay, and moq.WithWebSocketDelay changes the head start.
moq.Listen accepts sessions with per-request Accept/Reject; Request.Transport() returns the closed moq.Transport enum. Request.SetPublish/SetConsume return an error if the request is already answered, cancelled, or currently accepting; ErrBusy is the race with an in-flight Accept. JSON tracks take anything encoding/json handles and return json.RawMessage. The rest of the shared feature list maps one to one: FetchGroup/FetchMediaGroup, Dynamic() with Requests(ctx), Session.Bandwidth() to divide the send estimate, AppendDatagram/Datagrams(ctx), SetCatalogSection, Demand() for Used/Unused, Session().Stats(). moq.IsAuthError and moq.IsShutdown classify errors. moq.ProtocolError(err) is the structured protocol failure (scope, verbatim code, kind) when the peer sent one. err.Error() is the Rust error message.
EncodeAudio encodes raw PCM inside the binding. Its codec is OpusAudioCodec() or AacAudioCodec(), and AudioEncoderOutput.FrameDurationUs 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 VideoDecodedFrame from DecodeVideo owns its decoded picture until Close, including after the consumer is cancelled. Pixels(format) converts it on demand: VideoPixelFormatI420, or VideoPixelFormatRgba for four bytes a pixel. Close 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: a VideoSurfacePixelBuffer whose Pointer is the CVPixelBufferRef, valid until Close. Only macOS has one, so DecodeVideo fails with ErrUnsupported elsewhere.
Connection stats
Session().Stats() returns a ConnectionStats snapshot. Each field is a pointer, nil 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 |
|---|---|---|
RttUs | microseconds | Smoothed round-trip time. |
EstimatedSendRateBps | bits per second | Send bandwidth from the congestion controller. |
EstimatedRecvRateBps | bits per second | Receive bandwidth from MoQ PROBE. |
BytesSent | bytes | Total sent, including retransmissions and overhead. |
BytesReceived | bytes | Total received, including duplicates and overhead. |
BytesLost | bytes | Total lost, detected via retransmission or acknowledgement. |
PacketsSent | datagrams | Total datagrams sent. |
PacketsReceived | datagrams | Total datagrams received. |
PacketsLost | datagrams | Total datagrams detected as lost. |
- API reference: pkg.go.dev/moq.dev/moq
- Source:
go/;just go checkbuilds and tests locally - Mirrors the vanity path resolves to: moq-dev/moq-go (wrapper), moq-dev/moq-go-ffi (raw bindings and static libraries)
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.
session.Shutdown(ctx) drains finished tracks and returns a delivery error if the one-second deadline expires. Cancelling the context aborts immediately. client.Close() waits for shutdown and returns the same error; session.Cancel(code) remains immediate. Finish or abort live tracks before shutdown. IETF media streams are not drained yet.