From 1887c881b9c5b7bd8b2e21619e59b9ceeff6be58 Mon Sep 17 00:00:00 2001 From: Vsevolod Strukchinsky Date: Sat, 26 Sep 2026 11:12:02 +0500 Subject: [PATCH 1/2] =?UTF-8?q?fix(relay):=20use=20=C2=A711.4.3's=20same-u?= =?UTF-8?q?pstream=20and=20Prior=20Object=20ID=20Gap=20next-Object=20rules?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The relay kept an Object on an existing downstream subgroup stream only when its ID was one greater than the previous Object's, so a publisher skipping Object IDs cost a reset and a reopen per gap. §11.4.3 lists three ways to tell "the next Object"; the relay now also uses: - the Object was read next from the same upstream subgroup stream as the previously sent one, with only Objects that failed the subscriber's filters between (Forward State is one, §5.1.5); - its Prior Object ID Gap (§12.9) says the IDs between do not exist. Knowing from the cache or the Range Filters that the missing IDs are in other Subgroups or filtered out is not used (scope chosen with the user). runFanout numbers every Object read from an inbound stream, dedup losers included, since each is an Object between its neighbours. admit carries the run across filtered Objects; a queue overflow or a §12.3 expiry drop breaks it, since the dropped Object was not sent. New: message.PriorObjectIDGap (zero-alloc, shares the CheckObjectProperties walker). Tests: TestFanout_NextObject_GapOnOneUpstreamStaysOnStream, _FilteredObjectsBetween, _AcrossUpstreams (gap case) and TestRelay_ForwardStateOmissionKeepsStream were verified red on the unpatched relay; the other-upstream, duplicate-between and relay-drop cases guard against over-applying the rule and pass before and after. Each part of the change was removed in turn and a test failed each time. Tests that relied on a single-upstream gap forcing a reopen now use a second publisher (PUBLISH_DONE) or moved into the white-box drop test (FIRST_OBJECT on a reopened stream). Co-Authored-By: Claude Opus 5.5 (1M context) --- STATUS.md | 2 +- pkg/moqt/message/object_properties.go | 14 ++ pkg/moqt/message/properties_test.go | 26 ++ pkg/relay/forward_state_test.go | 38 ++- pkg/relay/handler_fanout.go | 91 ++++++- pkg/relay/handler_fanout_firstobject_test.go | 34 --- pkg/relay/handler_fanout_next_test.go | 153 +++++++++++ pkg/relay/handler_fanout_test.go | 136 ---------- pkg/relay/metrics.go | 11 +- pkg/relay/next_object_test.go | 252 +++++++++++++++++++ pkg/relay/publish_done_test.go | 20 +- 11 files changed, 572 insertions(+), 205 deletions(-) create mode 100644 pkg/relay/handler_fanout_next_test.go create mode 100644 pkg/relay/next_object_test.go diff --git a/STATUS.md b/STATUS.md index 7e91b850..ff64fde5 100644 --- a/STATUS.md +++ b/STATUS.md @@ -219,7 +219,7 @@ By package, bottom-up along the dependency stack: | 11.4 | Streams (subgroup / fetch) | DONE | Typed in/out subgroup + fetch streams. | | 11.4.1 | Stream cancellation | DONE | Bidi request-stream termination ends the request (handlers unregister on stream end); the relay sends PUBLISH_DONE on graceful subscription termination rather than abrupt reset. | | 11.4.2 | Subgroup header + delta object IDs | DONE | All subgroup-ID modes; `ReadDecoded` resolves deltas. A subgroup stream whose Track Alias is not bound yet is held unread (§11.4.2 MAY buffer): `session.Demux` parks it, count-bounded; the relay waits up to 1 s (`IncomingSubgroupStream.AwaitInboundTrack`, at most 32 streams per session, the rest reset with EXCESSIVE_LOAD). The §11.4.2 MUST to give control streams connection flow control first is not met by the bundled transports, so enough early data can delay the SUBSCRIBE_OK until the relay's wait runs out and the streams are reset. | -| 11.4.3 | Closing subgroup streams | DONE | Relay forwards only the next object on a stream (gap → reset+reopen), FINs on clean inbound EOF, resets on inbound reset, resets with MALFORMED_TRACK after a terminal EndOfGroup/EndOfTrack object (§2.4.2), marks reliable boundaries for RESET_STREAM_AT (`SetReliableBoundary`, transport-gated on `EnableStreamResetPartialDelivery`), and resets (not FINs) in-flight subgroups whose group falls out of range after a narrowing REQUEST_UPDATE. A subscription that skipped any Object of the Subgroup other than one before its Start Location (a filter, Forward State 0, a Start raised past Objects already sent, an inbox overflow with EXCESSIVE_LOAD, an expiry) gets resets, never a FIN, on that Subgroup's streams. Objects published before a subscription joined are treated as before its Start. | +| 11.4.3 | Closing subgroup streams | DONE | Relay forwards only the next object on a stream, otherwise reset+reopen: the next object is one ID greater, read next from the same upstream stream (only filtered-out objects between), or covered by its Prior Object ID Gap; an object the relay dropped, or one from another upstream, breaks the run. FINs on clean inbound EOF, resets on inbound reset, resets with MALFORMED_TRACK after a terminal EndOfGroup/EndOfTrack object (§2.4.2), marks reliable boundaries for RESET_STREAM_AT (`SetReliableBoundary`, transport-gated on `EnableStreamResetPartialDelivery`), and resets (not FINs) in-flight subgroups whose group falls out of range after a narrowing REQUEST_UPDATE. A subscription that skipped any Object of the Subgroup other than one before its Start Location (a filter, Forward State 0, a Start raised past Objects already sent, an inbox overflow with EXCESSIVE_LOAD, an expiry) gets resets, never a FIN, on that Subgroup's streams. Objects published before a subscription joined are treated as before its Start. | | 11.4.4 | Fetch header | DONE | Serialization Flags of 128 or more that are not an End of Range close the session. | | 11.4.4.1 | Fetch flags | DONE | All subgroup modes + delta/priority/properties/status flags. A first Object that references a prior Object's fields closes the session. | | 11.4.4.2 | End of range | DONE | Non-existent (0x8C) / unknown (0x10C) handled. An Object after a leading marker that references a prior Subgroup ID or Priority closes the session. | diff --git a/pkg/moqt/message/object_properties.go b/pkg/moqt/message/object_properties.go index 666eeffe..5ce173f9 100644 --- a/pkg/moqt/message/object_properties.go +++ b/pkg/moqt/message/object_properties.go @@ -76,3 +76,17 @@ func (c *objectPropertiesCheck) walk(raw []byte, nested bool) error { } return nil } + +// PriorObjectIDGap returns the Prior Object ID Gap (§12.9) in an Object's raw +// Properties and whether there is one, searching Immutable Properties too +// (§12.7). Properties that [CheckObjectProperties] rejects for a reason it can +// see without the Object's ID carry none. +// +// Must not allocate: per-Object path. +func PriorObjectIDGap(raw []byte) (uint64, bool) { + var c objectPropertiesCheck + if c.walk(raw, false) != nil || c.objectGaps != 1 { + return 0, false + } + return c.objectGap, true +} diff --git a/pkg/moqt/message/properties_test.go b/pkg/moqt/message/properties_test.go index 16e5ddd7..2ce47f63 100644 --- a/pkg/moqt/message/properties_test.go +++ b/pkg/moqt/message/properties_test.go @@ -285,6 +285,32 @@ func TestCheckObjectProperties(t *testing.T) { } } +// TestPriorObjectIDGap: the Prior Object ID Gap (§12.9) is found in either list +// (§12.7), and Properties that make the track malformed carry none. +func TestPriorObjectIDGap(t *testing.T) { + for _, tc := range []struct { + name string + raw []byte + gap uint64 + hasGap bool + }{ + {"empty", nil, 0, false}, + {"other properties", AppendTrackProperties([]wire.KVPair{kv(0x40, 1)}), 0, false}, + {"mutable", AppendTrackProperties([]wire.KVPair{kv(0x40, 1), kv(PropertyPriorObjectIDGap, 3)}), 3, true}, + {"zero", AppendTrackProperties([]wire.KVPair{kv(PropertyPriorObjectIDGap, 0)}), 0, true}, + {"inside Immutable", AppendTrackProperties([]wire.KVPair{immutable(kv(PropertyPriorObjectIDGap, 2))}), 2, true}, + {"two instances", AppendTrackProperties([]wire.KVPair{ + kv(PropertyPriorObjectIDGap, 1), immutable(kv(PropertyPriorObjectIDGap, 1)), + }), 0, false}, + {"unparseable", []byte{0x02}, 0, false}, + } { + gap, ok := PriorObjectIDGap(tc.raw) + if gap != tc.gap || ok != tc.hasGap { + t.Errorf("%s: PriorObjectIDGap = (%d, %v), want (%d, %v)", tc.name, gap, ok, tc.gap, tc.hasGap) + } + } +} + func BenchmarkCheckObjectProperties(b *testing.B) { raw := AppendTrackProperties([]wire.KVPair{ kv(0x40, 1), kv(PropertyPriorObjectIDGap, 1), immutable(kv(0x42, 7)), diff --git a/pkg/relay/forward_state_test.go b/pkg/relay/forward_state_test.go index d5d8bd5d..e4485c10 100644 --- a/pkg/relay/forward_state_test.go +++ b/pkg/relay/forward_state_test.go @@ -216,9 +216,11 @@ func TestRelay_SkipBeforeStartKeepsFIN(t *testing.T) { } } -// TestRelay_ForwardStateOmissionResetsReopenedStream: after a pause omitted an -// Object, the stream reopened on resume ends with a reset too. -func TestRelay_ForwardStateOmissionResetsReopenedStream(t *testing.T) { +// TestRelay_ForwardStateOmissionKeepsStream: an Object omitted while paused did +// not pass the subscriber's filters (§5.1.5: "The Forward parameter is also a +// type of filter"), so the Object after the resume is the next Object and +// stays on the stream (§11.4.3), which still ends with a reset. +func TestRelay_ForwardStateOmissionKeepsStream(t *testing.T) { t.Parallel() pubSess, teardown := connectRelay(t, relay.Config{}) defer teardown() @@ -240,7 +242,7 @@ func TestRelay_ForwardStateOmissionResetsReopenedStream(t *testing.T) { return } <-resumed - _ = sg.WriteObject(&message.SubgroupObject{Payload: []byte("2")}) // gap: a new stream + _ = sg.WriteObject(&message.SubgroupObject{Payload: []byte("2")}) _ = sg.Close() }() @@ -248,18 +250,25 @@ func TestRelay_ForwardStateOmissionResetsReopenedStream(t *testing.T) { if err != nil { t.Fatalf("AcceptDataStream: %v", err) } - first, ok := ds.(*session.IncomingSubgroupStream) + in, ok := ds.(*session.IncomingSubgroupStream) if !ok { t.Fatalf("AcceptDataStream = %T, want a subgroup stream", ds) } - if _, err := first.ReadObject(); err != nil { - t.Fatalf("ReadObject: %v", err) + type result struct { + ids []uint64 + end error } + done := make(chan result, 1) go func() { + var r result for { - if _, err := first.ReadObject(); err != nil { + o, err := in.ReadDecoded() + if err != nil { + r.end = err + done <- r return } + r.ids = append(r.ids, o.ObjectID) } }() if _, err := subSess.UpdateRequest( @@ -280,11 +289,16 @@ func TestRelay_ForwardStateOmissionResetsReopenedStream(t *testing.T) { } close(resumed) - ids, end := readUntilEnd(t, subSess) - if !slices.Equal(ids, []uint64{2}) { - t.Fatalf("the reopened stream carried Objects %v, want [2]", ids) + var r result + select { + case r = <-done: + case <-time.After(2 * time.Second): + t.Fatal("the subgroup stream did not end") + } + if !slices.Equal(r.ids, []uint64{0, 2}) { + t.Fatalf("the stream carried Objects %v, want [0 2]", r.ids) } - requireReset(t, end, "Object 1 was omitted while paused") + requireReset(t, r.end, "Object 1 was omitted while paused") } // TestRelay_StartRaisedToLaterGroupResetsPromptly: a Start raised past the diff --git a/pkg/relay/handler_fanout.go b/pkg/relay/handler_fanout.go index c8bd73ab..34771411 100644 --- a/pkg/relay/handler_fanout.go +++ b/pkg/relay/handler_fanout.go @@ -32,6 +32,19 @@ type fwdObject struct { // first marks the subgroup's true first object (§11.4.2 FIRST_OBJECT); // only an outbound stream beginning with it sets the bit. first bool + + // follows reports that the Object handed to the writer before this one + // was read just before it from the same inbound stream, with only + // Objects the subscriber's filters rejected between (see + // [subgroupWriter.admit]). + follows bool +} + +// inboundPos is an Object's place on its inbound subgroup stream: the stream, +// and how many Objects had been read from it, this one included. +type inboundPos struct { + src *session.IncomingSubgroupStream + seq uint64 } // subgroupWriterSet is the payload of a [registry.SharedSubgroup]: one @@ -265,6 +278,9 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming var ( firstObj = true + // pos counts every Object read, dedup losers included: each is an + // Object between its neighbours (§11.4.3). + pos = inboundPos{src: stream} // terminalSeen: an EndOfGroup/EndOfTrack was read on this inbound // stream; any later object makes the track malformed (§11.4.3, // §2.4.2). @@ -315,6 +331,7 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming isTrueFirst := firstObj && !hdr.ReplayingSubgroup firstObj = false + pos.seq++ objectID := stream.ObjectID() // resolved by ReadObject (§11.4.2) // Tracked whether or not this copy wins the dedup claim below. @@ -368,8 +385,14 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming if w == nil { continue } - if w.admit(hdr, objectID, obj.Properties) { - w.publish(fwdObject{obj: obj, absID: objectID, first: isTrueFirst, maxCacheAge: liveMaxAge}) + if take, follows := w.admit(pos, hdr, objectID, obj.Properties); take { + w.publish(fwdObject{ + obj: obj, + absID: objectID, + first: isTrueFirst, + maxCacheAge: liveMaxAge, + follows: follows, + }) } } sg.Mu.Unlock() @@ -435,7 +458,8 @@ func (h *sessionHandler) openWriterForSub( // subgroupWriter is the per-subscriber writer goroutine: it drains an inbox // onto outbound subgroup streams on the subscriber's session. // -// - A gap in Object IDs resets the stream and opens a fresh one (§11.4.3). +// - An Object that is not the next Object resets the stream and opens a +// fresh one (§11.4.3, see [isNextObject]). // - A clean inbound EOF FINs the stream; an inbound error resets it. // - A full inbox drops the object. An object that waited longer than // maxLag resets with TOO_FAR_BEHIND and terminates the subscription @@ -485,15 +509,27 @@ type subgroupWriter struct { // straggler from another upstream. Only touched by admit, under sg.Mu. lastAdmitted uint64 hasAdmitted bool + // lastPos is where the last Object admit let through was read, moved + // past the filtered Objects read after it; zero once the relay drops + // one. Only touched by admit and publish, under sg.Mu. + lastPos inboundPos } // admit decides whether w takes the Object at objectID of the subgroup hdr -// names, closing w when it will take none again. -func (w *subgroupWriter) admit(hdr message.SubgroupHeader, objectID uint64, props []byte) bool { +// names, read at pos, closing w when it will take none again. follows is +// [fwdObject.follows] for a taken Object. +func (w *subgroupWriter) admit( + pos inboundPos, + hdr message.SubgroupHeader, + objectID uint64, + props []byte, +) (take, follows bool) { + follows = w.lastPos.src == pos.src && w.lastPos.seq+1 == pos.seq switch w.sub.ForwardDecision(hdr.GroupID, objectID, hdr.SubgroupID, hdr.PublisherPriority, props) { case registry.Forward: w.lastAdmitted, w.hasAdmitted = objectID, true - return true + w.lastPos = pos + return true, follows case registry.SkipObject, registry.SkipPaused: // The stream stays open for later Objects, but the Subgroup is now // incomplete (§11.4.3). @@ -509,7 +545,13 @@ func (w *subgroupWriter) admit(hdr message.SubgroupHeader, objectID uint64, prop w.markIncomplete(moqt.StreamResetCancelled) } } - return false + // §11.4.3: an Object that "did not pass the subscriber's filters" does + // not separate the ones either side of it. Forward State is one of them + // (§5.1.5). + if follows { + w.lastPos = pos + } + return false, false } // markIncomplete sets subgroupWriter.incomplete; the first code recorded @@ -541,7 +583,8 @@ func (w *subgroupWriter) resetCode() moqt.StreamResetCode { // publish enqueues fwd without blocking, stamping its enqueue time for the // lag check. On overflow the object is dropped, and past maxDropsBeforeReset -// the writer is closed in reset mode. It is a no-op after close. +// the writer is closed in reset mode. It is a no-op after close. Callers hold +// sg.Mu. func (w *subgroupWriter) publish(fwd fwdObject) { w.dropsMu.Lock() if w.closed { @@ -556,6 +599,7 @@ func (w *subgroupWriter) publish(fwd fwdObject) { w.metrics.ObjectForwarded(w.ref, w.hdr.SubgroupID) default: w.metrics.ObjectDropped(w.ref, w.hdr.SubgroupID) + w.lastPos = inboundPos{} // the next Object does not follow a sent one w.dropsMu.Lock() w.drops++ w.markIncompleteLocked(moqt.StreamResetExcessiveLoad) @@ -634,6 +678,9 @@ func (w *subgroupWriter) run() { prevID uint64 hasWritten bool writeFailed bool + // dropped: an Object was dropped since the last one written, so the + // next is not known to follow it. + dropped bool ) // reopen resets the current outbound stream (if any) and opens a fresh @@ -712,10 +759,8 @@ func (w *subgroupWriter) run() { } } - // §11.4.3: only "the next Object" may go on an existing stream. Of - // the draft's three ways to tell, this relay uses only "one greater - // than the previous Object" (a choice) and reopens on any other gap. - if hasWritten && fwd.absID != prevID+1 { + // §11.4.3: only "the next Object" may go on an existing stream. + if hasWritten && !isNextObject(fwd, prevID, dropped) { w.metrics.SubgroupStreamReset(w.ref, w.hdr.SubgroupID, ResetCauseGap) if !reopen(fwd.first) { failWrites() @@ -727,6 +772,7 @@ func (w *subgroupWriter) run() { // the open above, which can block. if expired(fwd) { w.dropExpired(hasWritten) + dropped = true continue } @@ -761,6 +807,7 @@ func (w *subgroupWriter) run() { } prevID = fwd.absID hasWritten = true + dropped = false // §11.4.3: a later reset still delivers what was written. w.out.MarkReliable() } @@ -823,6 +870,26 @@ func (w *subgroupWriter) run() { w.closeOut(true, 0) } +// isNextObject reports whether fwd is "the next Object" (§11.4.3) on a stream +// whose last Object is prevID. Of the draft's ways to tell, the relay uses: +// the Object ID is one greater; fwd follows the last Object on its inbound +// stream; or its Prior Object ID Gap (§12.9) says the IDs between do not +// exist. Knowing from the cache or the subscriber's filters that they are in +// other Subgroups or filtered out is not used (a choice). +func isNextObject(fwd fwdObject, prevID uint64, dropped bool) bool { + if fwd.absID == prevID+1 { + return true + } + if dropped || fwd.absID <= prevID { + return false + } + if fwd.follows { + return true + } + gap, ok := message.PriorObjectIDGap(fwd.obj.Properties) + return ok && fwd.absID-gap <= prevID+1 +} + // openCounted opens a subgroup stream, counting it for the §10.12 Stream // Count; it fails once the subscription has terminated. func (w *subgroupWriter) openCounted(hdr message.SubgroupHeader) (*session.OutgoingSubgroupStream, error) { diff --git a/pkg/relay/handler_fanout_firstobject_test.go b/pkg/relay/handler_fanout_firstobject_test.go index 4836e753..778981dc 100644 --- a/pkg/relay/handler_fanout_firstobject_test.go +++ b/pkg/relay/handler_fanout_firstobject_test.go @@ -165,40 +165,6 @@ func TestFanout_FirstObjectBitClearedForFilteredHead(t *testing.T) { } } -// TestFanout_FirstObjectBitAcrossGapReopen: after a gap reopen (§11.4.3) only -// the original stream carries FIRST_OBJECT. -func TestFanout_FirstObjectBitAcrossGapReopen(t *testing.T) { - t.Parallel() - pub, sub := firstObjectTopology(t, nil) - // Objects 0, 1, then a jump to 5: the relay must not forward the - // non-consecutive object on the same outbound stream. - writeSubgroupObjects(t, pub, message.SubgroupHeader{ - SubgroupIDMode: message.SubgroupIDImplicitZero, TrackAlias: 7, GroupID: 0, - }, []uint64{0, 1, 5}) - - caps := captureSubgroups(t, sub, 2) - var origin, reopened *subgroupCapture - for i := range caps { - if len(caps[i].Objects) > 0 && caps[i].Objects[0] == 0 { - origin = &caps[i] - } else { - reopened = &caps[i] - } - } - if origin == nil || reopened == nil { - t.Fatalf("expected an origin and a reopened stream, got %+v", caps) - } - if origin.Header.ReplayingSubgroup { - t.Error("origin stream begins with the subgroup's first object: FIRST_OBJECT must be set") - } - if !reopened.Header.ReplayingSubgroup { - t.Error("gap-reopened stream is mid-subgroup: FIRST_OBJECT must be clear") - } - if len(reopened.Objects) != 1 || reopened.Objects[0] != 5 { - t.Errorf("reopened stream objects = %v, want [5]", reopened.Objects) - } -} - // TestFanout_FirstObjectBitNotInvented: an inbound replay stream (FIRST_OBJECT // clear) is not forwarded with the bit set. func TestFanout_FirstObjectBitNotInvented(t *testing.T) { diff --git a/pkg/relay/handler_fanout_next_test.go b/pkg/relay/handler_fanout_next_test.go new file mode 100644 index 00000000..ca83f998 --- /dev/null +++ b/pkg/relay/handler_fanout_next_test.go @@ -0,0 +1,153 @@ +package relay + +import ( + "context" + "errors" + "io" + "slices" + "testing" + "time" + + "github.com/floatdrop/moq-go/pkg/moqt/message" + "github.com/floatdrop/moq-go/pkg/moqt/session" + "github.com/floatdrop/moq-go/pkg/moqt/session/sessiontest" +) + +// forwarded is one Object as the subscriber received it, or with end set, the +// end of the stream it came on. +type forwarded struct { + stream int // 1-based index of the subgroup stream it came on + id uint64 + replay bool // the stream's header had FIRST_OBJECT clear (§11.4.2) + end error // io.EOF for a FIN +} + +// readForwarded emits every Object and stream end of every subgroup stream srv +// accepts. +func readForwarded(ctx context.Context, srv *session.Session) <-chan forwarded { + out := make(chan forwarded, 16) + go func() { + for stream := 1; ; stream++ { + ds, err := srv.AcceptDataStream(ctx) + if err != nil { + return + } + sg, ok := ds.(*session.IncomingSubgroupStream) + if !ok { + return + } + for { + if _, err := sg.ReadObject(); err != nil { + out <- forwarded{stream: stream, end: err} + break + } + out <- forwarded{stream: stream, id: sg.ObjectID(), replay: sg.Header.ReplayingSubgroup} + } + } + }() + return out +} + +// awaitForwarded waits for the next Object or stream end the subscriber +// receives. +func awaitForwarded(t *testing.T, in <-chan forwarded) forwarded { + t.Helper() + select { + case f := <-in: + return f + case <-time.After(2 * time.Second): + t.Fatal("nothing forwarded") + return forwarded{} + } +} + +// TestSubgroupWriter_RelayDropBreaksRun: an Object the relay dropped was not +// sent, so the upstream's next Object is not "the next Object" after the one +// before it (§11.4.3: "with no other Objects in between"). It goes on a new +// stream, without FIRST_OBJECT (§11.4.2), and both streams end with a reset, +// since an Object was omitted. +func TestSubgroupWriter_RelayDropBreaksRun(t *testing.T) { + t.Parallel() + src := new(session.IncomingSubgroupStream) // identity only + for _, tc := range []struct { + name string + // Objects 0, 1 and 2 are queued before the writer starts; 2 overflows + // a queue of 2, and expires with a maxAge. + queue int + maxAge time.Duration + want []forwarded // Objects then stream ends, in order + }{ + {"nothing dropped", 3, 0, []forwarded{ + {stream: 1, id: 0}, {stream: 1, id: 1}, {stream: 1, id: 2}, {stream: 1, id: 4}, + {stream: 1, end: io.EOF}, + }}, + {"queue overflow", 2, 0, []forwarded{ + {stream: 1, id: 0}, {stream: 1, id: 1}, {stream: 1, end: errReset}, + {stream: 2, id: 4, replay: true}, {stream: 2, end: errReset}, + }}, + {"expired", 3, time.Nanosecond, []forwarded{ + {stream: 1, id: 0}, {stream: 1, id: 1}, {stream: 1, end: errReset}, + {stream: 2, id: 4, replay: true}, {stream: 2, end: errReset}, + }}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + cli, srv := sessiontest.NewSessionPair(t) + got := readForwarded(t.Context(), srv) + w := newWedgeableWriter(t, cli) + w.inbox = make(chan fwdObject, tc.queue) + offer := func(seq, id uint64, maxAge time.Duration) { + take, follows := w.admit(inboundPos{src: src, seq: seq}, w.hdr, id, nil) + if take { + w.publish(fwdObject{ + obj: &message.SubgroupObject{Payload: []byte("x")}, + absID: id, + first: id == 0, + follows: follows, + maxCacheAge: maxAge, + }) + } + } + + // Upstream Objects 0, 1, 2 and 4, read in turn. + offer(1, 0, 0) + offer(2, 1, 0) + offer(3, 2, tc.maxAge) + go w.run() + var seen []forwarded + for { + f := awaitForwarded(t, got) + seen = append(seen, f) + if f.end == nil && f.id == 1 { + break + } + } + offer(4, 4, 0) + for { + f := awaitForwarded(t, got) + seen = append(seen, f) + if f.end == nil && f.id == 4 { + break + } + } + w.close(false, 0) + joinOrFatal(t, w) + seen = append(seen, awaitForwarded(t, got)) + + for i := range seen { + switch end := seen[i].end; { + case errors.Is(end, io.EOF): + seen[i].end = io.EOF + case end != nil: + seen[i].end = errReset + } + } + if !slices.Equal(seen, tc.want) { + t.Fatalf("forwarded %v,\nwant %v", seen, tc.want) + } + }) + } +} + +// errReset stands for any stream reset in [TestSubgroupWriter_RelayDropBreaksRun]. +var errReset = errors.New("reset") diff --git a/pkg/relay/handler_fanout_test.go b/pkg/relay/handler_fanout_test.go index d6f54c51..14a76cf0 100644 --- a/pkg/relay/handler_fanout_test.go +++ b/pkg/relay/handler_fanout_test.go @@ -608,142 +608,6 @@ func TestSubscribe_InstallsPriorityAndGroupOrder(t *testing.T) { // covered by the unit test on DownstreamSub setters. } -// TestFanout_GapInForwardedObjectIDsOpensNewStream: a gap in the Object IDs of -// one inbound subgroup resets the outbound stream and opens a new one for the -// next Object (§11.4.3). -func TestFanout_GapInForwardedObjectIDsOpensNewStream(t *testing.T) { - t.Parallel() - - pubSess, teardown := connectRelay(t, relay.Config{}) - defer teardown() - - const publisherAlias = uint64(7) - pubReqStream, err := pubSess.Publish(t.Context(), &message.Publish{ - Namespace: ns("video"), - Name: []byte("cam1"), - TrackAlias: publisherAlias, - }) - if err != nil { - t.Fatalf("Publish: %v", err) - } - defer pubReqStream.Close() - - subSess := dialAnotherClient(t, pubSess) - subReq, err := subSess.Subscribe(t.Context(), &message.Subscribe{ - Namespace: ns("video"), - Name: []byte("cam1"), - }) - if err != nil { - t.Fatalf("Subscribe: %v", err) - } - defer subReq.Close() - - // The subscriber-side reader collects every stream the relay opens - // for this subgroup, recording the (firstAbsID, objectCount, endErr) - // tuple for each. The relay should produce exactly two streams: one - // for absID=0 (then reset), and one for absID=2 (then FIN). - type streamSummary struct { - firstAbsID uint64 - count int - endErr error - } - streams := make(chan streamSummary, 4) - go func() { - for { - ds, err := subSess.AcceptDataStream(t.Context()) - if err != nil { - close(streams) - return - } - sg, ok := ds.(*session.IncomingSubgroupStream) - if !ok { - close(streams) - return - } - obj, err := sg.ReadObject() - if err != nil { - streams <- streamSummary{endErr: err} - continue - } - first := obj.ObjectIDDelta - count := 1 - for { - _, err := sg.ReadObject() - if err != nil { - streams <- streamSummary{firstAbsID: first, count: count, endErr: err} - break - } - count++ - } - } - }() - - pubSubgroup, err := pubSess.OpenSubgroup(message.SubgroupHeader{ - SubgroupIDMode: message.SubgroupIDExplicit, - TrackAlias: publisherAlias, - GroupID: 0, - SubgroupID: 0, - }) - if err != nil { - t.Fatalf("OpenSubgroup: %v", err) - } - - // Object with absolute ID 0. - if err := pubSubgroup.WriteObject(&message.SubgroupObject{ - ObjectIDDelta: 0, - Payload: []byte("first"), - }); err != nil { - t.Fatalf("WriteObject #0: %v", err) - } - // Object with absolute ID 2 — delta is (2 - 0 - 1) = 1. Skips ID 1. - if err := pubSubgroup.WriteObject(&message.SubgroupObject{ - ObjectIDDelta: 1, - Payload: []byte("third"), - }); err != nil { - t.Fatalf("WriteObject #2: %v", err) - } - if err := pubSubgroup.Close(); err != nil { - t.Fatalf("pubSubgroup.Close: %v", err) - } - - // Collect the first two summaries — that's all the relay should send. - var got []streamSummary - deadline := time.After(2 * time.Second) - for len(got) < 2 { - select { - case s, ok := <-streams: - if !ok { - t.Fatalf("subscriber stream channel closed after %d summaries: %v", len(got), got) - } - got = append(got, s) - case <-deadline: - t.Fatalf("subscriber did not receive two streams within deadline (got %d: %v)", len(got), got) - } - } - - // First stream: pre-gap object (absID=0), then reset. - if got[0].firstAbsID != 0 { - t.Errorf("stream 1 firstAbsID = %d, want 0", got[0].firstAbsID) - } - if got[0].count != 1 { - t.Errorf("stream 1 count = %d, want 1", got[0].count) - } - if errors.Is(got[0].endErr, io.EOF) { - t.Errorf("stream 1 ended with io.EOF, want a reset") - } - - // Second stream: post-gap object (absID=2), then clean FIN. - if got[1].firstAbsID != 2 { - t.Errorf("stream 2 firstAbsID = %d, want 2", got[1].firstAbsID) - } - if got[1].count != 1 { - t.Errorf("stream 2 count = %d, want 1", got[1].count) - } - if !errors.Is(got[1].endErr, io.EOF) { - t.Errorf("stream 2 ended with %v, want io.EOF (clean FIN)", got[1].endErr) - } -} - // TestFanout_InboundResetCancelsDownstream: a reset inbound subgroup stream // resets the downstream one rather than FINning it (§11.4.3). func TestFanout_InboundResetCancelsDownstream(t *testing.T) { diff --git a/pkg/relay/metrics.go b/pkg/relay/metrics.go index fa5394bb..d6f45ac3 100644 --- a/pkg/relay/metrics.go +++ b/pkg/relay/metrics.go @@ -71,12 +71,11 @@ type TrackRef struct { type ResetCause uint8 const ( - // ResetCauseGap is a §11.4.3 reopen: the next object to forward was not - // consecutive with the last one written, so the current outbound stream - // was reset and a fresh one opened. The relay MUST NOT forward a - // non-consecutive object on an existing subgroup stream, so this is - // correct behaviour — but it is also the direct consequence of an - // earlier drop or filter narrowing, and a subscriber sees the hole. + // ResetCauseGap is a §11.4.3 reopen: the object to forward was not known + // to be "the next Object" after the last one written, so the current + // outbound stream was reset and a fresh one opened. That is correct + // behaviour, but it usually follows an earlier drop, or Objects arriving + // from several upstreams, and a subscriber sees the hole. ResetCauseGap ResetCause = iota // ResetCauseDeliveryTimeout is §8: an object sat unsent past the diff --git a/pkg/relay/next_object_test.go b/pkg/relay/next_object_test.go new file mode 100644 index 00000000..1da7623e --- /dev/null +++ b/pkg/relay/next_object_test.go @@ -0,0 +1,252 @@ +package relay_test + +import ( + "errors" + "io" + "slices" + "testing" + "time" + + "github.com/floatdrop/moq-go/pkg/moqt/message" + "github.com/floatdrop/moq-go/pkg/moqt/session" + "github.com/floatdrop/moq-go/pkg/moqt/wire" + "github.com/floatdrop/moq-go/pkg/relay" +) + +// §11.4.3: the relay forwards an Object on an existing subgroup stream only if +// it is "the next Object", and otherwise resets the stream and opens another. + +// collectStreams reads events until objects Objects have arrived and every +// stream they came on has ended. +func collectStreams(t *testing.T, events <-chan objEvent, objects int) []objEvent { + t.Helper() + var ( + seen []objEvent + got int + streams int + ended int + ) + deadline := time.After(3 * time.Second) + for got < objects || ended < streams { + select { + case ev := <-events: + if ev.stream == 0 { + t.Fatalf("AcceptDataStream: %v", ev.err) + } + seen = append(seen, ev) + streams = max(streams, ev.stream) + if ev.err != nil { + ended++ + } else { + got++ + } + case <-deadline: + t.Fatalf("got %d/%d Objects, %d/%d streams ended: %v", got, objects, ended, streams, layout(seen)) + } + } + return seen +} + +// finned reports whether stream ended with a FIN in events. +func finned(events []objEvent, stream int) bool { + for _, ev := range events { + if ev.stream == stream && ev.err != nil { + return errors.Is(ev.err, io.EOF) + } + } + return false +} + +// TestFanout_NextObject_GapOnOneUpstreamStaysOnStream: Objects read in turn +// from one upstream stream are each the next Object (§11.4.3, second rule), +// so a publisher skipping IDs costs no reset. +func TestFanout_NextObject_GapOnOneUpstreamStaysOnStream(t *testing.T) { + t.Parallel() + pubSess, teardown := connectRelay(t, relay.Config{}) + defer teardown() + pub := publishVideoTrack(t, pubSess, "cam1", 1) + subSess := dialAnotherClient(t, pubSess) + subscribeCam1(t, subSess) + events := make(chan objEvent, 16) + go readSubgroups(t.Context(), subSess, events) + + sg, err := pub.OpenSubgroup(message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit}) + if err != nil { + t.Fatalf("OpenSubgroup: %v", err) + } + for _, id := range []uint64{0, 2, 5} { + if err := sg.WriteObjectAt(id, &message.SubgroupObject{Payload: []byte("x")}); err != nil { + t.Fatalf("WriteObjectAt %d: %v", id, err) + } + } + if err := sg.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + + seen := collectStreams(t, events, 3) + if got, want := layout(seen), [][]uint64{{0, 2, 5}}; !slices.EqualFunc(got, want, slices.Equal) { + t.Fatalf("streams = %v, want %v", got, want) + } + if !finned(seen, 1) { + t.Error("stream ended with a reset, want a FIN: no Object was omitted") + } +} + +// TestFanout_NextObject_FilteredObjectsBetween: Objects between two forwarded +// ones that did not pass the subscriber's filters keep the second the next +// Object (§11.4.3). +func TestFanout_NextObject_FilteredObjectsBetween(t *testing.T) { + t.Parallel() + pubSess, teardown := connectRelay(t, relay.Config{}) + defer teardown() + pub := publishVideoTrack(t, pubSess, "cam1", 1) + subSess := dialAnotherClient(t, pubSess) + subscribeCam1(t, subSess, message.RangeFilterParam(&message.RangeFilter{ + Type: message.ParamObjectIDFilter, + Ranges: []message.Range{{Start: 0, End: 0}, {Start: 2, End: 3}}, + })) + events := make(chan objEvent, 16) + go readSubgroups(t.Context(), subSess, events) + + sg, err := pub.OpenSubgroup(message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit}) + if err != nil { + t.Fatalf("OpenSubgroup: %v", err) + } + for range 4 { + if err := sg.WriteObject(&message.SubgroupObject{Payload: []byte("x")}); err != nil { + t.Fatalf("WriteObject: %v", err) + } + } + if err := sg.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + + seen := collectStreams(t, events, 3) + if got, want := layout(seen), [][]uint64{{0, 2, 3}}; !slices.EqualFunc(got, want, slices.Equal) { + t.Fatalf("streams = %v, want %v", got, want) + } + if finned(seen, 1) { + t.Error("stream ended with a FIN, want a reset: Object 1 was omitted") + } +} + +// twoPublishers PUBLISHes video/cam1 from two sessions, subscribes a third, and +// opens Subgroup 0 of Group 0, with Object Properties, from each publisher. +func twoPublishers(t *testing.T) (a, b *session.OutgoingSubgroupStream, events <-chan objEvent) { + t.Helper() + pubA, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + pubB := dialAnotherClient(t, pubA) + subSess := dialAnotherClient(t, pubA) + aPub := publishVideoTrack(t, pubA, "cam1", 1) + bPub := publishVideoTrack(t, pubB, "cam1", 2) + subscribeCam1(t, subSess) + ch := make(chan objEvent, 16) + go readSubgroups(t.Context(), subSess, ch) + + // Properties on both: the merged stream takes the first one's header. + hdr := message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit, Properties: true} + a, err := aPub.OpenSubgroup(hdr) + if err != nil { + t.Fatalf("A OpenSubgroup: %v", err) + } + b, err = bPub.OpenSubgroup(hdr) + if err != nil { + t.Fatalf("B OpenSubgroup: %v", err) + } + return a, b, ch +} + +// writeAndAwait writes the Object at id on sg and waits for the subscriber to +// receive it, so the next write is ordered after it at the relay. +func writeAndAwait( + t *testing.T, + sg *session.OutgoingSubgroupStream, + id uint64, + props []byte, + events <-chan objEvent, + streams *[]objEvent, +) { + t.Helper() + if err := sg.WriteObjectAt(id, &message.SubgroupObject{Properties: props, Payload: []byte("x")}); err != nil { + t.Fatalf("WriteObjectAt %d: %v", id, err) + } + for { + select { + case ev := <-events: + *streams = append(*streams, ev) + if ev.err == nil && ev.absID == id { + return + } + case <-time.After(2 * time.Second): + t.Fatalf("Object %d not forwarded", id) + } + } +} + +// layout is the Object IDs of each stream in events. +func layout(events []objEvent) [][]uint64 { + var out [][]uint64 + for _, ev := range events { + if ev.err != nil { + continue + } + for len(out) < ev.stream { + out = append(out, nil) + } + out[ev.stream-1] = append(out[ev.stream-1], ev.absID) + } + return out +} + +func priorObjectIDGap(gap uint64) []byte { + return message.AppendTrackProperties([]wire.KVPair{{Type: message.PropertyPriorObjectIDGap, IntVal: gap}}) +} + +// TestFanout_NextObject_AcrossUpstreams: with two upstreams feeding one +// Subgroup (§9.3), an Object from the other upstream after a gap is the next +// Object only if its Prior Object ID Gap (§12.9) covers the gap (§11.4.3, +// third rule). +func TestFanout_NextObject_AcrossUpstreams(t *testing.T) { + t.Parallel() + for _, tc := range []struct { + name string + props []byte + want [][]uint64 + }{ + {"no gap property", nil, [][]uint64{{0}, {3}}}, + {"gap short of the previous Object", priorObjectIDGap(1), [][]uint64{{0}, {3}}}, + {"gap reaching the previous Object", priorObjectIDGap(2), [][]uint64{{0, 3}}}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + a, b, events := twoPublishers(t) + var seen []objEvent + writeAndAwait(t, a, 0, nil, events, &seen) + writeAndAwait(t, b, 3, tc.props, events, &seen) + if got := layout(seen); !slices.EqualFunc(got, tc.want, slices.Equal) { + t.Fatalf("streams = %v, want %v", got, tc.want) + } + }) + } +} + +// TestFanout_NextObject_DuplicateBetweenBreaksRun: an upstream's Object that +// lost the §2.1 dedup claim was not sent on the downstream stream, so the +// upstream's following Object is not the next Object after its earlier one +// (§11.4.3: "with no other Objects in between"). +func TestFanout_NextObject_DuplicateBetweenBreaksRun(t *testing.T) { + t.Parallel() + a, b, events := twoPublishers(t) + var seen []objEvent + writeAndAwait(t, a, 1, nil, events, &seen) + writeAndAwait(t, b, 0, nil, events, &seen) + // B's 1 duplicates A's and is dropped. + if err := b.WriteObjectAt(1, &message.SubgroupObject{Payload: []byte("x")}); err != nil { + t.Fatalf("WriteObjectAt 1: %v", err) + } + writeAndAwait(t, b, 3, nil, events, &seen) + if got, want := layout(seen), [][]uint64{{1}, {0}, {3}}; !slices.EqualFunc(got, want, slices.Equal) { + t.Fatalf("streams = %v, want %v", got, want) + } +} diff --git a/pkg/relay/publish_done_test.go b/pkg/relay/publish_done_test.go index 91497c2b..257404e9 100644 --- a/pkg/relay/publish_done_test.go +++ b/pkg/relay/publish_done_test.go @@ -264,20 +264,28 @@ func TestPublishDone_AfterGapReopen(t *testing.T) { pubSess, teardown := connectRelay(t, relay.Config{}) defer teardown() pub := publishVideoTrack(t, pubSess, "cam1", 1) + other := publishVideoTrack(t, dialAnotherClient(t, pubSess), "cam1", 2) subSess := dialAnotherClient(t, pubSess) subReq := subscribeCam1(t, subSess) - sg, err := pub.OpenSubgroup(message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit}) + hdr := message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit} + sg, err := pub.OpenSubgroup(hdr) if err != nil { t.Fatalf("OpenSubgroup: %v", err) } + otherSg, err := other.OpenSubgroup(hdr) + if err != nil { + t.Fatalf("other OpenSubgroup: %v", err) + } wrote := make(chan struct{}) go func() { defer close(wrote) _ = sg.WriteObject(&message.SubgroupObject{Payload: []byte("0")}) - // Object 2 after object 0: the relay may not carry it on the same - // stream (§11.4.3), so it resets that one and opens another. - _ = sg.WriteObject(&message.SubgroupObject{ObjectIDDelta: 1, Payload: []byte("2")}) + // Objects 0 and 2 from two upstreams: in either order, neither is the + // next Object after the other (§11.4.3), so the relay resets the + // first stream and opens another. + _ = otherSg.WriteObjectAt(2, &message.SubgroupObject{Payload: []byte("2")}) + _ = otherSg.Close() }() ended := make(chan struct{}, 2) for range 2 { @@ -299,6 +307,10 @@ func TestPublishDone_AfterGapReopen(t *testing.T) { }() } + // PUBLISH_DONE goes out once both upstreams are done. + if err := other.Done(moqt.PublishDoneTrackEnded, "done"); err != nil { + t.Fatalf("other Done: %v", err) + } if err := pub.Done(moqt.PublishDoneTrackEnded, "done"); err != nil { t.Fatalf("Done: %v", err) } From 3e44dc5a81446e1ae77319de829f834833e95e5b Mon Sep 17 00:00:00 2001 From: Vsevolod Strukchinsky Date: Sat, 26 Sep 2026 11:14:19 +0500 Subject: [PATCH 2/2] fix(relay): trust a Prior Object ID Gap only when it ends at the previous Object MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review follow-up. A gap reaching past the previous Object covers an Object the relay already received, which §12.9 makes a malformed track, so it cannot show the Object is the next one (§11.4.3). The Object now goes on a new stream. Ending the track for it (§2.4.2) is left out: like the other §12.9 conditions that need earlier Objects, it is not detected. Also asserts message.PriorObjectIDGap does not allocate. The "gap covering the previous Object" case was verified red before the fix. Co-Authored-By: Claude Opus 5.5 (1M context) --- pkg/moqt/message/properties_test.go | 4 ++++ pkg/relay/handler_fanout.go | 7 +++++-- pkg/relay/next_object_test.go | 3 +++ 3 files changed, 12 insertions(+), 2 deletions(-) diff --git a/pkg/moqt/message/properties_test.go b/pkg/moqt/message/properties_test.go index 2ce47f63..5b3f8219 100644 --- a/pkg/moqt/message/properties_test.go +++ b/pkg/moqt/message/properties_test.go @@ -309,6 +309,10 @@ func TestPriorObjectIDGap(t *testing.T) { t.Errorf("%s: PriorObjectIDGap = (%d, %v), want (%d, %v)", tc.name, gap, ok, tc.gap, tc.hasGap) } } + raw := AppendTrackProperties([]wire.KVPair{immutable(kv(PropertyPriorObjectIDGap, 2))}) + if n := testing.AllocsPerRun(10, func() { PriorObjectIDGap(raw) }); n != 0 { + t.Errorf("PriorObjectIDGap allocates %v times, want 0", n) + } } func BenchmarkCheckObjectProperties(b *testing.B) { diff --git a/pkg/relay/handler_fanout.go b/pkg/relay/handler_fanout.go index 34771411..65648356 100644 --- a/pkg/relay/handler_fanout.go +++ b/pkg/relay/handler_fanout.go @@ -875,7 +875,10 @@ func (w *subgroupWriter) run() { // the Object ID is one greater; fwd follows the last Object on its inbound // stream; or its Prior Object ID Gap (§12.9) says the IDs between do not // exist. Knowing from the cache or the subscriber's filters that they are in -// other Subgroups or filtered out is not used (a choice). +// other Subgroups or filtered out is not used (a choice). A gap that also +// covers prevID is not trusted. It is a §12.9 malformed-track condition, but +// like the others that need earlier Objects it is not detected (see +// [message.CheckObjectProperties]); the Object just goes on a new stream. func isNextObject(fwd fwdObject, prevID uint64, dropped bool) bool { if fwd.absID == prevID+1 { return true @@ -887,7 +890,7 @@ func isNextObject(fwd fwdObject, prevID uint64, dropped bool) bool { return true } gap, ok := message.PriorObjectIDGap(fwd.obj.Properties) - return ok && fwd.absID-gap <= prevID+1 + return ok && fwd.absID-gap == prevID+1 } // openCounted opens a subgroup stream, counting it for the §10.12 Stream diff --git a/pkg/relay/next_object_test.go b/pkg/relay/next_object_test.go index 1da7623e..a837f57e 100644 --- a/pkg/relay/next_object_test.go +++ b/pkg/relay/next_object_test.go @@ -217,6 +217,9 @@ func TestFanout_NextObject_AcrossUpstreams(t *testing.T) { {"no gap property", nil, [][]uint64{{0}, {3}}}, {"gap short of the previous Object", priorObjectIDGap(1), [][]uint64{{0}, {3}}}, {"gap reaching the previous Object", priorObjectIDGap(2), [][]uint64{{0, 3}}}, + // A gap covering an Object received before is not trusted. §12.9 + // makes that a malformed track, which the relay does not detect. + {"gap covering the previous Object", priorObjectIDGap(3), [][]uint64{{0}, {3}}}, } { t.Run(tc.name, func(t *testing.T) { t.Parallel()