MoQ is under active development. APIs will change, but we keep backwards wire compatibility.

Skip to content

Go ​

Go Reference

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.

bash
go get moq.dev/moq@latest
go
import "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)
}
go
// 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 explicitly

For 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.

FieldUnitMeaning
RttUsmicrosecondsSmoothed round-trip time.
EstimatedSendRateBpsbits per secondSend bandwidth from the congestion controller.
EstimatedRecvRateBpsbits per secondReceive bandwidth from MoQ PROBE.
BytesSentbytesTotal sent, including retransmissions and overhead.
BytesReceivedbytesTotal received, including duplicates and overhead.
BytesLostbytesTotal lost, detected via retransmission or acknowledgement.
PacketsSentdatagramsTotal datagrams sent.
PacketsReceiveddatagramsTotal datagrams received.
PacketsLostdatagramsTotal datagrams detected as lost.

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.

Licensed under MIT or Apache-2.0