From 2079b8d38d018c4489e4f98f40b87e688123a556 Mon Sep 17 00:00:00 2001 From: Vsevolod Strukchinsky Date: Sat, 26 Sep 2026 16:12:07 +0500 Subject: [PATCH] fix(relay): remember a Subgroup's lowest forwarded Object across contributors MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit #103 honours a contributor's FIRST_OBJECT claim (§11.4.2) only if no lower Object ID was forwarded in the Subgroup (§2.2), but kept that in the Subgroup's writer set, which is torn down when its last contributor leaves: a later contributor claiming FIRST_OBJECT above an Object already forwarded was then honoured. The dedup ledger now records the lowest Object ID forwarded per Subgroup (status Objects included, datagrams not), and a new writer set starts from it (TrackEntry.LowestForwarded), within the ledger's 32-Group window. It reads the ledger before taking sg.Mu and keeps the minimum, since a contributor that joined meanwhile may already have forwarded a lower Object; that race is not reproducible in a test and was found in review. TestFanout_MultiPublisher_FirstObjectAfterTeardown and TestTrackEntry_LowestForwarded were verified red first. Co-Authored-By: Claude Opus 5.5 (1M context) --- STATUS.md | 2 +- pkg/relay/handler_fanout.go | 11 +++- pkg/relay/handler_fanout_multipub_test.go | 74 ++++++++++++++++++++++ pkg/relay/internal/registry/ledger_test.go | 21 ++++++ pkg/relay/internal/registry/track_entry.go | 24 +++++++ 5 files changed, 129 insertions(+), 3 deletions(-) diff --git a/STATUS.md b/STATUS.md index ff9cd496..d325b5fc 100644 --- a/STATUS.md +++ b/STATUS.md @@ -144,7 +144,7 @@ By package, bottom-up along the dependency stack: |-------|--------------------------------------|--------|-------| | 9.1 | Caching relays | DONE | LRU+TTL object cache (`cache/cache.go`); updates limited to non-existence/properties. A duplicate of a cached Object with a different Forwarding Preference, Subgroup ID, Priority or Payload, or different Immutable Properties (§2.4.2, §12.7), ends the track as malformed (`relay/handler_duplicate.go`); not checked once the first copy left the cache — see Limitations. | | 9.2 | Forward handling | DONE | FORWARD flag honoured; Forward=0 pauses delivery. Upstream Forward is set to 1 only when a downstream subscriber forwards, else the relay pauses it (Forward=0) and resumes on the first forwarding subscriber. | -| 9.3 | Multiple publishers | DONE | Per-track upstreams; dedup by `{GroupID, ObjectID}`. Upstreams of one Subgroup share one downstream stream per subscriber, with the first one's SUBGROUP_HEADER; a later one's Object Properties reopen it with PROPERTIES set, so none are dropped (§2.5). Like a §11.4.3 gap reopen, the reset keeps already-written Objects only where RESET_STREAM_AT is in use; always setting PROPERTIES would avoid it at a byte per Object. | +| 9.3 | Multiple publishers | DONE | Per-track upstreams; dedup by `{GroupID, ObjectID}`. Upstreams of one Subgroup share one downstream stream per subscriber, with the first one's SUBGROUP_HEADER; a later one's Object Properties reopen it with PROPERTIES set, so none are dropped (§2.5). Like a §11.4.3 gap reopen, the reset keeps already-written Objects only where RESET_STREAM_AT is in use; always setting PROPERTIES would avoid it at a byte per Object. A stream sets FIRST_OBJECT only if it begins below every Object forwarded in its Subgroup (§2.2), remembered for the last 32 Groups across contributors. | | 9.4 | Subscriber interactions | DONE | Upstream subscription established before SUBSCRIBE_OK; aggregation. | | 9.4.1 | Graceful subscriber switchover | DONE | GOAWAY grace period (`GoawayTimeout`). | | 9.5 | Publisher interactions | DONE | PUBLISH_NAMESPACE / PUBLISH with prefix matching (`namespace_registry.go`). | diff --git a/pkg/relay/handler_fanout.go b/pkg/relay/handler_fanout.go index 359dd4ee..0e8a02f5 100644 --- a/pkg/relay/handler_fanout.go +++ b/pkg/relay/handler_fanout.go @@ -66,8 +66,9 @@ type subgroupWriterSet struct { // lowest is the lowest Object ID forwarded, if forwarded. Objects are // published in ascending order (§2.2), so a FIRST_OBJECT claim (§11.4.2) // for a higher ID is wrong, whether or not a given subscriber got the - // lower one. It lives only as long as the set: after every contributor - // left, a new one's claim is not checked against earlier Objects. + // lower one. A new set starts from the ledger's + // ([registry.TrackEntry.LowestForwarded]), so it holds across contributors + // within the ledger's window. lowest uint64 forwarded bool @@ -247,8 +248,14 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming // double-open. The stream is drained even with no subscribers (§9.7). initialSubs, gen := entry.CopyDownstreamWithGen() pubTimeouts := entry.DeliveryTimeouts() + lowest, forwarded := entry.LowestForwarded(hdr.GroupID, hdr.SubgroupID) sg.Mu.Lock() set.gen = gen + // The minimum, not an assignment: a contributor that joined the set + // before this lock may already have forwarded a lower Object. + if forwarded && (!set.forwarded || lowest < set.lowest) { + set.lowest, set.forwarded = lowest, true + } for _, sub := range initialSubs { h.openWriterForSub(ctx, set.hdr, sub, set.writers, pubTimeouts, ref) } diff --git a/pkg/relay/handler_fanout_multipub_test.go b/pkg/relay/handler_fanout_multipub_test.go index a82ce795..1b690198 100644 --- a/pkg/relay/handler_fanout_multipub_test.go +++ b/pkg/relay/handler_fanout_multipub_test.go @@ -611,3 +611,77 @@ func TestFanout_MultiPublisher_FirstObjectOnlyForSubgroupsFirst(t *testing.T) { }) } } + +// TestFanout_MultiPublisher_FirstObjectAfterTeardown: the lowest Object +// forwarded in a Subgroup outlives its contributors, so a later contributor's +// FIRST_OBJECT claim above it is still not honoured (§11.4.2, §2.2). +func TestFanout_MultiPublisher_FirstObjectAfterTeardown(t *testing.T) { + t.Parallel() + pubA, teardown := connectRelay(t, relay.Config{}) + defer teardown() + pubB := dialAnotherClient(t, pubA) + aPub := publishVideoTrack(t, pubA, "cam1", 1) + bPub := publishVideoTrack(t, pubB, "cam1", 2) + subSess := newCam1Subscriber(t, pubA) + + type stream struct { + hdr message.SubgroupHeader + end error + } + streams := make(chan stream, 2) + go func() { + for { + ds, err := subSess.AcceptDataStream(t.Context()) + if err != nil { + return + } + sg, ok := ds.(*session.IncomingSubgroupStream) + if !ok { + return + } + for { + if _, err := sg.ReadObject(); err != nil { + streams <- stream{sg.Header, err} + break + } + } + } + }() + next := func() stream { + t.Helper() + select { + case s := <-streams: + return s + case <-time.After(2 * time.Second): + t.Fatal("no subgroup stream ended") + return stream{} + } + } + + // A sends Objects 0..2 and resets: a FIN would end the Subgroup there, + // making B's Object 3 malformed (§2.4.2). + hdr := message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit} + a, err := aPub.OpenSubgroup(hdr) + if err != nil { + t.Fatalf("A OpenSubgroup: %v", err) + } + for id := range uint64(3) { + if err := a.WriteObjectAt(id, &message.SubgroupObject{Payload: []byte("a")}); err != nil { + t.Fatalf("A WriteObjectAt %d: %v", id, err) + } + } + a.Cancel(moqt.StreamResetCancelled) + next() // the downstream stream ends once the Subgroup has no contributor + + b, err := bPub.OpenSubgroup(hdr) // FIRST_OBJECT set, starting at Object 3 + if err != nil { + t.Fatalf("B OpenSubgroup: %v", err) + } + if err := b.WriteObjectAt(3, &message.SubgroupObject{Payload: []byte("b")}); err != nil { + t.Fatalf("B WriteObjectAt 3: %v", err) + } + _ = b.Close() + if s := next(); !s.hdr.ReplayingSubgroup { + t.Fatal("the stream beginning with Object 3 sets FIRST_OBJECT, though Object 0 was forwarded before it") + } +} diff --git a/pkg/relay/internal/registry/ledger_test.go b/pkg/relay/internal/registry/ledger_test.go index 7317a992..0ce7deb9 100644 --- a/pkg/relay/internal/registry/ledger_test.go +++ b/pkg/relay/internal/registry/ledger_test.go @@ -448,3 +448,24 @@ func TestTrackEntry_LateEndPurgesCache(t *testing.T) { } } } + +// TestTrackEntry_LowestForwarded: the lowest Object ID forwarded per Subgroup, +// status Objects included and datagrams not, outlives the writers of the +// Subgroup so a later FIRST_OBJECT claim can be checked (§11.4.2, §2.2). +func TestTrackEntry_LowestForwarded(t *testing.T) { + t.Parallel() + e := newTestEntry("lowest") + if _, ok := e.LowestForwarded(1, 0); ok { + t.Fatal("a Subgroup nothing was forwarded in has a lowest Object") + } + mustClaim(t, e, registry.ObjectInfo{Group: 1, Object: 5}, true) + mustClaim(t, e, registry.ObjectInfo{Group: 1, Object: 3}, true) + mustClaim(t, e, registry.ObjectInfo{Group: 1, Object: 1, Datagram: true}, true) + mustClaim(t, e, registry.ObjectInfo{Group: 1, Object: 0, Subgroup: 1}, true) + mustClaim(t, e, registry.ObjectInfo{Group: 1, Object: 9, Subgroup: 2, Status: message.ObjectStatusEndOfGroup}, true) + for _, tc := range []struct{ subgroup, want uint64 }{{0, 3}, {1, 0}, {2, 9}} { + if low, ok := e.LowestForwarded(1, tc.subgroup); !ok || low != tc.want { + t.Errorf("LowestForwarded(1, %d) = (%d, %v), want (%d, true)", tc.subgroup, low, ok, tc.want) + } + } +} diff --git a/pkg/relay/internal/registry/track_entry.go b/pkg/relay/internal/registry/track_entry.go index 817e9abb..d3e6a323 100644 --- a/pkg/relay/internal/registry/track_entry.go +++ b/pkg/relay/internal/registry/track_entry.go @@ -387,6 +387,12 @@ type subgroupLedger struct { maxNormal, maxStatus uint64 hasNormal, hasStatus bool end end + // minObject is the lowest Object ID forwarded, if hasObject: a + // FIRST_OBJECT claim above it is wrong (§11.4.2, §2.2). A duplicate + // whose first copy was in another Subgroup or a datagram counts too; + // that can only clear FIRST_OBJECT, never set it. + minObject uint64 + hasObject bool } // SubgroupEnded records that an inbound subgroup stream ended with a FIN after @@ -478,6 +484,21 @@ func (e *TrackEntry) RecordDuplicate(o ObjectInfo) error { return nil } +// LowestForwarded reports the lowest Object ID forwarded in Subgroup +// (group, subgroup), if any, within the window. The writers of a Subgroup +// forget it when its last contributor leaves; this keeps it for a later +// contributor's FIRST_OBJECT claim (§11.4.2, §2.2). +func (e *TrackEntry) LowestForwarded(group, subgroup uint64) (uint64, bool) { + e.deliveredMu.Lock() + defer e.deliveredMu.Unlock() + g := e.delivered[group] + if g == nil { + return 0, false + } + sg := g.subgroups[subgroup] + return sg.minObject, sg.hasObject +} + // groupEnd reports where o ends its Group (see [end]), if it does: an // END_OF_GROUP or END_OF_TRACK status at M at M (§11.2.1.1); a datagram's // END_OF_GROUP bit on Object N at N+1, which §2.4.2's non-exhaustive list does @@ -631,6 +652,9 @@ func (e *TrackEntry) recordEndsLocked(g *deliveredGroup, o ObjectInfo) { } else { sg.maxStatus, sg.hasStatus = max(sg.maxStatus, o.Object), true } + if !sg.hasObject || o.Object < sg.minObject { + sg.minObject, sg.hasObject = o.Object, true + } g.subgroups[o.Subgroup] = sg }