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 }