Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion STATUS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`). |
Expand Down
11 changes: 9 additions & 2 deletions pkg/relay/handler_fanout.go
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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)
}
Expand Down
74 changes: 74 additions & 0 deletions pkg/relay/handler_fanout_multipub_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
}
21 changes: 21 additions & 0 deletions pkg/relay/internal/registry/ledger_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
}
24 changes: 24 additions & 0 deletions pkg/relay/internal/registry/track_entry.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
}

Expand Down
Loading