Skip to content
73 changes: 33 additions & 40 deletions STATUS.md

Large diffs are not rendered by default.

16 changes: 10 additions & 6 deletions pkg/moqt/message/datagram.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,14 +36,18 @@ type ObjectDatagram struct {
ObjectPayload []byte // Present when STATUS bit is 0
}

// IsValidDatagramType checks if a datagram type value is valid per §11.3.1
// Figure 24: 0x00..0x0F / 0x20..0x21 / 0x24..0x25 / 0x28..0x29 / 0x2C..0x2D.
// IsValidDatagramType checks if a datagram Type Flags value is valid per
// §11.3.1. The only bits with a specified meaning are the five Datagram*Bit
// flags, so the valid values are 0x00..0x0F / 0x20..0x21 / 0x24..0x25 /
// 0x28..0x29 / 0x2C..0x2D.
//
// The two invalid classes MUST cause a session PROTOCOL_VIOLATION:
// §11.3.1 lists the invalid values, which MUST close the session with a
// PROTOCOL_VIOLATION:
//
// - values outside the 0b00X0XXXX form (i.e. not 0x00..0x0F / 0x20..0x2F);
// - STATUS+END_OF_GROUP (0x22,0x23,0x26,0x27,0x2A,0x2B,0x2E,0x2F) — "an
// object status message cannot signal end of group".
// - bit 4 (0x10) set, or any bit set whose meaning is not specified (i.e.
// not 0x00..0x0F / 0x20..0x2F);
// - both STATUS and END_OF_GROUP set (0x22,0x23,0x26,0x27,0x2A,0x2B,0x2E,
// 0x2F).
//
// Note STATUS+PROPERTIES (0x21,0x25,0x29,0x2D) IS a valid type: it only
// becomes an error when the Object Status is not Normal (0x0) — a per-value
Expand Down
9 changes: 5 additions & 4 deletions pkg/moqt/message/grease.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,10 @@ import (
// GREASE values follow the pattern 0x7F * N + 0x9D for non-negative integer
// values of N (that is, 0x9D, 0x11C, 0x19B, ..., 0x3FFFFFFFFFFFFFDE).
//
// Implementations SHOULD send GREASE values in extensible fields to exercise
// recipient tolerance. Recipients MUST ignore unknown values and MUST NOT
// close the session solely because they received one.
// §14 reserves GREASE values in the Setup Options, Properties, error-code, and
// Auth Token Type registries: implementations "MUST handle unknown values
// gracefully", and endpoints "MUST NOT close the session solely because they
// received an unknown value".

// greaseBase and greaseStep define the GREASE value pattern: base + step*N.
const (
Expand All @@ -31,7 +32,7 @@ const maxGreaseN uint64 = (0x3FFFFFFFFFFFFFFF - greaseBase) / greaseStep
// returned value is suitable for use as a Setup Option type, Property type,
// or error code. Each call returns a fresh random value.
func GreaseValue() uint64 {
//nolint:gosec // G404: GREASE values are deliberately non-cryptographic (§1.4.3); randomness only spreads coverage.
//nolint:gosec // G404: GREASE values are deliberately non-cryptographic; randomness only spreads coverage.
n := rand.Uint64N(maxGreaseN + 1)
return greaseBase + greaseStep*n
}
Expand Down
7 changes: 4 additions & 3 deletions pkg/moqt/message/subgroup.go
Original file line number Diff line number Diff line change
Expand Up @@ -155,9 +155,10 @@ func IsSubgroupHeaderType(t uint64) bool {

// IsReservedSubgroupHeaderType reports whether t looks like a SUBGROUP_HEADER
// type byte (bit 4 set, bit 7 clear) but has the reserved SUBGROUP_ID_MODE
// value 0b11 in bits 1-2. Per §11.4.2, receiving such a value MUST be treated
// as a session-level PROTOCOL_VIOLATION — unlike a truly unknown stream type,
// which may be ignorable (GREASE).
// value 0b11 in bits 1-2. Per §11.4.2, receiving such a value MUST close the
// session with a PROTOCOL_VIOLATION. §3.4 also requires closing the session on
// a truly unknown stream type; this lets the caller report the reserved mode
// distinctly from an unknown type.
func IsReservedSubgroupHeaderType(t uint64) bool {
if t > 0x7F {
return false
Expand Down
4 changes: 2 additions & 2 deletions pkg/moqt/message/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -125,8 +125,8 @@ type validator interface {
}

// ErrUnknownType is returned for a message type not implemented by this
// package. Per §10 the receiver MUST close the session with
// PROTOCOL_VIOLATION; callers translate accordingly.
// package. Per §10 the receiver MUST close the session; §10 names no error
// code for this, and callers close with PROTOCOL_VIOLATION.
type ErrUnknownType Type

func (e ErrUnknownType) Error() string {
Expand Down
4 changes: 2 additions & 2 deletions pkg/moqt/session/datagram.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,8 @@ const paddingDatagramType uint64 = 0x132B3E29

// ReceiveDatagram blocks until a QUIC DATAGRAM frame arrives from the peer,
// parses it, and returns the contained ObjectDatagram. PADDING datagrams
// (§11.3) are silently consumed and the call retries. Unknown datagram types
// close the session with PROTOCOL_VIOLATION per §11.
// (§11.5.2) are silently consumed and the call retries. Unknown datagram types
// close the session (§11) with PROTOCOL_VIOLATION.
//
// Transport-level errors (session closed, ctx cancelled) are returned
// unwrapped so the caller can distinguish them from parse failures.
Expand Down
7 changes: 5 additions & 2 deletions pkg/moqt/session/datastream_in.go
Original file line number Diff line number Diff line change
Expand Up @@ -242,8 +242,11 @@ type IncomingFetchStream struct {
// newGroup = prevGroup + delta + 1; descending → newGroup =
// prevGroup - delta - 1. The §11.4.4 wire format does not encode
// the direction; the caller knows it from the GROUP_ORDER
// parameter it sent in FETCH (or from the publisher default).
// Defaults to ascending when unset (zero value).
// parameter it sent in FETCH (§10.2.8: Ascending when omitted) or,
// for a fill fetch stream, from FILL_PARAMETERS, else the
// subscription's group order (§10.2.15), which defaults to the
// Track's publisher preference (§10.2.8). Defaults to ascending when
// unset (zero value).
GroupOrder message.GroupOrder

// Decoder state used by ReadDecoded — running absolute values
Expand Down
31 changes: 20 additions & 11 deletions pkg/moqt/session/namespace.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,11 +10,13 @@ import (

// NamespacePublication is an established PUBLISH_NAMESPACE request (§10.16). It
// embeds the still-open request stream (so Close / writes / message.Marshal work
// directly on it) and carries the peer's REQUEST_OK. The caller announces tracks
// by writing NAMESPACE / NAMESPACE_DONE follow-ups to the embedded stream.
// directly on it) and carries the peer's REQUEST_OK. The namespace stays
// published until the request is cancelled ([NamespacePublication.Close],
// §6.2, §3.3.3); a FIN does not withdraw it (§3.3.2). NAMESPACE and
// NAMESPACE_DONE answer a SUBSCRIBE_NAMESPACE instead (§10.17, §10.18).
type NamespacePublication struct {
// Stream is the PUBLISH_NAMESPACE request stream, still open for
// NAMESPACE / NAMESPACE_DONE follow-ups. [NamespacePublication.Close]
// REQUEST_UPDATE follow-ups (§10.9). [NamespacePublication.Close]
// withdraws the publication.
Stream

Expand Down Expand Up @@ -85,8 +87,8 @@ func (t *TrackSubscription) Update(ctx context.Context, params message.Parameter
// caller supplies Namespace and optional Parameters.
//
// On success a [NamespacePublication] is returned whose embedded stream stays
// open (the caller may send NAMESPACE / NAMESPACE_DONE messages on it). On
// REQUEST_ERROR the stream is closed and a *RequestRejectedError is returned.
// open (the caller may send REQUEST_UPDATE on it, §10.9). On REQUEST_ERROR the
// stream is closed and a *RequestRejectedError is returned.
func (s *Session) PublishNamespace(
ctx context.Context,
m *message.PublishNamespace,
Expand Down Expand Up @@ -140,12 +142,19 @@ func (s *Session) SubscribeTracks(ctx context.Context, m *message.SubscribeTrack
// IncomingNamespacePublication is an accepted inbound PUBLISH_NAMESPACE (§10.16)
// — the receiving side of [Session.PublishNamespace]'s [NamespacePublication],
// returned by [Request.AcceptPublishNamespace]. REQUEST_OK has been sent; the
// announcer's follow-ups arrive on the embedded stream; read it with
// [RequestBroker.Serve], which enforces the session-level rules (§10,
// §10.2.1). Close it to end the publication.
// announcer's follow-ups (REQUEST_UPDATE, GOAWAY) arrive on the embedded
// stream. Read it with a [Session.NewRequestBroker]'s [RequestBroker.Serve],
// which closes the session on a malformed message (§10) or a GOAWAY violation
// (§10.4). Before Serve, call [RequestBroker.PeerMessages](true, false), since
// PUBLISH_STATE_NOTIFY applies only to subscriptions (§10.10), and
// [RequestBroker.UpdateScope](message.ScopeUpdatePublishNamespace), so an
// update's parameters are checked (§10.2.1) before it is answered; without
// [RequestBroker.HandleUpdates] each REQUEST_UPDATE is declined with
// NOT_SUPPORTED. Cancel the request (CancelRead and CancelWrite, §3.3.3) to
// revoke acceptance (§6.2); Close only FINs this side (§3.3.2).
type IncomingNamespacePublication struct {
// Stream is the PUBLISH_NAMESPACE request stream, still open to receive
// NAMESPACE / NAMESPACE_DONE notifications. Close it to end the publication.
// the announcer's follow-ups.
Stream
}

Expand Down Expand Up @@ -200,8 +209,8 @@ func acceptNamespaceRequest[M message.Message, T any](r *Request, op string, wra

// AcceptPublishNamespace accepts an inbound PUBLISH_NAMESPACE (§10.16), replies
// REQUEST_OK, and returns an [IncomingNamespacePublication] for receiving the
// announcer's NAMESPACE / NAMESPACE_DONE follow-ups — the accept-side
// counterpart of [Session.PublishNamespace]. r.First MUST be a
// announcer's follow-ups — the accept-side counterpart of
// [Session.PublishNamespace]. r.First MUST be a
// *message.PublishNamespace.
func (r *Request) AcceptPublishNamespace() (*IncomingNamespacePublication, error) {
return acceptNamespaceRequest[*message.PublishNamespace](r, "AcceptPublishNamespace",
Expand Down
3 changes: 2 additions & 1 deletion pkg/moqt/session/options.go
Original file line number Diff line number Diff line change
Expand Up @@ -152,7 +152,8 @@ func WithTokenVerifier(v TokenVerifier) Option {
// into the outbound SETUP message to exercise the peer's tolerance of unknown
// values. GREASE values follow the pattern 0x7F * N + 0x9D and are always
// larger than all currently defined SETUP option types, so appending preserves
// the ascending-Type ordering required by §1.4.3.
// the non-decreasing Type order that §1.4.3's Delta Type encoding requires
// (Setup Options are Key-Value-Pairs, not §10.2 Message Parameters).
func WithGrease() Option {
return func(c *config) {
c.setupOptions = append(c.setupOptions, message.GreaseSetupOption())
Expand Down
2 changes: 1 addition & 1 deletion pkg/relay/cache/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -333,7 +333,7 @@ func (c *ObjectCache) Len() int {
// - [message.GroupOrderDescending]: groups desc, objects asc within group.
//
// Within a group the inner order is always ascending by Object ID, matching
// §11.4.3's subgroup-stream constraint.
// §10.13: "Within each group, objects are sent in Object ID order".
//
// An empty or inverted range (end < start) returns nil.
//
Expand Down
13 changes: 12 additions & 1 deletion pkg/relay/export_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,10 @@
package relay

import "github.com/floatdrop/moq-go/pkg/moqt/track"
import (
"time"

"github.com/floatdrop/moq-go/pkg/moqt/track"
)

// SetTestHookAfterAliasRegistered installs hook, to be called at the moment a Track Alias
// becomes routable on the SUBSCRIBE and PUBLISH paths, and returns a function restoring the previous value. See
Expand Down Expand Up @@ -48,3 +52,10 @@ func SetTestHookBeforeForwardClaim(hook func(track.FullTrackName)) (restore func
testHookBeforeForwardClaim.Store(&hook)
return func() { testHookBeforeForwardClaim.Store(prev) }
}

// SetTrackStatusTimeout bounds a forwarded TRACK_STATUS's upstream round trip
// at d, and returns the restore.
func SetTrackStatusTimeout(d time.Duration) (restore func()) {
prev := trackStatusTimeout.Swap(int64(d))
return func() { trackStatusTimeout.Store(prev) }
}
3 changes: 2 additions & 1 deletion pkg/relay/handler_datagram.go
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,8 @@ func (h *sessionHandler) handleDatagram(ctx context.Context, d *message.ObjectDa

downstream := entry.CopyDownstream()
for _, sub := range downstream {
// §5.1.4: a datagram counts as subgroup 0.
// A datagram belongs to no Subgroup (§2.2), and §5.1.4 does not say
// how a SUBGROUP_FILTER treats one; the relay filters it as Subgroup 0.
if sub.ForwardDecision(d.GroupID, d.ObjectID, 0, d.PublisherPriority, d.Properties) != registry.Forward {
continue
}
Expand Down
3 changes: 2 additions & 1 deletion pkg/relay/handler_fanout.go
Original file line number Diff line number Diff line change
Expand Up @@ -374,7 +374,8 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming

if created {
// Under sg.Mu so a concurrent contributor's joiner scan can't
// double-open. The stream is drained even with no subscribers (§9.7).
// double-open. The stream is drained even with no subscribers, so its
// Objects still reach the cache (§9.1) and any later joiner.
initialSubs, gen := entry.CopyDownstreamWithGen()
pubTimeouts := entry.DeliveryTimeouts()
lowest, forwarded := entry.LowestForwarded(hdr.GroupID, hdr.SubgroupID)
Expand Down
4 changes: 2 additions & 2 deletions pkg/relay/handler_fetch.go
Original file line number Diff line number Diff line change
Expand Up @@ -308,7 +308,7 @@ func (h *sessionHandler) stitchedFetchObjects(
return fetchElements(cached, unknown, nil, order), nil
}
span := registry.LocRange{Lo: unknown[0].Lo, Hi: unknown[len(unknown)-1].Hi}
done := h.tracks.BeginFetch(up.Session, fullName.Key())
done := h.tracks.BeginRequest(message.TypeFetch, up.Session, fullName.Key())
ans, refusal := h.fetchUpstreamRange(ctx, up, fullName, span, order, fillTimeout)
done()
if errors.Is(refusal, session.ErrMalformedTrack) {
Expand Down Expand Up @@ -379,7 +379,7 @@ func (h *sessionHandler) pickFetchUpstream(entry *registry.TrackEntry) *registry
if !u.FetchCapable || !u.IsEstablished() || u.Session == nil || goingAway(u.Session) {
continue
}
if u.Session == h.sess && h.tracks.FetchPending(u.Session, entry.FullName.Key()) {
if u.Session == h.sess && h.tracks.RequestPending(message.TypeFetch, u.Session, entry.FullName.Key()) {
continue
}
return u
Expand Down
31 changes: 0 additions & 31 deletions pkg/relay/handler_fetch_session_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,37 +61,6 @@ func TestTrackStatus_ReplyForKnownTrack(t *testing.T) {
}
}

// TestTrackStatus_ReplyEmptyPropertiesForKnownNamespace: a TRACK_STATUS for a
// track under an advertised namespace with no upstream yet gets
// TRACK_STATUS_OK with empty Properties.
func TestTrackStatus_ReplyEmptyPropertiesForKnownNamespace(t *testing.T) {
t.Parallel()
pubSess, teardown := connectRelay(t, relay.Config{})
defer teardown()

pnsStream, err := pubSess.PublishNamespace(t.Context(), &message.PublishNamespace{
Namespace: ns("video"),
})
if err != nil {
t.Fatalf("PublishNamespace: %v", err)
}
defer pnsStream.Close()

querySess := dialAnotherClient(t, pubSess)
tsStream, err := querySess.TrackStatus(t.Context(), &message.TrackStatus{
Namespace: ns("video"),
Name: []byte("cam-anything"),
})
if err != nil {
t.Fatalf("TrackStatus: %v", err)
}
defer tsStream.Close()

if len(tsStream.OK.TrackProperties) != 0 {
t.Fatalf("TrackProperties = %q, want empty", tsStream.OK.TrackProperties)
}
}

// TestTrackStatus_RejectsUnknownTrack pins the no-publisher-no-namespace
// case: TRACK_STATUS for a name no one has claimed returns
// RequestDoesNotExist.
Expand Down
3 changes: 2 additions & 1 deletion pkg/relay/handler_forward.go
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,8 @@ func (h *sessionHandler) serveForwardedPublish(
if _, err := h.sess.AwaitPublishOK(ctx, stream); err != nil {
h.log.LogAttrs(ctx, slog.LevelDebug, "forwarded PUBLISH refused",
slog.String("name", string(fullName.Name)), slog.String("err", err.Error()))
// §3.3.3: no PUBLISH_DONE after the subscriber's REQUEST_ERROR.
// §5.1, §5.1.1: the subscriber's REQUEST_ERROR terminates the
// subscription, so no PUBLISH_DONE follows.
sub.EndRefused()
return
}
Expand Down
Loading
Loading