diff --git a/STATUS.md b/STATUS.md index a5abae7b..1758082f 100644 --- a/STATUS.md +++ b/STATUS.md @@ -34,10 +34,10 @@ By package, bottom-up along the dependency stack: streaming `Decoder` over one control-frame interface. - **`message`** — typed control, request-stream, and data-stream messages with parameter negotiation: SETUP, GOAWAY, SUBSCRIBE, PUBLISH (+DONE/SKIPPED), - FETCH (standalone + relative/absolute joining), TRACK_STATUS, REQUEST_UPDATE, - the namespace messages, §11 object framing (subgroup/fetch/datagram), - location filters, GREASE, and a parse-time `Validate` hook that rejects - structurally-malformed messages. + FETCH (draft-20 replaced Joining FETCH with fill fetch streams), + TRACK_STATUS, REQUEST_UPDATE, the namespace messages, §11 object framing + (subgroup/fetch/datagram), location filters, GREASE, and a parse-time + `Validate` hook that rejects structurally-malformed messages. - **`session`** — the SETUP handshake with version negotiation, control multiplexing and request-ID allocation, §3.5 Track-Alias management with collision detection, the request openers (`Publish`/`Subscribe`/`Fetch`/…) and @@ -164,7 +164,7 @@ By package, bottom-up along the dependency stack: | 10.2.3 | SUBGROUP_DELIVERY_TIMEOUT | 0x06 | PARTIAL| Parsed and resolved; the stream reset is not enforced on the bundled transports, and datagrams are not dropped (see §8). | | 10.2.4 | OBJECT_DELIVERY_TIMEOUT | 0x02 | DONE | | | 10.2.5 | FILL_TIMEOUT | 0x0A | DONE | The budget for a FETCH's or fill's upstream FETCH, its response included: when it runs out, what arrived is served and the rest is an End of Timed-Out Range; 0 asks no upstream. Default 5s. | -| 10.2.6 | RENDEZVOUS_TIMEOUT | 0x04 | DONE | The relay holds a SUBSCRIBE with no publisher, capped by `Config.MaxRendezvousTimeout` (30s), and answers TIMEOUT when the hold runs out. No publisher means none matched, or each answered DOES_NOT_EXIST or TIMEOUT, or is draining; any other refusal ends the hold with its error. A PUBLISH of the track, or a covering namespace newly published here or advertised through Discovery, wakes the hold; a publisher that already answered is not asked again. Upstream SUBSCRIBEs to other relays carry what is left of the hold, and are cut short when a publisher arrives here. | +| 10.2.6 | RENDEZVOUS_TIMEOUT | 0x04 | DONE | The relay holds a SUBSCRIBE with no publisher, capped by `Config.MaxRendezvousTimeout` (30s), and answers TIMEOUT when the hold runs out. No publisher means none matched, or each answered DOES_NOT_EXIST, TIMEOUT or GOING_AWAY, failed at the transport, or is draining; any other refusal ends the hold with its error. The refusal a subscriber gets does not depend on the order candidates answer in: the highest ranked, and of equal rank the soonest retry. A PUBLISH of the track, or a covering namespace newly published here or advertised through Discovery, wakes the hold; a publisher that already answered is not asked again. Upstream SUBSCRIBEs to other relays carry what is left of the hold, and are cut short when a publisher arrives here. | | 10.2.7 | SUBSCRIBER_PRIORITY | 0x20 | DONE | | | 10.2.8 | GROUP_ORDER | 0x22 | DONE | A value outside {1, 2} closes the session, wherever it appears, FILL_PARAMETERS included. | | 10.2.9 | LOCATION_FILTER | 0x21 | DONE | An end Group overflowing 2^64-1 closes the session with PROTOCOL_VIOLATION (§5.1.2); a value that does not parse, with KEY_VALUE_FORMATTING_ERROR (§1.4.3). | @@ -173,13 +173,13 @@ By package, bottom-up along the dependency stack: | 10.2.12 | PRIORITY_FILTER | 0x27 | DONE | Enforced per object (subgroup priority); >255 rejected INVALID_FILTER. | | 10.2.13 | OBJECT_PROPERTY_FILTER | 0x28 | DONE | Enforced per object against Object Properties; even property type. | | 10.2.14 | TRACK_PROPERTY_FILTER | 0x29 | DONE | Gates PUBLISH forwarding on SUBSCRIBE_TRACKS against Track Properties; even property type. | -| 10.2.15 | FILL_PARAMETERS | 0x23 | PARTIAL| Inner Table 6 scope and duplicates checked; omitted Range Filters are not inherited from the subscription (see §5.1.3). | +| 10.2.15 | FILL_PARAMETERS | 0x23 | DONE | Inner Table 6 scope and duplicates checked; the fill inherits the subscription's Range Filters, those inside overriding per type (§5.1.3). SUBSCRIBER_PRIORITY inside it has no effect: fill streams carry no priority input (§7.2). | | 10.2.16 | EXPIRES | 0x08 | DONE | | | 10.2.17 | LARGEST_OBJECT | 0x09 | DONE | Monotonic constraint applied. | | 10.2.18 | FORWARD | 0x10 | DONE | A value above 1 closes the session, in every message that may carry it. | | 10.2.19 | NEW_GROUP_REQUEST | 0x32 | DONE | | | 10.2.20 | TRACK_NAMESPACE_PREFIX | 0x34 | DONE | Applied on REQUEST_UPDATE; SUBSCRIBE_NAMESPACE reconciles its announced set. | -| 10.2.21 | INCLUDE_PROPERTIES | 0x35 | DONE | A value other than 0 or 1 closes the session. With 0 the relay sends empty Track Properties in SUBSCRIBE_OK, FETCH_OK, TRACK_STATUS_OK and forwarded PUBLISH, and writes the priority inline on that subscription's subgroups and datagrams, since the subscriber cannot inherit DEFAULT_PUBLISHER_PRIORITY. | +| 10.2.21 | INCLUDE_PROPERTIES | 0x35 | DONE | A value other than 0 or 1 closes the session. With 0 the relay sends empty Track Properties in SUBSCRIBE_OK, FETCH_OK, TRACK_STATUS_OK and forwarded PUBLISH, and writes the priority inline on that subscription's subgroups and datagrams, since the subscriber cannot inherit DEFAULT_PUBLISHER_PRIORITY. Nor can the subscriber learn a fill's Group Order when its SUBSCRIBE omitted GROUP_ORDER, a draft gap (see Limitations). | | 10.3 | SETUP | 0x2F00 | DONE | Bidirectional handshake; options as KV pairs. | | 10.3.1.1| AUTHORITY option | 0x05 | PARTIAL| Sent (`WithAuthority`); refused from a server or over WebTransport (INVALID_AUTHORITY) and when not RFC 3986 syntax (MALFORMED_AUTHORITY, `uri.CheckAuthority`). Whether the server serves it is not checked — see Limitations. | | 10.3.1.2| PATH option | 0x01 | PARTIAL| Sent (`WithPath`); refused from a server or over WebTransport (INVALID_PATH) and when not RFC 3986 syntax (MALFORMED_PATH, `uri.CheckPathAndQuery`). Whether the server serves it is not checked — see Limitations. | @@ -199,7 +199,7 @@ By package, bottom-up along the dependency stack: | 10.12 | PUBLISH_DONE | 0x0B | DONE | Sent once every stream of the subscription has closed and no datagram send is in progress, with the exact Stream Count; written on its own goroutine, so subscribers do not wait on each other. When a track's last upstream ends, its PUBLISH_DONE code reaches subscribers if it is about the track (TRACK_ENDED, MALFORMED_TRACK); codes about the relay's own upstream subscription become INTERNAL_ERROR. Session `Publication.Done` resets the subgroups still open with CANCELLED, refuses later opens and writes (`ErrPublicationEnded`), and counts every subgroup opened, however the opens race it. | | 10.13 | FETCH | 0x16 | DONE | Standalone, the only kind in draft-20. From the cache, a Location is non-existent only on a signal: a Prior Group or Object ID Gap, a Group's or the Track's end, or an upstream's FETCH. Other uncached Locations are FETCHed from a fetch-capable upstream in one span, within FILL_TIMEOUT, or else marked End of Unknown (or Timed-Out) Range. | | 10.14 | FETCH_OK | 0x18 | DONE | The relay sets End Of Track when the End Location is the Object an END_OF_TRACK status made the Track's final one; it does not learn a Track's end from an upstream FETCH_OK's End Of Track. An End Location before the FETCH's Start closes the session. A Start relative to the Largest Object is compared through End ≤ Largest; an End of {0,0} is let through, as it cannot be told apart from "no content yet". | -| 10.15 | TRACK_STATUS | 0x0D | DONE | Reply via REQUEST_OK, then FIN; any follow-up from the requester closes the session. | +| 10.15 | TRACK_STATUS | 0x0D | DONE | Reply via REQUEST_OK, then FIN; any follow-up from the requester closes the session. The relay answers from an Established subscription, else forwards TRACK_STATUS to every candidate SUBSCRIBE would try, concurrently, and combines the answers as a SUBSCRIBE_OK would (the entry's Track Properties, else the first answer's; the largest LARGEST_OBJECT), or gives the refusal SUBSCRIBE would. A forwarding TRACK_STATUS counts against `Config.MaxSubscriptionsPerSession` (§13.1), and each upstream round trip is bounded at 5s, counting as TIMEOUT (§13.6). Relay policy: not to the requester's own session while a TRACK_STATUS for the track to it is in flight (a loop, §6.2). | | 10.16 | PUBLISH_NAMESPACE | 0x06 | DONE | | | 10.17 | NAMESPACE | 0x08 | DONE | Per namespace, counted over local and remote sources. | | 10.18 | NAMESPACE_DONE | 0x0E | DONE | Never before its NAMESPACE. | @@ -354,8 +354,9 @@ Known protocol gaps, roughly ordered by how load-bearing they are: adapters absorb the knob and quic-go round-robins instead. A REQUEST_UPDATE that changes priority mid-stream applies only to subsequently opened subgroups. - **LOC encryption / SecureObjects and Private Properties** — intentionally out - of scope pending a chosen SecureObjects revision. Some property IDs are - draft-tentative (e.g. `PropAudioLevel = 0x0A`, pending IANA assignment). + of scope pending a chosen SecureObjects revision. The property IDs are + those draft-ietf-moq-loc-04 requests from IANA (§6.1), e.g. `PropAudioLevel = + 0x0C`, and may change until assigned. - **MSF** — no timeline GZIP compression, content protection (§4.3), token authorization, or logs/analytics. No built-in ABR helper: every catalog field a selector needs is surfaced (AltGroup, Width/Height, Bitrate, RenderGroup, @@ -396,9 +397,13 @@ Known protocol gaps, roughly ordered by how load-bearing they are: - **Handles the application reads itself** — REQUEST_UPDATE / PUBLISH_STATE_NOTIFY roles (§10.9, §10.10) and Message Parameter scope (§10.2.1) are enforced by brokers from typed handles' `Broker()`, by the - session's own reads, and by the relay. Handles the application reads itself - (the namespace handles, `FetchResponder`, or any stream read with - `message.Parse`) are checked only if it calls `Session.CheckPeerParams`. + session's own reads, and by the relay. On handles with no broker of their + own (`NamespacePublication`, `IncomingNamespacePublication`, + `IncomingNamespaceSubscription`, `IncomingTrackSubscription`, + `FetchResponder`) or any stream read with `message.Parse`, the application + checks both: roles itself, parameter scope with `Session.CheckPeerParams`. A + bare `Session.NewRequestBroker` checks them only as configured + (`PeerMessages`, `UpdateScope`). Likewise for §10 framing: such a reader must close the session itself on an error wrapping `message.ErrMalformedMessage`, or read through `Session.NewRequestBroker(stream).Serve`, which does. @@ -424,7 +429,10 @@ Known protocol gaps, roughly ordered by how load-bearing they are: subscriber SHOULD individually unsubscribe from each existing subscription"), nor migrates to the New Session URI, nor closes the session once no subscriptions remain (§3.6 RECOMMENDED). It waits for the sender to - close. + close. A GOAWAY on a request stream is checked (§10.4) and otherwise + ignored: the relay neither re-issues that request (at the New Session URI, + or on this session when none is given) nor closes the old stream, which the + recipient SHOULD do. - **Malformed tracks (§2.4.2, §9.1, §12.8, §12.9)** — the session reports Object Properties that make a track malformed (`session.ErrMalformedTrack`), and the relay then ends the track: PUBLISH_DONE MALFORMED_TRACK to every @@ -489,6 +497,10 @@ Known protocol gaps, roughly ordered by how load-bearing they are: Properties, and neither SUBSCRIBE_OK nor the FETCH_HEADER carries a Group Order, so it cannot tell a Descending fill from an Ascending one. A gap in the draft; a subscriber that asks for a fill avoids it by sending GROUP_ORDER. +- **SUBGROUP_FILTER on datagrams (§5.1.4, §2.2)** — an interpretation. A + Datagram Object "does not belong to a Subgroup in any way" (§2.2), and + §5.1.4 does not say how a Subgroup filter treats one; the relay filters it + as Subgroup 0. - **FETCH End of Range markers in Descending order (§11.4.4.2)** — an interpretation. A marker covers "Locations between the last serialized Object, if any, and this Location"; the relay reads "between" in the order @@ -508,10 +520,6 @@ Known protocol gaps, roughly ordered by how load-bearing they are: Group Order differs, and the fill-delivered one first within a Group. The relay writes each stream as it is fed; ordering across them is a scheduler the relay does not have. -- **Duplicate Objects from redundant upstreams are not compared (§9.1)** — - the first copy of each {Group, Object} is forwarded and later ones are - dropped unread. Comparing them would detect a malformed track (§2.4.2 - condition 6), at a cost on every Object. ### Draft-20 compliance review backlog @@ -523,11 +531,6 @@ Relay: - Objects age from when they were read whole rather than their beginning (§12.3). -- TRACK_STATUS is not answered as SUBSCRIBE would be in two cases (§10.15): - with only the namespace advertised it answers an empty OK where SUBSCRIBE - goes upstream and may fail, and a leftover entry with Track Properties or a - LARGEST_OBJECT but no established upstream answers OK where SUBSCRIBE would - go upstream. - A client's SUBSCRIBE to a track it publishes gets DOES_NOT_EXIST while an upstream SUBSCRIBE for that track to the client is still pending, and its FETCH marks a hole unknown while a stitch FETCH for the track to it is in @@ -535,23 +538,13 @@ Relay: by a relay peer (§6.2 has no loop protection), so it declines the second hop rather than loop. Self-subscriptions are otherwise "identical" (§5.1). - Filters are not aggregated upstream (§6.3.1 SHOULD). -- Among SUBSCRIBE candidates that only say the track has no publisher yet - (DOES_NOT_EXIST, TIMEOUT, draining), the last to answer sets the refusal - code, so a draining local publisher and a remote relay's DOES_NOT_EXIST - yield DOES_NOT_EXIST while the reverse yields GOING_AWAY. An upstream relay - that answers GOING_AWAY itself, before its GOAWAY reaches this relay, is - passed on as INTERNAL_ERROR. - -Documentation: - -- Limitations: "Duplicate Objects … are not compared" is stale; the LOC entry names - `PropAudioLevel = 0x0A` (it is 0x0C); "Handles the application reads itself" - says `CheckPeerParams` checks roles; "Inbound GOAWAY" omits request streams. -- Table rows 10.2.15 and 10.2.21 overstate what is done (see - the items above), and the package summary still lists joining FETCH. -- `session/namespace.go` says NAMESPACE / NAMESPACE_DONE go on a - PUBLISH_NAMESPACE stream (§10.17, §10.18). -- About a dozen stale `§` citations (padding, grease, fetch ordering, caching). +- Concurrent TRACK_STATUS requests for one track with no Established + subscription are each forwarded upstream, not coalesced. +- A SUBSCRIBE the relay cannot open to a candidate for want of bidi-stream + credit is answered DOES_NOT_EXIST with no Retry Interval ("SHOULD NOT be + retried"), and outranks another candidate's GOING_AWAY or TIMEOUT, though + the publisher is there; EXCESSIVE_LOAD with a retry would fit better + (§10.6.2). Open questions for interop: whether an End of Range marker carries an Object Payload Length (Figure 28 vs §11.4.4.2), and whether EXPIRES may appear in diff --git a/pkg/moqt/message/datagram.go b/pkg/moqt/message/datagram.go index d6561951..4a55f7e6 100644 --- a/pkg/moqt/message/datagram.go +++ b/pkg/moqt/message/datagram.go @@ -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 diff --git a/pkg/moqt/message/grease.go b/pkg/moqt/message/grease.go index 1c6b6324..60d520d0 100644 --- a/pkg/moqt/message/grease.go +++ b/pkg/moqt/message/grease.go @@ -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 ( @@ -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 } diff --git a/pkg/moqt/message/subgroup.go b/pkg/moqt/message/subgroup.go index 4af6b7e1..86952f9e 100644 --- a/pkg/moqt/message/subgroup.go +++ b/pkg/moqt/message/subgroup.go @@ -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 diff --git a/pkg/moqt/message/types.go b/pkg/moqt/message/types.go index 2da0b960..0b0306c8 100644 --- a/pkg/moqt/message/types.go +++ b/pkg/moqt/message/types.go @@ -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 { diff --git a/pkg/moqt/session/datagram.go b/pkg/moqt/session/datagram.go index e8dc7262..7b3b45f7 100644 --- a/pkg/moqt/session/datagram.go +++ b/pkg/moqt/session/datagram.go @@ -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. diff --git a/pkg/moqt/session/datastream_in.go b/pkg/moqt/session/datastream_in.go index 60fe4461..fc4405aa 100644 --- a/pkg/moqt/session/datastream_in.go +++ b/pkg/moqt/session/datastream_in.go @@ -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 diff --git a/pkg/moqt/session/namespace.go b/pkg/moqt/session/namespace.go index 702c2b72..4278472b 100644 --- a/pkg/moqt/session/namespace.go +++ b/pkg/moqt/session/namespace.go @@ -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 @@ -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, @@ -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 } @@ -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", diff --git a/pkg/moqt/session/options.go b/pkg/moqt/session/options.go index 150f9776..a216b308 100644 --- a/pkg/moqt/session/options.go +++ b/pkg/moqt/session/options.go @@ -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()) diff --git a/pkg/relay/cache/cache.go b/pkg/relay/cache/cache.go index e2975f34..b6bb2133 100644 --- a/pkg/relay/cache/cache.go +++ b/pkg/relay/cache/cache.go @@ -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. // diff --git a/pkg/relay/export_test.go b/pkg/relay/export_test.go index 0d978d84..478a5561 100644 --- a/pkg/relay/export_test.go +++ b/pkg/relay/export_test.go @@ -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 @@ -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) } +} diff --git a/pkg/relay/handler_datagram.go b/pkg/relay/handler_datagram.go index db08c939..e4b0da53 100644 --- a/pkg/relay/handler_datagram.go +++ b/pkg/relay/handler_datagram.go @@ -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 } diff --git a/pkg/relay/handler_fanout.go b/pkg/relay/handler_fanout.go index 96068c51..98126e39 100644 --- a/pkg/relay/handler_fanout.go +++ b/pkg/relay/handler_fanout.go @@ -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) diff --git a/pkg/relay/handler_fetch.go b/pkg/relay/handler_fetch.go index cd92c2e2..5cd9b329 100644 --- a/pkg/relay/handler_fetch.go +++ b/pkg/relay/handler_fetch.go @@ -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) { @@ -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 diff --git a/pkg/relay/handler_fetch_session_test.go b/pkg/relay/handler_fetch_session_test.go index b97c6a5f..9b052498 100644 --- a/pkg/relay/handler_fetch_session_test.go +++ b/pkg/relay/handler_fetch_session_test.go @@ -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. diff --git a/pkg/relay/handler_forward.go b/pkg/relay/handler_forward.go index 70f84ff2..488bc4c2 100644 --- a/pkg/relay/handler_forward.go +++ b/pkg/relay/handler_forward.go @@ -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 } diff --git a/pkg/relay/handler_subscribe.go b/pkg/relay/handler_subscribe.go index 644f2a69..a22c3452 100644 --- a/pkg/relay/handler_subscribe.go +++ b/pkg/relay/handler_subscribe.go @@ -513,15 +513,111 @@ func (l *holdLook) end() { // awaitsPublisher reports whether err, from [sessionHandler.subscribeUpstream], // leaves the track without a current publisher: none matched (nil), or every -// one that did answered DOES_NOT_EXIST, or TIMEOUT for an upstream relay's own -// hold, or is draining (§10.4), since subscribeUpstream reports such an error -// only when no other kind occurred. +// one that did answered DOES_NOT_EXIST, TIMEOUT for an upstream relay's own +// hold, or GOING_AWAY for its own drain, failed at the transport without an +// answer (see [transportFailure]), or is draining (§10.4). subscribeUpstream +// reports such an error only when no other kind occurred. func awaitsPublisher(err error) bool { - if err == nil || errors.Is(err, errGoingAway) { + if err == nil || errors.Is(err, errGoingAway) || transportFailure(err) { return true } rej, ok := errors.AsType[*session.RequestRejectedError](err) - return ok && (rej.Code == moqt.RequestDoesNotExist || rej.Code == moqt.RequestTimeout) + return ok && (rej.Code == moqt.RequestDoesNotExist || rej.Code == moqt.RequestTimeout || + rej.Code == moqt.RequestGoingAway) +} + +// transportFailure reports whether err is a candidate's request failing at +// the transport without any answer about the track: a stream reset, a FIN, +// or the session ending. Not a request the relay could not open for want of +// stream credit (session.ErrNoStreamCredit), which says the publisher is +// there, only busy. +func transportFailure(err error) bool { + if errors.Is(err, errGoingAway) || isTrackPropertiesErr(err) || errors.Is(err, session.ErrNoStreamCredit) { + return false + } + _, rejected := errors.AsType[*session.RequestRejectedError](err) + return !rejected +} + +// preferCandidateErr returns whichever of err and last, two candidates' +// errors, the subscriber is answered with (see [upstreamRejection]), so the +// answer does not depend on the order candidates answer in: the higher +// ranked ([candidateErrRank]), and of equal rank the one allowing the soonest +// retry, "SHOULD NOT be retried" (Retry Interval 0, §10.6.2) only when both +// say so. +func preferCandidateErr(last, err error) error { + if r, l := candidateErrRank(err), candidateErrRank(last); r != l { + if r > l { + return err + } + return last + } + if a, b := candidateRetry(err), candidateRetry(last); a != b { + if a != 0 && (b == 0 || a < b) { + return err + } + return last + } + // Of equal rank and retry, the answers differ only among "any other + // refusal" (rank 4): prefer a specific code to INTERNAL_ERROR, then the + // lower code. + a, b := upstreamRejection(err).Code, upstreamRejection(last).Code + if a != b && (b == moqt.RequestInternalError || (a != moqt.RequestInternalError && a < b)) { + return err + } + return last +} + +// candidateRetry is the Retry Interval err is answered with: for a +// GOING_AWAY without one, the most the relay's own jittered one can be (see +// [upstreamRejection]), else the upstream's. +func candidateRetry(err error) uint64 { + if errors.Is(err, errGoingAway) { + return goingAwayRetry + } + rej, ok := errors.AsType[*session.RequestRejectedError](err) + if !ok { + return 0 + } + if rej.Code == moqt.RequestGoingAway && rej.RetryInterval == 0 { + return goingAwayRetry + } + return rej.RetryInterval +} + +// candidateErrRank orders candidates' errors for [preferCandidateErr]: an +// unknown Mandatory Track Property (§2.5.1: UNSUPPORTED_EXTENSION, a MUST), +// then Track Properties that do not parse, then any other refusal, which ends +// a RENDEZVOUS_TIMEOUT hold (see awaitsPublisher), then the answers saying +// the track has no publisher yet, the most actionable first: GOING_AWAY, +// which says to retry (§10.6.2), then TIMEOUT, then DOES_NOT_EXIST and a +// request that failed at the transport. GOING_AWAY is taken to outrank +// §10.2.6's DOES_NOT_EXIST for "no publisher is available": the publisher is +// known, only draining. +func candidateErrRank(err error) int { + if err == nil { + return 0 + } + if _, ok := errors.AsType[*session.ErrUnsupportedMandatoryTrackProperty](err); ok { + return 6 + } + if isTrackPropertiesErr(err) { + return 5 + } + if !awaitsPublisher(err) { + return 4 + } + if errors.Is(err, errGoingAway) { + return 3 + } + rej, _ := errors.AsType[*session.RequestRejectedError](err) + if rej != nil && rej.Code == moqt.RequestGoingAway { + return 3 + } + if rej != nil && rej.Code == moqt.RequestTimeout { + return 2 + } + return 1 // DOES_NOT_EXIST, or no REQUEST_ERROR at all } // subscribeUpstream subscribes fullName on every matching source (§9.5): @@ -562,6 +658,7 @@ func (h *sessionHandler) subscribeUpstream( resultEntry *registry.TrackEntry anyEstab bool lastErr error + retry []*session.Session ) establish := func(sess *session.Session, src string, params message.Parameters) { if subscribed[sess] || ctx.Err() != nil { @@ -591,13 +688,15 @@ func (h *sessionHandler) subscribeUpstream( delete(subscribed, sess) return } - // Keep going. A Track Properties refusal outranks other errors - // (§2.5.1 fixes its downstream code), and any error outranks - // one that only says the track has no publisher yet, so a - // held SUBSCRIBE ends on it (see awaitsPublisher). - if isTrackPropertiesErr(err) || (!isTrackPropertiesErr(lastErr) && awaitsPublisher(lastErr)) { - lastErr = err + // A request that failed at the transport said nothing about + // the track: sess is asked again on a held SUBSCRIBE's next + // look, though not twice in this one. + if transportFailure(err) { + retry = append(retry, sess) } + // Keep going, with the error the subscriber is to get (see + // preferCandidateErr). + lastErr = preferCandidateErr(lastErr, err) h.log.LogAttrs(ctx, slog.LevelDebug, "subscribeUpstream: candidate failed, continuing", slog.String("source", src), slog.String("err", err.Error())) return @@ -618,16 +717,17 @@ func (h *sessionHandler) subscribeUpstream( remotes, draining := h.upstreams.resolveUpstreams(ctx, fullName.Namespace) // A draining relay was sent no request (§10.4); it answers as a draining - // publisher would, ranked with the other candidates' errors. Taken to - // outrank §10.2.6's DOES_NOT_EXIST for "no publisher is available": the - // publisher is known, and GOING_AWAY (§10.6.2) says to retry. - if draining && !isTrackPropertiesErr(lastErr) && awaitsPublisher(lastErr) { - lastErr = errGoingAway + // publisher would, ranked with the other candidates' errors. + if draining { + lastErr = preferCandidateErr(lastErr, errGoingAway) } for _, remote := range remotes { establish(remote, "discovery-remote", remoteExtra) } + for _, sess := range retry { + delete(subscribed, sess) + } if anyEstab { return resultEntry, true, nil } @@ -656,8 +756,11 @@ func (h *sessionHandler) subscribeUpstreamOnSession( if goingAway(sess) { return nil, nil, errGoingAway } - // §9.4: always Next Object (§5.1.2), so one upstream serves every - // downstream filter; the fanout applies those. + // Always Next Object (§5.1.2: StartGroup and StartObject both 0): every + // downstream SUBSCRIBE aggregates onto this one upstream (§9.4 MAY), and + // the fanout applies each downstream filter. The upstream passes only + // Objects after Largest Object as of when it is processed, so on its own + // it cannot serve a downstream range that starts earlier. filter := &message.LocationFilter{Fields: 2} params := message.Parameters{message.LocationFilterParam(filter)} @@ -873,7 +976,7 @@ func upstreamRejection(err error) *session.RequestRejectedError { // elsewhere. return &session.RequestRejectedError{ Code: moqt.RequestGoingAway, - RetryInterval: retryIntervalAfter(time.Second), + RetryInterval: retryIntervalAfter(goingAwayRetryAfter), } } up, ok := errors.AsType[*session.RequestRejectedError](err) @@ -882,11 +985,19 @@ func upstreamRejection(err error) *session.RequestRejectedError { } rej := &session.RequestRejectedError{Code: moqt.RequestInternalError, RetryInterval: up.RetryInterval} switch up.Code { + // GOING_AWAY: an upstream relay draining before its GOAWAY reached this + // one, which answers a draining publisher the same way. case moqt.RequestDoesNotExist, moqt.RequestTimeout, moqt.RequestExcessiveLoad, - moqt.RequestUnsupportedExtension: + moqt.RequestUnsupportedExtension, moqt.RequestGoingAway: rej.Code = up.Code + // Relay policy: a GOING_AWAY saying not to retry (Retry Interval + // 0, §10.6.2) is taken to speak for the upstream's own draining + // hop, and the track may be reached another way. + if up.Code == moqt.RequestGoingAway && up.RetryInterval == 0 { + rej.RetryInterval = retryIntervalAfter(goingAwayRetryAfter) + } case moqt.RequestInternalError, moqt.RequestUnauthorized, moqt.RequestNotSupported, - moqt.RequestMalformedAuthToken, moqt.RequestExpiredAuthToken, moqt.RequestGoingAway, + moqt.RequestMalformedAuthToken, moqt.RequestExpiredAuthToken, moqt.RequestInvalidRange, moqt.RequestInvalidFilter, moqt.RequestRedirect, moqt.RequestMalformedTrack, moqt.RequestUninterested, moqt.RequestPrefixOverlap, moqt.RequestNamespaceTooLarge: @@ -895,6 +1006,14 @@ func upstreamRejection(err error) *session.RequestRejectedError { return rej } +// goingAwayRetryAfter is how long the relay tells a subscriber to wait before +// retrying past a draining upstream; goingAwayRetry is the largest Retry +// Interval its jitter can make of it (see [retryIntervalAfter]). +const ( + goingAwayRetryAfter = time.Second + goingAwayRetry = uint64(goingAwayRetryAfter/time.Millisecond) * 3 / 2 +) + // isTrackPropertiesErr reports whether err is a Track Properties validation // failure from [session.Session.Subscribe] or [session.Session.Fetch]: an // unknown Mandatory Track Property, or Track Properties that do not parse. diff --git a/pkg/relay/handler_track_status.go b/pkg/relay/handler_track_status.go index 5f4cdaae..80886c0c 100644 --- a/pkg/relay/handler_track_status.go +++ b/pkg/relay/handler_track_status.go @@ -4,6 +4,10 @@ import ( "context" "errors" "log/slog" + "slices" + "sync" + "sync/atomic" + "time" "github.com/floatdrop/moq-go/pkg/moqt" "github.com/floatdrop/moq-go/pkg/moqt/message" @@ -11,16 +15,30 @@ import ( "github.com/floatdrop/moq-go/pkg/moqt/track" ) +// trackStatusTimeout, when set by a test, replaces trackStatusUpstreamTimeout. +var trackStatusTimeout atomic.Int64 + +// trackStatusUpstreamTimeout bounds a forwarded TRACK_STATUS's upstream round +// trip, as FILL_TIMEOUT's default bounds a stitch FETCH's (§13.6): a +// candidate that does not answer within it counts as TIMEOUT. +const trackStatusUpstreamTimeout = defaultUpstreamFetchTimeout + // handleTrackStatus implements TRACK_STATUS (§10.15): a metadata-only query for -// a track's Properties and existence, answered without creating a subscription -// or round-tripping upstream. The reply is REQUEST_OK (aliased as -// [message.TrackStatusOK]) carrying the same Track Properties block SUBSCRIBE_OK -// would, plus §10.2.17 LARGEST_OBJECT when objects have been forwarded. +// a track's Properties and existence, which the relay "treats ... identically +// as if it had received a SUBSCRIBE", without creating a subscription. The +// reply is TRACK_STATUS_OK (a REQUEST_OK, [message.TrackStatusOK]) carrying +// what the relay's SUBSCRIBE_OK would: the Track Properties and §10.2.17 +// LARGEST_OBJECT. // -// It answers from the track registry when the track has an established -// upstream or metadata, falls back to an empty TRACK_STATUS_OK when only the -// namespace is advertised locally, and otherwise rejects with -// [moqt.RequestDoesNotExist]. +// A track with an Established subscription is answered from the track +// registry. Otherwise, as SUBSCRIBE would go upstream, TRACK_STATUS is +// forwarded to every candidate SUBSCRIBE would try (§10.15: relays "MAY +// forward TRACK_STATUS to one or more publishers"; see +// [sessionHandler.trackStatusUpstream]), and their answers are combined as a +// SUBSCRIBE_OK's would be: the entry's Track Properties if it has any, else +// the first answer's (§9.6), and the largest LARGEST_OBJECT of all, the +// entry's included (§10.2.17). If none answers OK, the refusal is the one +// SUBSCRIBE would give. func (h *sessionHandler) handleTrackStatus(ctx context.Context, req *session.Request, msg *message.TrackStatus) { if err := h.auth.AuthorizeTrackStatus(ctx, h.sess, msg); err != nil { h.rejectAuth(ctx, req, "TrackStatus", err) @@ -28,54 +46,170 @@ func (h *sessionHandler) handleTrackStatus(ctx context.Context, req *session.Req } fullName := track.FullTrackName{Namespace: msg.Namespace, Name: msg.Name} - entry, known := h.tracks.Get(fullName.Key()) - // Answer TRACK_STATUS_OK for an entry with an established upstream, which - // SUBSCRIBE would accept (§10.15: "treats it identically as if it had - // received a SUBSCRIBE"), or with metadata to surface: Properties or a - // §10.2.17 LargestObject watermark. var ( + properties []byte largest message.Location hasLargest bool ) + entry, known := h.tracks.Get(fullName.Key()) if known { + properties = entry.GetProperties() largest, hasLargest = entry.GetLargest() } - hasProperties := known && len(entry.GetProperties()) > 0 - if known && (hasProperties || hasLargest || hasEstablishedUpstream(entry)) { - reply := &message.TrackStatusOK{} - // §10.2.21: INCLUDE_PROPERTIES=0 empties the Track Properties only. - if hasProperties && includeProperties(msg.Parameters) { - reply.TrackProperties = entry.GetProperties() - } - if hasLargest { - // §10.2.17: omit LARGEST_OBJECT when no objects have - // been observed; emit it (and the watermark) otherwise. - reply.Parameters = append(reply.Parameters, - message.LargestObjectParam(largest.Group, largest.Object)) + if !known || !hasEstablishedUpstream(entry) { + // §13.1: a forwarded TRACK_STATUS holds upstream requests open, so + // it counts against the subscription cap while it does. + if !h.limiter.acquireSub() { + h.rejectExcessiveLoad(ctx, req, "subscription") + return } - // AcceptTrackStatus FINs after the reply (§10.15). - if err := req.AcceptTrackStatus(reply); err != nil { - h.log.LogAttrs(ctx, slog.LevelDebug, "TRACK_STATUS_OK write failed", - slog.String("err", err.Error())) + oks, err := h.trackStatusUpstream(ctx, req, fullName) + h.limiter.releaseSub() + if len(oks) == 0 { + h.rejectTrackStatus(ctx, req, err) + return } - return + properties, largest, hasLargest = mergeTrackStatus(oks, properties, largest, hasLargest) } - // No local track entry with properties. Check the namespace - // registry — if a publisher has advertised the namespace, the track - // at least *might* exist, so we reply with an empty Properties block. - // No LARGEST_OBJECT either: nothing has been observed. - if len(h.names.MatchPublishers(msg.Namespace)) > 0 { - if err := req.AcceptTrackStatus(nil); err != nil { - h.log.LogAttrs(ctx, slog.LevelDebug, "TRACK_STATUS_OK (empty) write failed", - slog.String("err", err.Error())) + reply := &message.TrackStatusOK{} + // §10.2.21: INCLUDE_PROPERTIES=0 empties the Track Properties only. + if includeProperties(msg.Parameters) { + reply.TrackProperties = properties + } + // §10.2.17: LARGEST_OBJECT only once Objects have been published. + if hasLargest { + reply.Parameters = message.Parameters{message.LargestObjectParam(largest.Group, largest.Object)} + } + // AcceptTrackStatus FINs after the reply (§10.15). + if err := req.AcceptTrackStatus(reply); err != nil { + h.log.LogAttrs(ctx, slog.LevelDebug, "TRACK_STATUS_OK write failed", + slog.String("err", err.Error())) + } +} + +// mergeTrackStatus combines TRACK_STATUS_OKs with what the relay already +// knows of the track, as a SUBSCRIBE_OK's would be: its Track Properties if +// any, else the first answer's (§9.6), and the largest LARGEST_OBJECT +// (§10.2.17). +func mergeTrackStatus( + oks []*message.TrackStatusOK, + properties []byte, + largest message.Location, + hasLargest bool, +) ([]byte, message.Location, bool) { + for _, ok := range oks { + if len(properties) == 0 { + properties = ok.TrackProperties + } + p, found := ok.Parameters.Find(message.ParamLargestObject) + if l := (message.Location{Group: p.Group, Object: p.Object}); found && (!hasLargest || largest.Less(l)) { + largest, hasLargest = l, true } - return } + return properties, largest, hasLargest +} - if err := req.RejectError(moqt.RequestDoesNotExist, "relay: track not known"); err != nil && - !errors.Is(err, context.Canceled) { +// rejectTrackStatus answers a TRACK_STATUS no candidate accepted: with the +// refusal SUBSCRIBE would give for err (see [upstreamRejection]), or +// DOES_NOT_EXIST when there was no candidate. +func (h *sessionHandler) rejectTrackStatus(ctx context.Context, req *session.Request, err error) { + rej := &session.RequestRejectedError{Code: moqt.RequestDoesNotExist, Reason: "relay: track not known"} + if err != nil { + rej = upstreamRejection(err) + rej.Reason = "relay: no upstream for track: " + err.Error() + } + if werr := req.Reject(rej); werr != nil && !errors.Is(werr, context.Canceled) { h.log.LogAttrs(ctx, slog.LevelDebug, "TRACK_STATUS reject write failed", - slog.String("err", err.Error())) + slog.String("err", werr.Error())) + } +} + +// trackStatusUpstream forwards TRACK_STATUS for fullName, concurrently, to +// every candidate SUBSCRIBE would try (see [sessionHandler.subscribeUpstream]): +// each local publisher of a covering namespace and each relay Discovery +// resolves. It returns their TRACK_STATUS_OKs in that order and, when there +// are none, the refusal SUBSCRIBE would give ([preferCandidateErr]), nil when there +// was no candidate. +// +// trackStatusUpstreamTimeout bounds the whole forwarding, resolving the +// Discovery candidates included, and it all ends when the requester cancels +// (STOP_SENDING). A draining candidate is sent nothing +// (§10.4). Relay policy, as for SUBSCRIBE and FETCH: the requester's own +// session is skipped while a TRACK_STATUS for the track to it is in flight, +// since this request may be that one routed back, and a second would loop +// (§6.2). +func (h *sessionHandler) trackStatusUpstream( + ctx context.Context, + req *session.Request, + fullName track.FullTrackName, +) ([]*message.TrackStatusOK, error) { + timeout := trackStatusUpstreamTimeout + if d := time.Duration(trackStatusTimeout.Load()); d > 0 { + timeout = d + } + upCtx, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + defer context.AfterFunc(req.Stream.Context(), cancel)() + + key := fullName.Key() + var lastErr error + fail := func(err error) { + lastErr = preferCandidateErr(lastErr, err) + } + var candidates []*session.Session + seen := map[*session.Session]bool{} + add := func(sess *session.Session) { + if seen[sess] { + return + } + seen[sess] = true + switch { + case goingAway(sess): + fail(errGoingAway) + case sess == h.sess && h.tracks.RequestPending(message.TypeTrackStatus, sess, key): + default: + candidates = append(candidates, sess) + } + } + for _, pub := range h.names.MatchPublishers(fullName.Namespace) { + add(pub.Session) + } + remotes, draining := h.upstreams.resolveUpstreams(upCtx, fullName.Namespace) + if draining { + fail(errGoingAway) + } + for _, remote := range remotes { + add(remote) + } + + oks := make([]*message.TrackStatusOK, len(candidates)) + errs := make([]error, len(candidates)) + var wg sync.WaitGroup + for i, sess := range candidates { + done := h.tracks.BeginRequest(message.TypeTrackStatus, sess, key) + wg.Go(func() { + defer done() + ts, err := sess.TrackStatus(upCtx, &message.TrackStatus{Namespace: fullName.Namespace, Name: fullName.Name}) + if err != nil { + if errors.Is(err, context.DeadlineExceeded) { + err = &session.RequestRejectedError{ + Code: moqt.RequestTimeout, + Reason: "relay: TRACK_STATUS timed out", + } + } + errs[i] = err + return + } + _ = ts.Close() + oks[i] = ts.OK + }) + } + wg.Wait() + for _, err := range errs { + if err != nil { + fail(err) + } } + return slices.DeleteFunc(oks, func(ok *message.TrackStatusOK) bool { return ok == nil }), lastErr } diff --git a/pkg/relay/internal/registry/track_registry.go b/pkg/relay/internal/registry/track_registry.go index b000e885..8d6a91b5 100644 --- a/pkg/relay/internal/registry/track_registry.go +++ b/pkg/relay/internal/registry/track_registry.go @@ -96,9 +96,10 @@ type TrackRegistry struct { // claims holds the upstream SUBSCRIBEs in flight, one per (session, // track); see [TrackRegistry.ClaimUpstream]. Guarded by mu. claims map[upstreamClaim]struct{} - // fetches counts the stitch FETCHes in flight, per (session, track); see - // [TrackRegistry.BeginFetch]. Guarded by mu. - fetches map[upstreamClaim]int + // inflight counts the relay's own FETCHes and TRACK_STATUSes in flight, + // per (type, session, track); see [TrackRegistry.BeginRequest]. Guarded + // by mu. + inflight map[inflightRequest]int // arrivals are the waiters for a track's next upstream; see // [TrackRegistry.AwaitUpstream]. Guarded by mu. arrivals arrivals[track.Key] @@ -225,30 +226,39 @@ func (r *TrackRegistry) ClaimUpstream(sess *session.Session, key track.Key) (rel }, true } -// BeginFetch marks a stitch FETCH for key on sess as in flight until done; -// see [TrackRegistry.FetchPending]. -func (r *TrackRegistry) BeginFetch(sess *session.Session, key track.Key) (done func()) { - c := upstreamClaim{sess: sess, key: key} +// inflightRequest is one kind of request the relay sends for a track on a +// session. +type inflightRequest struct { + upstreamClaim + + typ message.Type +} + +// BeginRequest marks a request of type typ for key on sess as in flight until +// done; see [TrackRegistry.RequestPending]. +func (r *TrackRegistry) BeginRequest(typ message.Type, sess *session.Session, key track.Key) (done func()) { + c := inflightRequest{typ: typ, sess: sess, key: key} r.mu.Lock() defer r.mu.Unlock() - if r.fetches == nil { - r.fetches = make(map[upstreamClaim]int) + if r.inflight == nil { + r.inflight = make(map[inflightRequest]int) } - r.fetches[c]++ + r.inflight[c]++ return func() { r.mu.Lock() defer r.mu.Unlock() - if r.fetches[c]--; r.fetches[c] == 0 { - delete(r.fetches, c) + if r.inflight[c]--; r.inflight[c] == 0 { + delete(r.inflight, c) } } } -// FetchPending reports whether a stitch FETCH for key on sess is in flight. -func (r *TrackRegistry) FetchPending(sess *session.Session, key track.Key) bool { +// RequestPending reports whether a request of type typ for key on sess is in +// flight. +func (r *TrackRegistry) RequestPending(typ message.Type, sess *session.Session, key track.Key) bool { r.mu.RLock() defer r.mu.RUnlock() - return r.fetches[upstreamClaim{sess: sess, key: key}] > 0 + return r.inflight[inflightRequest{typ: typ, sess: sess, key: key}] > 0 } // ReleaseIfUnsubscribed removes the on-demand upstream up from the entry for diff --git a/pkg/relay/limiter.go b/pkg/relay/limiter.go index 3acfe8b6..f33ef594 100644 --- a/pkg/relay/limiter.go +++ b/pkg/relay/limiter.go @@ -11,7 +11,8 @@ import "sync" // // The counts track in-flight request handlers: acquire is called at dispatch // before a handler is spawned, release when it returns (a handler runs for its -// request's whole lifetime). An over-limit request is rejected with +// request's whole lifetime); a TRACK_STATUS holds a subscription slot only +// while it is forwarded upstream. An over-limit request is rejected with // REQUEST_ERROR EXCESSIVE_LOAD before any shared state is mutated. type sessionLimiter struct { mu sync.Mutex diff --git a/pkg/relay/pool_goaway_test.go b/pkg/relay/pool_goaway_test.go index 41d6a6ab..fd2340cb 100644 --- a/pkg/relay/pool_goaway_test.go +++ b/pkg/relay/pool_goaway_test.go @@ -100,3 +100,31 @@ func testDrainingUpstreamRelay(t *testing.T, localRefuser bool) { return ok && rej.Code == moqt.RequestGoingAway }, "a SUBSCRIBE whose only upstream relay is draining was never refused with GOING_AWAY") } + +// TestRendezvous_UpstreamGoingAwayKeepsHold: a publisher answering the relay's +// SUBSCRIBE with GOING_AWAY, as an upstream relay draining does before its +// GOAWAY arrives, leaves the track without a publisher yet, so a +// RENDEZVOUS_TIMEOUT hold goes on (§10.2.6). Passing the code on is covered by +// TestSubscribe_UpstreamRejects_PropagatesRejection. +func TestRendezvous_UpstreamGoingAwayKeepsHold(t *testing.T) { + t.Parallel() + subSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + draining := dialAnotherClient(t, subSess) + publishNS(t, draining, "video") + go func() { + for { + req, err := draining.AcceptRequest(t.Context()) + if err != nil { + return + } + _ = req.RejectError(moqt.RequestGoingAway, "draining") + } + }() + done := subscribeRendezvous(t.Context(), subSess, 5*time.Second) + requireHeld(t, done) + publishVideoTrack(t, dialAnotherClient(t, subSess), "cam1", 7) + if err := awaitAnswer(t, done); err != nil { + t.Fatalf("held SUBSCRIBE: %v", err) + } +} diff --git a/pkg/relay/relay.go b/pkg/relay/relay.go index 03561bb9..db62783f 100644 --- a/pkg/relay/relay.go +++ b/pkg/relay/relay.go @@ -165,7 +165,8 @@ type Config struct { // MaxSubscriptionsPerSession bounds the number of concurrently-active // SUBSCRIBE requests a single session may hold (§13.1, subscription - // amplification). Excess SUBSCRIBEs are rejected with REQUEST_ERROR + // amplification), counting a TRACK_STATUS while the relay forwards it + // upstream. Excess requests are rejected with REQUEST_ERROR // EXCESSIVE_LOAD before any state is mutated. Zero (the default) means // unlimited — limits are a deployment policy the operator opts into. MaxSubscriptionsPerSession int diff --git a/pkg/relay/relaynet/relaynet.go b/pkg/relay/relaynet/relaynet.go index b2d46028..f625e3e4 100644 --- a/pkg/relay/relaynet/relaynet.go +++ b/pkg/relay/relaynet/relaynet.go @@ -32,14 +32,15 @@ import ( "github.com/floatdrop/moq-go/pkg/moqt/session/quicconn" ) -// MOQTQUICALPNs lists the raw-QUIC MOQT ALPNs the relay accepts. Draft-19 -// SETUP carries no version field (§3.1), so the "moqt-NN" ALPN is itself the -// draft-version signal — the negotiated ALPN fixes the draft. We advertise -// only "moqt-20", the draft this implementation speaks. The older -// "moqt-18"/"-17"/"-16" and the pre-15 "moq-00" (which expected in-SETUP -// version negotiation, removed in -19) are deliberately not offered: our -19 -// wire behavior can't complete a SETUP with a peer that selected any of them, -// so advertising them would only let such a peer clear TLS and then fail. +// MOQTQUICALPNs lists the raw-QUIC MOQT ALPNs the relay accepts. SETUP +// carries no version field (§10.3); MOQT negotiates the version with ALPN +// (§3.1), so the "moqt-NN" ALPN is itself the draft-version signal — the +// negotiated ALPN fixes the draft. We advertise only "moqt-20", the draft this +// implementation speaks. The older "moqt-19"/"-18"/"-17"/"-16" and "moq-00" +// (used before -15, followed by version negotiation in SETUP — §3.1) are +// deliberately not offered: our -20 wire behavior can't complete a SETUP with +// a peer that selected any of them, so advertising them would only let such a +// peer clear TLS and then fail. var MOQTQUICALPNs = []string{"moqt-20"} // defaultQUICConfig returns the QUIC tuning the relay listens and dials with: diff --git a/pkg/relay/session_upstream_test.go b/pkg/relay/session_upstream_test.go index 917eb207..7db7423f 100644 --- a/pkg/relay/session_upstream_test.go +++ b/pkg/relay/session_upstream_test.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "slices" + "sync/atomic" "testing" "time" @@ -304,9 +305,11 @@ func TestSubscribe_UpstreamAliasReusableAfterTeardown(t *testing.T) { } // TestSubscribe_UpstreamRejects_PropagatesRejection: an upstream REQUEST_ERROR -// code about the track passes downstream; one about the relay's own hop, or not -// defined for SUBSCRIBE (MALFORMED_TRACK answers a FETCH), becomes -// INTERNAL_ERROR (§10.6.2). The Retry Interval is kept either way. +// code about the track passes downstream, GOING_AWAY among them (an upstream +// relay draining, as the relay answers for a draining publisher itself); one +// about the relay's own hop, or not defined for SUBSCRIBE (MALFORMED_TRACK +// answers a FETCH), becomes INTERNAL_ERROR (§10.6.2). The Retry Interval is +// kept either way. func TestSubscribe_UpstreamRejects_PropagatesRejection(t *testing.T) { t.Parallel() for _, tc := range []struct { @@ -319,7 +322,7 @@ func TestSubscribe_UpstreamRejects_PropagatesRejection(t *testing.T) { {moqt.RequestMalformedTrack, moqt.RequestInternalError, 0}, {moqt.RequestUnauthorized, moqt.RequestInternalError, 0}, {moqt.RequestExpiredAuthToken, moqt.RequestInternalError, 2001}, - {moqt.RequestGoingAway, moqt.RequestInternalError, 0}, + {moqt.RequestGoingAway, moqt.RequestGoingAway, 7001}, {moqt.RequestInvalidRange, moqt.RequestInternalError, 0}, {moqt.RequestErrorCode(0x7777), moqt.RequestInternalError, 31}, } { @@ -345,3 +348,265 @@ func TestSubscribe_UpstreamRejects_PropagatesRejection(t *testing.T) { }) } } + +// TestSubscribe_NoPublisherYetRanking: when every candidate only says the +// track has no publisher yet, the refusal is the most actionable of their +// answers whatever the order they answer in: GOING_AWAY (retry soon), then +// TIMEOUT, then DOES_NOT_EXIST. +func TestSubscribe_NoPublisherYetRanking(t *testing.T) { + t.Parallel() + for _, tc := range []struct { + a, b, want moqt.RequestErrorCode + }{ + {moqt.RequestDoesNotExist, moqt.RequestGoingAway, moqt.RequestGoingAway}, + {moqt.RequestDoesNotExist, moqt.RequestTimeout, moqt.RequestTimeout}, + {moqt.RequestTimeout, moqt.RequestGoingAway, moqt.RequestGoingAway}, + } { + for _, order := range [][2]moqt.RequestErrorCode{{tc.a, tc.b}, {tc.b, tc.a}} { + t.Run(fmt.Sprintf("%#x then %#x", uint64(order[0]), uint64(order[1])), func(t *testing.T) { + t.Parallel() + subSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + for _, code := range order { + pub := dialAnotherClient(t, subSess) + publishNS(t, pub, "video") + go func() { + for { + req, err := pub.AcceptRequest(t.Context()) + if err != nil { + return + } + _ = req.RejectError(code, "no cam1") + } + }() + } + _, err := subSess.Subscribe( + t.Context(), + &message.Subscribe{Namespace: ns("video"), Name: []byte("cam1")}, + ) + requireRejectedWithCode(t, err, tc.want) + }) + } + } +} + +// refusingPublishers adds, on subSess's relay, one PUBLISH_NAMESPACE video +// publisher per answer, in order; each answers every SUBSCRIBE with it. +func refusingPublishers(t *testing.T, subSess *session.Session, answers ...func(*session.Request)) { + t.Helper() + for _, answer := range answers { + pub := dialAnotherClient(t, subSess) + publishNS(t, pub, "video") + go func() { + for { + req, err := pub.AcceptRequest(t.Context()) + if err != nil { + return + } + answer(req) + } + }() + } +} + +func refuse(code moqt.RequestErrorCode, retry uint64) func(*session.Request) { + return func(r *session.Request) { + _ = r.Reject(&session.RequestRejectedError{Code: code, RetryInterval: retry, Reason: "no"}) + } +} + +// reset cancels the request at the transport (§3.3.3): no REQUEST_ERROR. +func reset(r *session.Request) { + r.Stream.CancelRead(uint64(moqt.StreamResetCancelled)) + r.Stream.CancelWrite(uint64(moqt.StreamResetCancelled)) +} + +func subscribeCam1Rejected(t *testing.T, subSess *session.Session) *session.RequestRejectedError { + t.Helper() + _, err := subSess.Subscribe(t.Context(), &message.Subscribe{Namespace: ns("video"), Name: []byte("cam1")}) + rej, ok := errors.AsType[*session.RequestRejectedError](err) + if !ok { + t.Fatalf("Subscribe = %v, want a REQUEST_ERROR", err) + } + return rej +} + +// TestSubscribe_UpstreamGoingAwayWithoutRetry: an upstream GOING_AWAY that says +// not to retry (Retry Interval 0, §10.6.2) speaks for the upstream's own +// draining hop, so the subscriber is told to retry after the relay's own +// interval, as when the relay sees the drain itself. +func TestSubscribe_UpstreamGoingAwayWithoutRetry(t *testing.T) { + t.Parallel() + subSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + refusingPublishers(t, subSess, refuse(moqt.RequestGoingAway, 0)) + rej := subscribeCam1Rejected(t, subSess) + if rej.Code != moqt.RequestGoingAway || rej.RetryInterval < 1001 || rej.RetryInterval > 1501 { + t.Fatalf("got %#x with Retry Interval %d, want GOING_AWAY retrying after about 1s", + uint64(rej.Code), rej.RetryInterval) + } +} + +// TestSubscribe_TiedRefusalsSoonestRetry: of refusals with the same code, the +// subscriber gets the one allowing the soonest retry, whatever the order; +// "SHOULD NOT be retried" (0, §10.6.2) only when every one says so. +func TestSubscribe_TiedRefusalsSoonestRetry(t *testing.T) { + t.Parallel() + for _, tc := range []struct { + code moqt.RequestErrorCode + a, b, want uint64 + }{ + {moqt.RequestDoesNotExist, 0, 501, 501}, + {moqt.RequestDoesNotExist, 0, 0, 0}, + {moqt.RequestTimeout, 3001, 1001, 1001}, + } { + for _, order := range [][2]uint64{{tc.a, tc.b}, {tc.b, tc.a}} { + t.Run(fmt.Sprintf("%#x %d then %d", uint64(tc.code), order[0], order[1]), func(t *testing.T) { + t.Parallel() + subSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + refusingPublishers(t, subSess, refuse(tc.code, order[0]), refuse(tc.code, order[1])) + if rej := subscribeCam1Rejected(t, subSess); rej.Code != tc.code || rej.RetryInterval != tc.want { + t.Fatalf("got %#x with Retry Interval %d, want %#x with %d", + uint64(rej.Code), rej.RetryInterval, uint64(tc.code), tc.want) + } + }) + } + } +} + +// TestSubscribe_TransportFailureIsNoPublisherYet: a candidate whose request +// fails at the transport, with no REQUEST_ERROR, says nothing about the +// track, so it ranks as DOES_NOT_EXIST (the code it is answered with): it +// does not mask another candidate's GOING_AWAY, and a RENDEZVOUS_TIMEOUT hold +// goes on (§10.2.6). +func TestSubscribe_TransportFailureIsNoPublisherYet(t *testing.T) { + t.Parallel() + for _, order := range []string{"reset first", "reset last"} { + t.Run(order, func(t *testing.T) { + t.Parallel() + subSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + answers := []func(*session.Request){reset, refuse(moqt.RequestGoingAway, 2001)} + if order == "reset last" { + answers[0], answers[1] = answers[1], answers[0] + } + refusingPublishers(t, subSess, answers...) + if rej := subscribeCam1Rejected(t, subSess); rej.Code != moqt.RequestGoingAway { + t.Fatalf("got %#x, want the other candidate's GOING_AWAY", uint64(rej.Code)) + } + }) + } + t.Run("hold", func(t *testing.T) { + t.Parallel() + subSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + refusingPublishers(t, subSess, reset) + done := subscribeRendezvous(t.Context(), subSess, 5*time.Second) + requireHeld(t, done) + publishVideoTrack(t, dialAnotherClient(t, subSess), "cam1", 7) + if err := awaitAnswer(t, done); err != nil { + t.Fatalf("held SUBSCRIBE: %v", err) + } + }) +} + +// TestSubscribe_UnsupportedMandatoryPropertyOutranksMalformed: of two +// candidates whose SUBSCRIBE_OKs carry bad Track Properties, an unknown +// Mandatory Track Property (§2.5.1: UNSUPPORTED_EXTENSION, a MUST) wins over +// Properties that do not parse, whatever the order. +func TestSubscribe_UnsupportedMandatoryPropertyOutranksMalformed(t *testing.T) { + t.Parallel() + accept := func(props []byte) func(*session.Request) { + return func(r *session.Request) { + _, _ = r.AcceptSubscribe(&message.SubscribeOK{TrackAlias: 5, TrackProperties: props}) + } + } + for _, order := range []string{"mandatory first", "mandatory last"} { + t.Run(order, func(t *testing.T) { + t.Parallel() + subSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + answers := []func(*session.Request){accept(mandatoryProps()), accept(malformedProps)} + if order == "mandatory last" { + answers[0], answers[1] = answers[1], answers[0] + } + refusingPublishers(t, subSess, answers...) + if rej := subscribeCam1Rejected(t, subSess); rej.Code != moqt.RequestUnsupportedExtension { + t.Fatalf("got %#x, want UNSUPPORTED_EXTENSION", uint64(rej.Code)) + } + }) + } +} + +// TestRendezvous_NoStreamCreditEndsHold: a SUBSCRIBE the relay cannot even +// open to a live publisher, for want of bidi-stream credit, is not a sign the +// track has no publisher, so a RENDEZVOUS_TIMEOUT hold does not wait out its +// budget on it: the subscriber is answered at once, as without a hold. +func TestRendezvous_NoStreamCreditEndsHold(t *testing.T) { + t.Parallel() + subSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + pub := dialAnotherClientWithLimits(t, subSess, -1, 0) // the relay may open no stream to it + publishNS(t, pub, "video") + done := subscribeRendezvous(t.Context(), subSess, 5*time.Second) + select { + case err := <-done: + requireRejectedWithCode(t, err, moqt.RequestDoesNotExist) + case <-time.After(time.Second): + t.Fatal("SUBSCRIBE held against a live publisher the relay had no stream credit for") + } +} + +// TestSubscribe_OtherRefusalTieIsOrderFree: two refusals of the "any other" +// kind with the same Retry Interval, one answered as INTERNAL_ERROR (an +// UNAUTHORIZED about the relay's hop) and one passed on (EXCESSIVE_LOAD), +// give the subscriber the specific code whatever the order. +func TestSubscribe_OtherRefusalTieIsOrderFree(t *testing.T) { + t.Parallel() + for _, order := range []string{"specific first", "specific last"} { + t.Run(order, func(t *testing.T) { + t.Parallel() + subSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + answers := []func(*session.Request){ + refuse(moqt.RequestExcessiveLoad, 0), + refuse(moqt.RequestUnauthorized, 0), + } + if order == "specific last" { + answers[0], answers[1] = answers[1], answers[0] + } + refusingPublishers(t, subSess, answers...) + if rej := subscribeCam1Rejected(t, subSess); rej.Code != moqt.RequestExcessiveLoad { + t.Fatalf("got %#x, want EXCESSIVE_LOAD", uint64(rej.Code)) + } + }) + } +} + +// TestRendezvous_TransportFailureAskedAgain: a publisher whose SUBSCRIBE failed +// at the transport is asked again on the hold's next look, and serves the +// subscriber if it has the track by then. +func TestRendezvous_TransportFailureAskedAgain(t *testing.T) { + t.Parallel() + subSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + var asked atomic.Int32 + refusingPublishers(t, subSess, func(r *session.Request) { + if asked.Add(1) == 1 { + reset(r) + return + } + _, _ = r.AcceptSubscribe(&message.SubscribeOK{TrackAlias: 9}) + }) + done := subscribeRendezvous(t.Context(), subSess, 5*time.Second) + requireHeld(t, done) + // Another publisher of the namespace arrives, waking the hold. + refusingPublishers(t, subSess, refuse(moqt.RequestDoesNotExist, 0)) + if err := awaitAnswer(t, done); err != nil { + t.Fatalf("held SUBSCRIBE: %v", err) + } + if n := asked.Load(); n != 2 { + t.Fatalf("the reset publisher was asked %d times, want 2", n) + } +} diff --git a/pkg/relay/track_status_forward_test.go b/pkg/relay/track_status_forward_test.go new file mode 100644 index 00000000..37ab838a --- /dev/null +++ b/pkg/relay/track_status_forward_test.go @@ -0,0 +1,350 @@ +package relay_test + +import ( + "bytes" + "context" + "errors" + "fmt" + "sync/atomic" + "testing" + "time" + + "github.com/floatdrop/moq-go/pkg/moqt" + "github.com/floatdrop/moq-go/pkg/moqt/message" + "github.com/floatdrop/moq-go/pkg/moqt/session" + "github.com/floatdrop/moq-go/pkg/relay" + "github.com/floatdrop/moq-go/pkg/relay/discovery" +) + +// A relay with no Established subscription "MAY forward TRACK_STATUS to one or +// more publishers" (§10.15). This one forwards it to the candidates SUBSCRIBE +// would try, and answers with the first TRACK_STATUS_OK or the refusal +// SUBSCRIBE would give. + +// answerTrackStatus answers every TRACK_STATUS sess receives with reply, and +// counts them. +func answerTrackStatus(t *testing.T, sess *session.Session, reply func(*session.Request)) *atomic.Int32 { + t.Helper() + var n atomic.Int32 + go func() { + for { + req, err := sess.AcceptRequest(t.Context()) + if err != nil { + return + } + if _, ok := req.First.(*message.TrackStatus); ok { + n.Add(1) + reply(req) + } + } + }() + return &n +} + +func trackStatusCam1(t *testing.T, sess *session.Session, params ...message.Parameter) (*message.TrackStatusOK, error) { + t.Helper() + ts, err := sess.TrackStatus(t.Context(), &message.TrackStatus{ + Namespace: ns("video"), Name: []byte("cam1"), Parameters: params, + }) + if err != nil { + return nil, err + } + _ = ts.Close() + return ts.OK, nil +} + +// TestTrackStatus_ForwardsToNamespacePublisher: with only the namespace +// advertised, the publisher's TRACK_STATUS_OK is passed on: its Track +// Properties, unless INCLUDE_PROPERTIES is 0 (§10.2.21), and its +// LARGEST_OBJECT (§10.2.17), as a SUBSCRIBE_OK would carry them. +func TestTrackStatus_ForwardsToNamespacePublisher(t *testing.T) { + t.Parallel() + for _, include := range []bool{true, false} { + t.Run(map[bool]string{true: "with properties", false: "INCLUDE_PROPERTIES=0"}[include], func(t *testing.T) { + t.Parallel() + pub, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + publishNS(t, pub, "video") + asked := answerTrackStatus(t, pub, func(r *session.Request) { + _ = r.AcceptTrackStatus(&message.TrackStatusOK{ + TrackProperties: opaqueProps("h265"), + Parameters: message.Parameters{message.LargestObjectParam(3, 4)}, + }) + }) + var params []message.Parameter + if !include { + params = append(params, message.IncludePropertiesParam(false)) + } + ok, err := trackStatusCam1(t, dialAnotherClient(t, pub), params...) + if err != nil { + t.Fatalf("TrackStatus: %v", err) + } + if asked.Load() != 1 { + t.Fatalf("the publisher was asked %d times, want 1", asked.Load()) + } + wantProps := opaqueProps("h265") + if !include { + wantProps = nil + } + if !bytes.Equal(ok.TrackProperties, wantProps) { + t.Errorf("TrackProperties = %x, want %x", ok.TrackProperties, wantProps) + } + if p, found := ok.Parameters.Find(message.ParamLargestObject); !found || p.Group != 3 || p.Object != 4 { + t.Errorf("LARGEST_OBJECT = %+v (found %v), want {3, 4}", p, found) + } + }) + } +} + +// TestTrackStatus_PassesUpstreamRefusal: a refusal from the publisher is +// answered as SUBSCRIBE would answer it: DOES_NOT_EXIST passed on, a code about +// the relay's own hop as INTERNAL_ERROR (§10.6.2). +func TestTrackStatus_PassesUpstreamRefusal(t *testing.T) { + t.Parallel() + for _, tc := range []struct{ upstream, want moqt.RequestErrorCode }{ + {moqt.RequestDoesNotExist, moqt.RequestDoesNotExist}, + {moqt.RequestUnauthorized, moqt.RequestInternalError}, + } { + t.Run(fmt.Sprintf("%#x", uint64(tc.upstream)), func(t *testing.T) { + t.Parallel() + pub, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + publishNS(t, pub, "video") + answerTrackStatus(t, pub, func(r *session.Request) { _ = r.RejectError(tc.upstream, "no") }) + _, err := trackStatusCam1(t, dialAnotherClient(t, pub)) + requireRejectedWithCode(t, err, tc.want) + }) + } +} + +// TestTrackStatus_LeftoverEntryAsksUpstream: a track whose publisher left, +// with its Track Properties still on the relay's entry, has no Established +// subscription, so TRACK_STATUS is answered as SUBSCRIBE would be: with no +// publisher left to ask, DOES_NOT_EXIST. +func TestTrackStatus_LeftoverEntryAsksUpstream(t *testing.T) { + t.Parallel() + pubSess, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + pub := publishVideoTrackProps(t, pubSess, "cam1", 7, trackProp(0x10, 1)) + querySess := dialAnotherClient(t, pubSess) + subscribeCam1(t, querySess) // keeps the entry once the publisher leaves + if _, err := trackStatusCam1(t, querySess); err != nil { + t.Fatalf("TrackStatus while published: %v", err) + } + _ = pub.Close() + waitFor(t, 2*time.Second, func() bool { + _, err := subscribeCam1Err(t, dialAnotherClient(t, pubSess)) + return err != nil + }, "the track still has an upstream after its publisher left") + + _, err := trackStatusCam1(t, querySess) + requireRejectedWithCode(t, err, moqt.RequestDoesNotExist) +} + +// subscribeCam1Err SUBSCRIBEs to video/cam1, returning the result. +func subscribeCam1Err(t *testing.T, sess *session.Session) (*session.Subscription, error) { + t.Helper() + sub, err := sess.Subscribe(t.Context(), &message.Subscribe{Namespace: ns("video"), Name: []byte("cam1")}) + if err == nil { + t.Cleanup(func() { _ = sub.Close() }) + } + return sub, err +} + +// TestCrossRelay_TrackStatusForwardedToRemote: a relay Discovery resolves is a +// candidate as it is for SUBSCRIBE, and its answer is passed on. +func TestCrossRelay_TrackStatusForwardedToRemote(t *testing.T) { + t.Parallel() + store := discovery.NewMemoryStore() + defer store.Close() + relayA, relayB := startRelayPair(t.Context(), store) + defer func() { relayA.stop(t); relayB.stop(t) }() + + pub := dialClient(t, relayB) + publishNS(t, pub, "video") + answerTrackStatus(t, pub, func(r *session.Request) { + _ = r.AcceptTrackStatus(&message.TrackStatusOK{TrackProperties: opaqueProps("remote")}) + }) + ok, err := trackStatusCam1(t, dialClient(t, relayA)) + if err != nil { + t.Fatalf("TrackStatus: %v", err) + } + if !bytes.Equal(ok.TrackProperties, opaqueProps("remote")) { + t.Errorf("TrackProperties = %x, want the remote publisher's", ok.TrackProperties) + } +} + +// TestRelay_TrackStatusLoopStopsAtSecondHop: a peer relay that routes the +// relay's TRACK_STATUS straight back on the same session (§6.2 has no loop +// protection) is asked once: the routed-back request finds the relay's own +// still pending on that session and is not forwarded to it again. +func TestRelay_TrackStatusLoopStopsAtSecondHop(t *testing.T) { + t.Parallel() + peer, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + publishNS(t, peer, "video") + var routed atomic.Int32 + asked := answerTrackStatus(t, peer, func(r *session.Request) { + go func() { + if routed.Add(1) > loopCap { // a relay that does not stop still ends + _ = r.RejectError(moqt.RequestDoesNotExist, "loop cap") + return + } + if _, err := trackStatusCam1(t, peer); err != nil { + _ = r.RejectError(moqt.RequestDoesNotExist, "routed back: "+err.Error()) + return + } + _ = r.AcceptTrackStatus(nil) + }() + }) + _, _ = trackStatusCam1(t, dialAnotherClient(t, peer)) + if n := asked.Load(); n != 1 { + t.Fatalf("the relay forwarded TRACK_STATUS to the peer %d times, want 1", n) + } +} + +// namespacePeer is a fresh session on the relay sess is on that +// PUBLISH_NAMESPACEs video and answers each TRACK_STATUS with reply. +func namespacePeer(t *testing.T, sess *session.Session, reply func(*session.Request)) *session.Session { + t.Helper() + p := dialAnotherClient(t, sess) + answerTrackStatus(t, p, reply) + publishNS(t, p, "video") + return p +} + +func acceptTrackStatusWith(props []byte, largest *message.Location) func(*session.Request) { + return func(r *session.Request) { + ok := &message.TrackStatusOK{TrackProperties: props} + if largest != nil { + ok.Parameters = message.Parameters{message.LargestObjectParam(largest.Group, largest.Object)} + } + _ = r.AcceptTrackStatus(ok) + } +} + +// TestTrackStatus_AsksEveryCandidate: every candidate is asked, as SUBSCRIBE +// subscribes on every matching publisher (§9.5), and the answers combine as a +// SUBSCRIBE_OK's would: the largest LARGEST_OBJECT (§10.2.17), the first +// answer's Track Properties (§9.6). A refusal from one does not keep another's +// answer from the subscriber. +func TestTrackStatus_AsksEveryCandidate(t *testing.T) { + t.Parallel() + t.Run("largest of two", func(t *testing.T) { + t.Parallel() + client, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + namespacePeer(t, client, acceptTrackStatusWith(opaqueProps("first"), &message.Location{Group: 3})) + namespacePeer(t, client, acceptTrackStatusWith(opaqueProps("second"), &message.Location{Group: 7, Object: 2})) + ok, err := trackStatusCam1(t, client) + if err != nil { + t.Fatalf("TrackStatus: %v", err) + } + if p, found := ok.Parameters.Find(message.ParamLargestObject); !found || p.Group != 7 || p.Object != 2 { + t.Errorf("LARGEST_OBJECT = %+v (found %v), want the larger {7, 2}", p, found) + } + if !bytes.Equal(ok.TrackProperties, opaqueProps("first")) { + t.Errorf("TrackProperties = %x, want the first answer's", ok.TrackProperties) + } + }) + t.Run("one refuses", func(t *testing.T) { + t.Parallel() + client, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + namespacePeer(t, client, func(r *session.Request) { _ = r.RejectError(moqt.RequestDoesNotExist, "no") }) + namespacePeer(t, client, acceptTrackStatusWith(opaqueProps("second"), nil)) + ok, err := trackStatusCam1(t, client) + if err != nil { + t.Fatalf("TrackStatus: %v", err) + } + if !bytes.Equal(ok.TrackProperties, opaqueProps("second")) { + t.Errorf("TrackProperties = %x, want the answering publisher's", ok.TrackProperties) + } + }) +} + +// TestTrackStatus_UpstreamMandatoryPropertyRefused: an upstream +// TRACK_STATUS_OK with an unknown Mandatory Track Property is refused with +// UNSUPPORTED_EXTENSION, as a SUBSCRIBE_OK carrying it would be (§2.5.1). +func TestTrackStatus_UpstreamMandatoryPropertyRefused(t *testing.T) { + t.Parallel() + client, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + namespacePeer(t, client, acceptTrackStatusWith(mandatoryProps(), nil)) + _, err := trackStatusCam1(t, client) + requireRejectedWithCode(t, err, moqt.RequestUnsupportedExtension) +} + +// TestTrackStatus_DrainingCandidateGoingAway: a candidate that sent GOAWAY is +// sent nothing (§10.4), and the refusal is GOING_AWAY, as for SUBSCRIBE. +func TestTrackStatus_DrainingCandidateGoingAway(t *testing.T) { + t.Parallel() + client, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + draining := namespacePeer(t, client, func(r *session.Request) { _ = r.RejectError(moqt.RequestDoesNotExist, "no") }) + if err := draining.SendGoaway(0, ""); err != nil { + t.Fatalf("SendGoaway: %v", err) + } + waitFor(t, 2*time.Second, func() bool { + _, err := trackStatusCam1(t, client) + rej, ok := errors.AsType[*session.RequestRejectedError](err) + return ok && rej.Code == moqt.RequestGoingAway + }, "TRACK_STATUS to a draining publisher never refused with GOING_AWAY") +} + +// TestTrackStatus_SilentCandidateTimesOut: a candidate that does not answer +// counts as TIMEOUT once the upstream round trip's bound runs out (§13.6), +// and the requester is answered. +func TestTrackStatus_SilentCandidateTimesOut(t *testing.T) { + defer relay.SetTrackStatusTimeout(200 * time.Millisecond)() + client, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + namespacePeer(t, client, func(*session.Request) {}) // never answers + _, err := trackStatusCam1(t, client) + requireRejectedWithCode(t, err, moqt.RequestTimeout) +} + +// TestTrackStatus_CancelEndsUpstream: the requester cancelling its +// TRACK_STATUS (§3.3.3) cancels the relay's upstream one. +func TestTrackStatus_CancelEndsUpstream(t *testing.T) { + t.Parallel() + client, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + upstream := make(chan session.Stream, 1) + namespacePeer(t, client, func(r *session.Request) { upstream <- r.Stream }) + ctx, cancel := context.WithCancel(t.Context()) + go func() { + _, _ = client.TrackStatus(ctx, &message.TrackStatus{Namespace: ns("video"), Name: []byte("cam1")}) + }() + var s session.Stream + select { + case s = <-upstream: + case <-time.After(2 * time.Second): + t.Fatal("the relay never forwarded the TRACK_STATUS") + } + cancel() + select { + case <-s.Context().Done(): + case <-time.After(2 * time.Second): + t.Fatal("the relay's upstream TRACK_STATUS outlived the requester's cancel") + } +} + +// TestTrackStatus_ForwardingCountsAgainstSubscriptionCap: a TRACK_STATUS +// waiting on upstreams counts against MaxSubscriptionsPerSession (§13.1), so +// one more past the cap is refused with EXCESSIVE_LOAD. +func TestTrackStatus_ForwardingCountsAgainstSubscriptionCap(t *testing.T) { + t.Parallel() + client, teardown := connectRelay(t, relay.Config{MaxSubscriptionsPerSession: 1}) + t.Cleanup(teardown) + asked := make(chan struct{}, 2) + namespacePeer(t, client, func(*session.Request) { asked <- struct{}{} }) // never answers + go func() { _, _ = trackStatusCam1(t, client) }() + select { + case <-asked: + case <-time.After(2 * time.Second): + t.Fatal("the first TRACK_STATUS was never forwarded") + } + _, err := trackStatusCam1(t, client) + requireRejectedWithCode(t, err, moqt.RequestExcessiveLoad) +} diff --git a/pkg/relay/track_status_test.go b/pkg/relay/track_status_test.go index fc7bee10..2c98c427 100644 --- a/pkg/relay/track_status_test.go +++ b/pkg/relay/track_status_test.go @@ -5,6 +5,7 @@ import ( "time" "github.com/floatdrop/moq-go/pkg/moqt/message" + "github.com/floatdrop/moq-go/pkg/moqt/session" "github.com/floatdrop/moq-go/pkg/relay" ) @@ -108,6 +109,8 @@ func TestRelay_TrackStatusRequestUpdateClosesSession(t *testing.T) { if _, err := pubSess.PublishNamespace(t.Context(), &message.PublishNamespace{Namespace: video}); err != nil { t.Fatalf("PublishNamespace: %v", err) } + // The relay forwards the TRACK_STATUS to the namespace's publisher. + answerTrackStatus(t, pubSess, func(r *session.Request) { _ = r.AcceptTrackStatus(nil) }) peer, conn := dialRaw(t, l) stream, err := conn.OpenStream()