From 22e4d01800f122eab9afb03d50dd5916221adca1 Mon Sep 17 00:00:00 2001 From: Vsevolod Strukchinsky Date: Sat, 26 Sep 2026 14:58:02 +0500 Subject: [PATCH] =?UTF-8?q?fix(relay):=20detect=20the=20rest=20of=20=C2=A7?= =?UTF-8?q?2.4.2's=20malformed-track=20conditions?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The relay ended a track for Object Properties and for an Object after END_OF_GROUP on the same stream. It now also detects, on live subgroup and datagram Objects of any upstream, against the last 32 Groups: - a Subgroup's Publisher Priority changing (item 1); - an Object past a Subgroup's, Group's or Track's end, or two different ends (items 2 to 5). An end is kept as the first missing ID: an END_OF_GROUP or END_OF_TRACK status at M ends the Group at M (END_OF_TRACK the Track too); a FIN (§11.4.3; the Group's too with the END_OF_GROUP header bit, §11.4.2) or a datagram's END_OF_GROUP bit after Object N at N+1. Each end is also checked against the Objects already received, duplicates included (TrackEntry.RecordDuplicate, after the §9.1 check), and the Objects past it are removed from the cache (§2.4.2 MUST NOT be cached), as is a first copy that fails the §9.1 duplicate check. Interpretations, chosen with the user and marked in STATUS.md: - §2.4.2 calls both the status Object at M and Object N the "final Object"; they are one end when M = N+1, as §9.1 lets a relay turn one into the other; - Objects past an end make the track malformed, not dropped; - a Normal Object at an end is past it only if a FIN or bit set it: with status Objects alone it is §9.1's existing-to-not-existing change or §2.1's late Object; a status end at M and a FIN end at M+1 agree for the same reason; - a datagram's END_OF_GROUP bit counts as a Group's end; - item 7 is read per Object (§11.2.1), which the §9.1 check covers. A claimed-but-not-yet-cached race with a later end is documented. ClaimDelivered takes an ObjectInfo; SubgroupEnded is called on a clean inbound FIN via the new sessionHandler.inboundEnded, extracted from runFanout. TestFanout_MultiPublisher_MergesDisjointObjects FIN'd one Subgroup on two streams with different finals (item 3); A now resets. Tests, each verified red first: TestTrackEntry_FinalObjects (priority, ends, their agreement and every order the reviews probed), _RecordDuplicateChecksAgain, _LateEndPurgesCache, TestCheckDuplicatePurgesFirstCopy, six TestRelay_MalformedObjectEndsTrack cases, TestRelay_EndSignalsAgree, TestRelay_DuplicateEndOfGroupEndsGroup. Co-Authored-By: Claude Opus 5.5 (1M context) --- STATUS.md | 53 ++- pkg/relay/duplicate_consistency_test.go | 38 +++ pkg/relay/handler_datagram.go | 15 +- pkg/relay/handler_duplicate.go | 4 + pkg/relay/handler_duplicate_test.go | 25 ++ pkg/relay/handler_fanout.go | 79 +++-- pkg/relay/handler_fanout_multipub_test.go | 19 +- pkg/relay/internal/registry/track_entry.go | 356 ++++++++++++++++++++- pkg/relay/internal/registry/track_test.go | 347 +++++++++++++++++++- pkg/relay/malformed_track_test.go | 127 ++++++++ 10 files changed, 1010 insertions(+), 53 deletions(-) create mode 100644 pkg/relay/handler_duplicate_test.go diff --git a/STATUS.md b/STATUS.md index fba65b55..ff9cd496 100644 --- a/STATUS.md +++ b/STATUS.md @@ -428,22 +428,43 @@ Known protocol gaps, roughly ordered by how load-bearing they are: 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. -- **Malformed tracks: per-Object conditions (§2.4.2, §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 downstream subscriber, its subscription to that - publisher cancelled, the Object not cached. The relay also ends it for two - Prior Group ID Gap values in one Group, and for a duplicate that differs from - the cached first copy (§9.1; Payloads and Immutable Properties compared only - when both copies are Normal, since Normal may become End of Group). A - duplicate is not compared once the first copy left the cache (evicted or - expired), nor while a concurrent contributor has yet to cache it. Not - detected: the rest of §2.4.2's list other than an Object after END_OF_GROUP - on the same stream. A downstream FETCH already being served from the cache - when the track is found malformed is not reset: the relay does not track - fetch streams per track. One interpretation: an Object with two Immutable - Properties is treated as malformed, although §12.7 states "MUST NOT contain - more than one instance" outside its list of malformed conditions. +- **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 + downstream subscriber, its subscription to that publisher cancelled, the + Objects triggering it not cached (removed, if earlier ones were). The relay + also detects, on live subgroup and datagram Objects of any upstream, against + the last 32 Groups (`registry.TrackEntry.ClaimDelivered`, `SubgroupEnded`, + `RecordDuplicate`): §2.4.2's list — a Subgroup's Publisher Priority + changing; an Object past a Subgroup's, Group's or Track's end, or two + different ends; a duplicate that differs from the cached first copy (§9.1; + Payloads and Immutable Properties compared only when both copies are + Normal, since Normal may become End of Group) — and two Prior Group ID Gap + values in one Group. An end is the first missing ID: an END_OF_GROUP or + END_OF_TRACK status at M ends the Group at M (END_OF_TRACK the Track too), + and a FIN (§11.4.3; the Group's too with END_OF_GROUP set, §11.4.2) or a + datagram's END_OF_GROUP bit after Object N at N+1. Each end is also checked + against the Objects already received, duplicates included. Interpretations, + the first three chosen with the maintainer: §2.4.2 calls both the status + Object at M and the Object N the "final Object", which would make the + §9.1-equivalent M = N+1 two different finals; a Normal Object at an end is + past it only if a FIN or bit set it — with status Objects alone it is the + §9.1 existing-to-not-existing change or the §2.1 late Object; for the same + reason a status end at M and a FIN or bit end at M+1 agree, ending at M; a + datagram's END_OF_GROUP bit counts as a Group's end, which §2.4.2's + non-exhaustive list does not name; and §2.4.2 item 7 (another Forwarding + Preference) is read per Object, since it may vary within a Track (§11.2.1), + so it is the §9.1 duplicate check. Not detected: in upstream FETCH + responses (whose End of Track is not used either), in a session that is not + a relay's, past the window, and a duplicate once the first copy left the + cache (evicted or expired) or before a concurrent contributor cached it. An + Object another upstream had claimed but not yet cached when a later end put + it past that end stays cached. A downstream FETCH already being served from + the cache when the track is found malformed is not reset: the relay does not + track fetch streams per track. One interpretation: an Object with two + Immutable Properties is treated as malformed, although §12.7 states "MUST + NOT contain more than one instance" outside its list of malformed + conditions. - **Objects inside an announced gap are dropped, not malformed (§2.1, §9.1, §12.8, §12.9)** — an interpretation. §12.8 and §12.9 list "an Object with an ID within a previously communicated gap" and "a gap covering an Object it diff --git a/pkg/relay/duplicate_consistency_test.go b/pkg/relay/duplicate_consistency_test.go index ece03b6c..a263be63 100644 --- a/pkg/relay/duplicate_consistency_test.go +++ b/pkg/relay/duplicate_consistency_test.go @@ -175,3 +175,41 @@ func TestRelay_DuplicateConsistency(t *testing.T) { }) } } + +// TestRelay_DuplicateEndOfGroupEndsGroup: an END_OF_GROUP arriving as a +// duplicate of a Normal Object (§9.1: existing to not existing) still ends the +// Group, so an Object past it makes the track malformed (§2.4.2). +func TestRelay_DuplicateEndOfGroupEndsGroup(t *testing.T) { + t.Parallel() + pubA, teardown := connectRelay(t, relay.Config{}) + defer teardown() + pubB := dialAnotherClient(t, pubA) + publishVideoTrack(t, pubA, "cam1", 1) + publishVideoTrack(t, pubB, "cam1", 2) + subSess := dialAnotherClient(t, pubA) + subReq := subscribeCam1(t, subSess) + received := make(chan got, 8) + go receiveDatagrams(t.Context(), subSess, received) + + send := func(sess *session.Session, alias uint64, d *message.ObjectDatagram) { + t.Helper() + d.TrackAlias, d.GroupID = alias, 1 + d.Type |= message.DatagramDefaultPriorityBit + if err := sess.SendDatagram(d); err != nil { + t.Fatalf("SendDatagram: %v", err) + } + } + send(pubA, 1, &message.ObjectDatagram{ObjectID: 2, ObjectPayload: []byte("x")}) + select { + case <-received: // forwarded, so cached + case <-time.After(2 * time.Second): + t.Fatal("Object 2 not forwarded") + } + send(pubB, 2, &message.ObjectDatagram{ + Type: message.DatagramStatusBit, ObjectID: 2, ObjectStatus: message.ObjectStatusEndOfGroup, + }) + send(pubB, 2, &message.ObjectDatagram{ObjectID: 3, ObjectPayload: []byte("x")}) + if pd := awaitPublishDone(t, subReq); pd.StatusCode != moqt.PublishDoneMalformedTrack { + t.Fatalf("PUBLISH_DONE %#x, want MALFORMED_TRACK", uint64(pd.StatusCode)) + } +} diff --git a/pkg/relay/handler_datagram.go b/pkg/relay/handler_datagram.go index 6ff3a6cd..8e1e9cc8 100644 --- a/pkg/relay/handler_datagram.go +++ b/pkg/relay/handler_datagram.go @@ -63,7 +63,16 @@ func (h *sessionHandler) handleDatagram(ctx context.Context, d *message.ObjectDa // §9.3: the first copy of {GroupID, ObjectID} wins, unless an announced // gap says it does not exist (§2.1, §9.1). - fresh, err := entry.ClaimDelivered(d.GroupID, d.ObjectID, message.ObjectPriorGaps(d.Properties)) + info := registry.ObjectInfo{ + Group: d.GroupID, + Object: d.ObjectID, + Datagram: true, + Priority: d.PublisherPriority, + Status: d.ObjectStatus, + EndOfGroup: d.HasEndOfGroup(), + Gaps: message.ObjectPriorGaps(d.Properties), + } + fresh, err := entry.ClaimDelivered(info) if err != nil { h.endMalformedTrack(ctx, entry, h.sess, err) return @@ -79,6 +88,10 @@ func (h *sessionHandler) handleDatagram(ctx context.Context, d *message.ObjectDa Payload: d.ObjectPayload, }); err != nil { h.endMalformedTrack(ctx, entry, h.sess, err) + return + } + if err := entry.RecordDuplicate(info); err != nil { + h.endMalformedTrack(ctx, entry, h.sess, err) } return } diff --git a/pkg/relay/handler_duplicate.go b/pkg/relay/handler_duplicate.go index 3d238224..c30bcb63 100644 --- a/pkg/relay/handler_duplicate.go +++ b/pkg/relay/handler_duplicate.go @@ -23,6 +23,9 @@ import ( // A copy an announced gap says does not exist is compared too, when the // first copy is cached: §2.1 excuses its arrival, not a different content. // +// On a difference the first copy is removed from the cache: it triggered the +// Malformed Track status too, and such Objects MUST NOT be cached (§2.4.2). +// // Limitation: nothing is compared when the first copy is not in the cache: // evicted, expired (§12.3), or not yet put there by a concurrent contributor. func checkDuplicate(c *cache.ObjectCache, dup *cache.CachedObject) error { @@ -47,6 +50,7 @@ func checkDuplicate(c *cache.ObjectCache, dup *cache.CachedObject) error { default: return nil } + c.Delete(dup.GroupID, dup.ObjectID) return fmt.Errorf("%w: a duplicate of Object %d in Group %d has a different %s (§9.1, §2.4.2)", session.ErrMalformedTrack, dup.ObjectID, dup.GroupID, field) } diff --git a/pkg/relay/handler_duplicate_test.go b/pkg/relay/handler_duplicate_test.go new file mode 100644 index 00000000..b6f2da49 --- /dev/null +++ b/pkg/relay/handler_duplicate_test.go @@ -0,0 +1,25 @@ +package relay + +import ( + "testing" + + "github.com/floatdrop/moq-go/pkg/relay/cache" +) + +// TestCheckDuplicatePurgesFirstCopy: a first copy that a differing duplicate +// makes malformed is removed from the cache: Object(s) triggering Malformed +// Track status MUST NOT be cached (§2.4.2). +func TestCheckDuplicatePurgesFirstCopy(t *testing.T) { + t.Parallel() + c := cache.NewObjectCache(8, 0) + c.Put(&cache.CachedObject{GroupID: 1, ObjectID: 0, Payload: []byte("a")}) + if err := checkDuplicate(c, &cache.CachedObject{GroupID: 1, ObjectID: 0, Payload: []byte("a")}); err != nil { + t.Fatalf("an identical copy: %v", err) + } + if err := checkDuplicate(c, &cache.CachedObject{GroupID: 1, ObjectID: 0, Payload: []byte("b")}); err == nil { + t.Fatal("a copy with another Payload is not malformed") + } + if _, ok := c.Get(1, 0); ok { + t.Error("the first copy is still cached") + } +} diff --git a/pkg/relay/handler_fanout.go b/pkg/relay/handler_fanout.go index bb200b02..359dd4ee 100644 --- a/pkg/relay/handler_fanout.go +++ b/pkg/relay/handler_fanout.go @@ -297,6 +297,8 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming var ( firstObj = true + // last is the last Object read, which a FIN ends the Subgroup after. + last *message.SubgroupObject // pos counts every Object read, dedup losers included: each is an // Object between its neighbours (§11.4.3). pos = inboundPos{src: stream} @@ -320,25 +322,10 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming obj, err = stream.ReadObject() } if err != nil { - if errors.Is(err, io.EOF) { - return // clean end of stream — last contributor will FIN. - } - if errors.Is(err, context.Canceled) { - // The session is going away: reset, never FIN. - inboundReset = true - return - } - if errors.Is(err, session.ErrMalformedTrack) { - malformed(err) - return - } - h.log.LogAttrs(ctx, slog.LevelDebug, "fanout: inbound ReadObject failed", - slog.String("err", err.Error())) - // An unparseable object leaves the publisher writing; stop it. - stream.Cancel(moqt.StreamResetInternalError) - inboundReset = true + inboundReset, inboundResetCode = h.inboundEnded(ctx, entry, stream, hdr, last, err) return } + last = obj if terminalSeen { h.log.LogAttrs(ctx, slog.LevelDebug, @@ -360,7 +347,15 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming // §9.3: the first upstream to deliver {GroupID, ObjectID} forwards it, // unless an announced gap says it does not exist (§2.1, §9.1). Outside // sg.Mu, so dedup losers never touch the writer set. - fresh, err := entry.ClaimDelivered(hdr.GroupID, objectID, message.ObjectPriorGaps(obj.Properties)) + info := registry.ObjectInfo{ + Group: hdr.GroupID, + Object: objectID, + Subgroup: hdr.SubgroupID, + Priority: hdr.PublisherPriority, + Status: obj.ObjectStatus, + Gaps: message.ObjectPriorGaps(obj.Properties), + } + fresh, err := entry.ClaimDelivered(info) if err != nil { malformed(err) return @@ -379,6 +374,10 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming malformed(err) return } + if err := entry.RecordDuplicate(info); err != nil { + malformed(err) + return + } continue // redundant copy already forwarded by a peer upstream. } @@ -437,6 +436,50 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming } } +// inboundEnded handles err, which ended the reads of stream (whose header is +// hdr, and whose last Object read was last, nil if none), and reports whether +// the stream counts as reset, and with which code, for the Subgroup's FIN +// (§9.3). +func (h *sessionHandler) inboundEnded( + ctx context.Context, + entry *registry.TrackEntry, + stream *session.IncomingSubgroupStream, + hdr message.SubgroupHeader, + last *message.SubgroupObject, + err error, +) (reset bool, code moqt.StreamResetCode) { + switch { + case errors.Is(err, io.EOF): + // Clean end of stream; the last contributor will FIN. It ends the + // Subgroup after the last Object read (§2.4.2). + if last == nil { + return false, 0 + } + err = entry.SubgroupEnded(registry.ObjectInfo{ + Group: hdr.GroupID, + Object: stream.ObjectID(), + Subgroup: hdr.SubgroupID, + Priority: hdr.PublisherPriority, + Status: last.ObjectStatus, + }, hdr.EndOfGroup) + if err == nil { + return false, 0 + } + case errors.Is(err, context.Canceled): + // The session is going away: reset, never FIN. + return true, moqt.StreamResetCancelled + case !errors.Is(err, session.ErrMalformedTrack): + h.log.LogAttrs(ctx, slog.LevelDebug, "fanout: inbound ReadObject failed", + slog.String("err", err.Error())) + // An unparseable object leaves the publisher writing; stop it. + stream.Cancel(moqt.StreamResetInternalError) + return true, moqt.StreamResetCancelled + } + stream.Cancel(moqt.StreamResetMalformedTrack) + h.endMalformedTrack(ctx, entry, h.sess, err) + return true, moqt.StreamResetMalformedTrack +} + // openWriterForSub starts a subgroupWriter for sub and records it in writers // (nil when sub is not Established, so it is not retried). // diff --git a/pkg/relay/handler_fanout_multipub_test.go b/pkg/relay/handler_fanout_multipub_test.go index 7c997ddc..1f7c1a82 100644 --- a/pkg/relay/handler_fanout_multipub_test.go +++ b/pkg/relay/handler_fanout_multipub_test.go @@ -364,15 +364,8 @@ func TestFanout_MultiPublisher_MergesDisjointObjects(t *testing.T) { t.Fatalf("B WriteObjectAt %d: %v", id, err) } } - if err := aSg.Close(); err != nil { - t.Fatalf("A Close: %v", err) - } - if err := bSg.Close(); err != nil { - t.Fatalf("B Close: %v", err) - } - - // Collect every delivered object until the merged stream(s) finish. Six - // distinct objects must arrive, each exactly once. + // Collect every delivered object. Six distinct objects must arrive, each + // exactly once. seen := map[uint64]int{} deadline := time.After(3 * time.Second) for len(seen) < 6 { @@ -396,6 +389,14 @@ func TestFanout_MultiPublisher_MergesDisjointObjects(t *testing.T) { t.Fatalf("object %d missing from delivered set %v", id, slices.Sorted(maps.Keys(seen))) } } + + // A ends with a reset: a FIN would make its last Object, 4, the + // Subgroup's final one, and B's 5 past it would make the track malformed + // (§2.4.2). + aSg.Cancel(moqt.StreamResetCancelled) + if err := bSg.Close(); err != nil { + t.Fatalf("B Close: %v", err) + } } // awaitObject waits for the next delivered object event (failing on a stream end diff --git a/pkg/relay/internal/registry/track_entry.go b/pkg/relay/internal/registry/track_entry.go index b014d137..817e9abb 100644 --- a/pkg/relay/internal/registry/track_entry.go +++ b/pkg/relay/internal/registry/track_entry.go @@ -3,6 +3,7 @@ package registry import ( "fmt" "maps" + "math" "slices" "sync" "time" @@ -143,8 +144,12 @@ type TrackEntry struct { // It also holds the Prior Group and Object ID Gaps (§12.8, §12.9) seen: // each Group its Object ID gaps and Group gap value, and groupGaps the Group // ID ranges announced absent, pruned with the same window. - delivered map[uint64]*deliveredGroup - groupGaps []idRange + delivered map[uint64]*deliveredGroup + groupGaps []idRange + // trackEnd is where the Track ends, its END_OF_TRACK (§2.4.2), if + // hasTrackEnd; kept past the window. + trackEnd message.Location + hasTrackEnd bool deliveredMax uint64 deliveredHasMax bool @@ -250,9 +255,11 @@ const deliveredGroupWindow = 32 // drop objects received on a multi-object stream". A cached copy of an Object // a gap later covers is kept (§9.1 makes updating the cache a MAY). The one gap rule // it does report, as an error wrapping [session.ErrMalformedTrack], is a Group -// carrying two Prior Group ID Gap values (§12.8). The rejected Object leaves no -// state in the ledger. -func (e *TrackEntry) ClaimDelivered(group, object uint64, gaps message.PriorGaps) (bool, error) { +// carrying two Prior Group ID Gap values (§12.8). An Object ClaimDelivered +// rejects leaves no state in the ledger, and a duplicate records only its +// gaps, since the caller's §9.1 check may still reject it. +func (e *TrackEntry) ClaimDelivered(o ObjectInfo) (bool, error) { + group, object, gaps := o.Group, o.Object, o.Gaps e.deliveredMu.Lock() defer e.deliveredMu.Unlock() @@ -267,6 +274,9 @@ func (e *TrackEntry) ClaimDelivered(group, object uint64, gaps message.PriorGaps return false, fmt.Errorf("%w: Group %d carries Prior Group ID Gaps %d and %d (§12.8)", session.ErrMalformedTrack, group, g.groupGap, gaps.Group) } + if err := e.checkEndsLocked(g, o); err != nil { + return false, err + } if e.announcedAbsentLocked(g, group, object) { return false, nil } @@ -295,14 +305,345 @@ func (e *TrackEntry) ClaimDelivered(group, object uint64, gaps message.PriorGaps if _, ok := g.objects[object]; ok { return false, nil } + e.recordEndsLocked(g, o) g.objects[object] = struct{}{} return true, nil } +// ObjectInfo is what [TrackEntry.ClaimDelivered] checks of an Object against +// the earlier Objects of the track. +type ObjectInfo struct { + Group, Object uint64 + // Subgroup is the Object's Subgroup ID, unless Datagram (§11.2.1). + Subgroup uint64 + Datagram bool + // Priority is the resolved Publisher Priority (§7). + Priority uint8 + // Status is the Object Status (§11.2.1.1). + Status uint64 + // EndOfGroup is a datagram's END_OF_GROUP bit (§11.3.1). + EndOfGroup bool + // Gaps are the Object's Prior Group and Object ID Gaps (§12.8, §12.9). + Gaps message.PriorGaps +} + +// end is where a Subgroup, Group or Track ends: its first missing Object ID +// (§2.4.2). A status Object at M ends it at M (§11.2.1.1); a FIN after Object N +// (§11.4.3), or an END_OF_GROUP bit on it (§11.4.2, §11.3.1), at N+1. The two +// agree when M = N+1, as §9.1 lets a relay turn one into the other. +// +// Interpretation: §2.4.2 calls both the "final Object", the status Object at +// M and the Object N, which would make M = N+1 two different finals. +type end struct { + at uint64 + set bool + // hard reports that a FIN or END_OF_GROUP bit set it. Only then is a + // Normal Object at at past the end: with status Objects alone, it is the + // Object going from existing to not existing (§9.1), or the late Object of + // §2.1, in either order. + hard bool +} + +// past reports whether an Object at id is past e: a status Object strictly +// after it, a Normal one also at it if e is hard. +func (e end) past(id uint64, normal bool) bool { + return e.set && (id > e.at || id == e.at && normal && e.hard) +} + +// conflicts reports whether e already ends somewhere other than at (set by a +// FIN or bit if hard). A status end at M and a hard end at M+1 agree: Object M +// went from existing to not existing (§9.1), and the end is M. +func (e end) conflicts(at uint64, hard bool) bool { + switch { + case !e.set, e.at == at: + return false + case hard && !e.hard: + return at != e.at+1 + case !hard && e.hard: + return at+1 != e.at + } + return true +} + +// with is e also ending at at, set by a FIN or bit if hard, once checked with +// conflicts. +func (e end) with(at uint64, hard bool) end { + switch { + case !e.set: + return end{at: at, set: true, hard: hard} + case at == e.at: + return end{at: at, set: true, hard: e.hard || hard} + case at < e.at: + return end{at: at, set: true} // a status end below a hard one + } + return e +} + +// subgroupLedger is one Subgroup of a [deliveredGroup]. +type subgroupLedger struct { + priority uint8 + // maxNormal and maxStatus are the largest Normal and status Object IDs + // received, if hasNormal and hasStatus. + maxNormal, maxStatus uint64 + hasNormal, hasStatus bool + end end +} + +// SubgroupEnded records that an inbound subgroup stream ended with a FIN after +// lastObj: that ends its Subgroup (see [end]; on a status Object, at it) and, +// if the stream's header set END_OF_GROUP (§11.4.2), the Group. An error +// wrapping [session.ErrMalformedTrack] reports a §2.4.2 condition: another +// stream of the Subgroup ended elsewhere, or an Object past this end was +// received; the Objects past the lower end are then removed from the cache. +func (e *TrackEntry) SubgroupEnded(lastObj ObjectInfo, endOfGroup bool) error { + group, subgroup := lastObj.Group, lastObj.Subgroup + e.deliveredMu.Lock() + defer e.deliveredMu.Unlock() + g := e.delivered[group] + if g == nil { + return nil // aged out of the window, or every Object dropped + } + hard := lastObj.Status == message.ObjectStatusNormal + at := lastObj.Object + if hard { + at++ + } + ends := end{at: at, set: true, hard: hard} + sg, seen := g.subgroups[subgroup] + if !seen { + sg.priority = lastObj.Priority // every Object of it was dropped + } + sgEnd, gEnd := sg.end.with(at, hard), g.end.with(at, hard) + switch { + case sg.end.conflicts(at, hard): + e.purgePastLocked(group, &subgroup, lower(sg.end, ends)) + return fmt.Errorf("%w: Subgroup %d of Group %d ends at Objects %d and %d (§2.4.2)", + session.ErrMalformedTrack, subgroup, group, sg.end.at, at) + case sg.hasStatus && sgEnd.past(sg.maxStatus, false), + sg.hasNormal && sgEnd.past(sg.maxNormal, true): + e.purgePastLocked(group, &subgroup, sgEnd) + return fmt.Errorf("%w: Subgroup %d of Group %d ends at Object %d below an Object received (§2.4.2)", + session.ErrMalformedTrack, subgroup, group, at) + case !endOfGroup: + case g.end.conflicts(at, hard): + e.purgePastLocked(group, nil, lower(g.end, ends)) + return fmt.Errorf("%w: Group %d ends at Objects %d and %d (§2.4.2)", + session.ErrMalformedTrack, group, g.end.at, at) + case g.pastEnd(gEnd): + e.purgePastLocked(group, nil, gEnd) + return fmt.Errorf("%w: Group %d ends at Object %d below an Object received (§2.4.2)", + session.ErrMalformedTrack, group, at) + } + sg.end = sgEnd + if g.subgroups == nil { + g.subgroups = make(map[uint64]subgroupLedger) + } + g.subgroups[subgroup] = sg + if endOfGroup { + g.end = gEnd + } + return nil +} + +// lower is whichever of a and b ends first. +func lower(a, b end) end { + if b.at < a.at { + return b + } + return a +} + +// RecordDuplicate records what o, a copy [TrackEntry.ClaimDelivered] reported +// as already delivered, says about its Subgroup, Group and Track, once the +// caller's §9.1 check found it consistent with the first copy: an +// END_OF_GROUP arriving after a Normal Object at its ID still ends the Group, +// and a Normal copy of a status Object still counts as received. +// +// Objects claimed since o's [TrackEntry.ClaimDelivered] are checked too: an +// error wrapping [session.ErrMalformedTrack] reports a §2.4.2 condition. +func (e *TrackEntry) RecordDuplicate(o ObjectInfo) error { + e.deliveredMu.Lock() + defer e.deliveredMu.Unlock() + g := e.delivered[o.Group] + if g == nil { + return nil + } + if _, ok := g.objects[o.Object]; !ok { // one dropped as announced absent + return nil + } + if err := e.checkEndsLocked(g, o); err != nil { + return err + } + e.recordEndsLocked(g, o) + return nil +} + +// 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 +// not name but §11.3.1 defines alike. +func groupEnd(o ObjectInfo) (at uint64, hard, ok bool) { + switch { + case o.Status == message.ObjectStatusEndOfGroup, o.Status == message.ObjectStatusEndOfTrack: + return o.Object, false, true + case o.Datagram && o.EndOfGroup: + return o.Object + 1, true, true + } + return 0, false, false +} + +// checkEndsLocked reports the §2.4.2 conditions an Object o makes against the +// earlier ones, g being its Group's ledger entry: a Publisher Priority other +// than its Subgroup's, an Object past its Subgroup's, Group's or Track's end, +// or an end placed elsewhere or below an Object received; the Objects past the +// end are then removed from the cache. +func (e *TrackEntry) checkEndsLocked(g *deliveredGroup, o ObjectInfo) error { + normal := o.Status == message.ObjectStatusNormal + loc := message.Location{Group: o.Group, Object: o.Object} + switch { + case e.hasTrackEnd && e.trackEnd.Less(loc): + return fmt.Errorf("%w: Object %d of Group %d is past END_OF_TRACK at %d/%d (§2.4.2)", + session.ErrMalformedTrack, o.Object, o.Group, e.trackEnd.Group, e.trackEnd.Object) + case o.Status != message.ObjectStatusEndOfTrack: + case e.hasTrackEnd && e.trackEnd != loc: + e.purgeTrackPastLocked(loc) // below the earlier one: not past it + return fmt.Errorf("%w: END_OF_TRACK at Object %d of Group %d and at %d/%d (§2.4.2)", + session.ErrMalformedTrack, o.Object, o.Group, e.trackEnd.Group, e.trackEnd.Object) + case e.purgeTrackPastLocked(loc): + return fmt.Errorf("%w: END_OF_TRACK at Object %d of Group %d is below an Object received (§2.4.2)", + session.ErrMalformedTrack, o.Object, o.Group) + } + if g == nil { + return nil + } + if g.end.past(o.Object, normal) { + return fmt.Errorf("%w: Object %d of Group %d is past its end at %d (§2.4.2)", + session.ErrMalformedTrack, o.Object, o.Group, g.end.at) + } + if at, hard, ok := groupEnd(o); ok { + gEnd := g.end.with(at, hard) + switch { + case g.end.conflicts(at, hard): + e.purgePastLocked(o.Group, nil, lower(g.end, end{at: at, set: true, hard: hard})) + return fmt.Errorf("%w: Group %d ends at Objects %d and %d (§2.4.2)", + session.ErrMalformedTrack, o.Group, g.end.at, at) + case g.pastEnd(gEnd): + e.purgePastLocked(o.Group, nil, gEnd) + return fmt.Errorf("%w: Group %d ends at Object %d below an Object received (§2.4.2)", + session.ErrMalformedTrack, o.Group, at) + } + } + if o.Datagram { + return nil + } + sg, seen := g.subgroups[o.Subgroup] + switch { + case !seen: + case sg.priority != o.Priority: + return fmt.Errorf("%w: Subgroup %d of Group %d has Publisher Priorities %d and %d (§2.4.2)", + session.ErrMalformedTrack, o.Subgroup, o.Group, sg.priority, o.Priority) + case sg.end.past(o.Object, normal): + return fmt.Errorf("%w: Object %d of Subgroup %d in Group %d is past its end at %d (§2.4.2)", + session.ErrMalformedTrack, o.Object, o.Subgroup, o.Group, sg.end.at) + } + return nil +} + +// purgeTrackPastLocked reports whether an Object past END_OF_TRACK at loc was +// received, in the window, and removes those from the cache. END_OF_TRACK is a +// status end (see [end]): a Normal Object at loc is the late Object of §2.1. +func (e *TrackEntry) purgeTrackPastLocked(loc message.Location) bool { + found := false + ended := end{at: loc.Object, set: true} + for id, g := range e.delivered { + switch { + case id > loc.Group: + found = true + e.purgeGroupLocked(id) + case id == loc.Group && g.pastEnd(ended): + found = true + e.purgePastLocked(id, nil, ended) + } + } + return found +} + +// purgeGroupLocked removes every Object of group from the cache, status +// Objects included (§2.4.2). +func (e *TrackEntry) purgeGroupLocked(group uint64) { + objs := e.Cache.GetRange( + message.Location{Group: group}, + message.Location{Group: group, Object: math.MaxUint64}, + message.GroupOrderAscending, + ) + for _, o := range objs { + e.Cache.Delete(group, o.ObjectID) + } +} + +// purgePastLocked removes from the cache the Objects of group, in subgroup if +// not nil, past ended: Object(s) triggering Malformed Track status MUST NOT be +// cached (§2.4.2). +func (e *TrackEntry) purgePastLocked(group uint64, subgroup *uint64, ended end) { + objs := e.Cache.GetRange( + message.Location{Group: group, Object: ended.at}, + message.Location{Group: group, Object: math.MaxUint64}, + message.GroupOrderAscending, + ) + for _, o := range objs { + if !ended.past(o.ObjectID, o.Status == message.ObjectStatusNormal) { + continue + } + if subgroup != nil && (o.ForwardingPref != cache.ForwardingSubgroup || o.SubgroupID != *subgroup) { + continue + } + e.Cache.Delete(group, o.ObjectID) + } +} + +// recordEndsLocked records what o, once checked, says about its Subgroup, +// Group and Track. +func (e *TrackEntry) recordEndsLocked(g *deliveredGroup, o ObjectInfo) { + normal := o.Status == message.ObjectStatusNormal + if normal { + g.maxNormal, g.hasNormal = max(g.maxNormal, o.Object), true + } else { + g.maxStatus, g.hasStatus = max(g.maxStatus, o.Object), true + } + if at, hard, ok := groupEnd(o); ok { + g.end = g.end.with(at, hard) + } + if o.Status == message.ObjectStatusEndOfTrack { + e.trackEnd, e.hasTrackEnd = message.Location{Group: o.Group, Object: o.Object}, true + } + if o.Datagram { + return + } + if g.subgroups == nil { + g.subgroups = make(map[uint64]subgroupLedger) + } + sg, seen := g.subgroups[o.Subgroup] + if !seen { + sg.priority = o.Priority + } + if normal { + sg.maxNormal, sg.hasNormal = max(sg.maxNormal, o.Object), true + } else { + sg.maxStatus, sg.hasStatus = max(sg.maxStatus, o.Object), true + } + g.subgroups[o.Subgroup] = sg +} + // deliveredGroup is one Group of the [TrackEntry.ClaimDelivered] ledger; g is // nil for a Group not yet seen. type deliveredGroup struct { objects map[uint64]struct{} + // maxNormal and maxStatus are the largest Normal and status Object IDs + // received, if hasNormal and hasStatus. + maxNormal, maxStatus uint64 + hasNormal, hasStatus bool + subgroups map[uint64]subgroupLedger // created on first use + end end // objectGaps are the Object ID ranges announced absent (§12.9). objectGaps []idRange // groupGap is the Group's Prior Group ID Gap (§12.8), if hasGroupGap. @@ -310,6 +651,11 @@ type deliveredGroup struct { hasGroupGap bool } +// pastEnd reports whether an Object received is past ended. +func (g *deliveredGroup) pastEnd(ended end) bool { + return g.hasNormal && ended.past(g.maxNormal, true) || g.hasStatus && ended.past(g.maxStatus, false) +} + // idRange is an inclusive range of Group or Object IDs. type idRange struct{ lo, hi uint64 } diff --git a/pkg/relay/internal/registry/track_test.go b/pkg/relay/internal/registry/track_test.go index 783d5937..4a49c956 100644 --- a/pkg/relay/internal/registry/track_test.go +++ b/pkg/relay/internal/registry/track_test.go @@ -34,7 +34,7 @@ func TestTrackEntry_ClaimDelivered(t *testing.T) { e := r.GetOrCreate(newTestTrackName("dedup")) claim := func(group, object uint64) bool { t.Helper() - fresh, err := e.ClaimDelivered(group, object, message.PriorGaps{}) + fresh, err := e.ClaimDelivered(registry.ObjectInfo{Group: group, Object: object}) if err != nil { t.Fatalf("ClaimDelivered(%d, %d): %v", group, object, err) } @@ -119,7 +119,7 @@ func TestTrackEntry_ClaimDeliveredGapProperties(t *testing.T) { e := registry.NewTrackRegistry().GetOrCreate(newTestTrackName("gaps")) last := len(tc.claims) - 1 for i, c := range tc.claims { - fresh, err := e.ClaimDelivered(c.group, c.object, c.gaps) + fresh, err := e.ClaimDelivered(registry.ObjectInfo{Group: c.group, Object: c.object, Gaps: c.gaps}) got := forwarded switch { case err != nil: @@ -149,19 +149,358 @@ func TestTrackEntry_ClaimDeliveredMalformedLeavesNoTrace(t *testing.T) { e := registry.NewTrackRegistry().GetOrCreate(newTestTrackName("gaps")) mustClaim := func(group, object uint64, gaps message.PriorGaps, wantFresh bool) { t.Helper() - fresh, err := e.ClaimDelivered(group, object, gaps) + fresh, err := e.ClaimDelivered(registry.ObjectInfo{Group: group, Object: object, Gaps: gaps}) if err != nil || fresh != wantFresh { t.Fatalf("ClaimDelivered(%d, %d, %+v) = (%v, %v), want (%v, nil)", group, object, gaps, fresh, err, wantFresh) } } mustClaim(9, 0, message.PriorGaps{Group: 2, HasGroup: true}, true) // Groups 7-8 absent - if _, err := e.ClaimDelivered(9, 1, message.PriorGaps{Group: 3, HasGroup: true}); err == nil { + if _, err := e.ClaimDelivered( + registry.ObjectInfo{Group: 9, Object: 1, Gaps: message.PriorGaps{Group: 3, HasGroup: true}}, + ); err == nil { t.Fatal("a second Prior Group ID Gap value in Group 9 is not malformed") } mustClaim(9, 1, message.PriorGaps{}, true) // Object 1 was not recorded mustClaim(6, 0, message.PriorGaps{}, true) // nor the gap of 3 (Groups 6-8) mustClaim(8, 0, message.PriorGaps{}, false) + + // Nor a Subgroup's Publisher Priority. + if _, err := e.ClaimDelivered(registry.ObjectInfo{Group: 9, Object: 5, Priority: 7}); err == nil { + t.Fatal("another Publisher Priority in Subgroup 0 of Group 9 is not malformed") + } + if fresh, err := e.ClaimDelivered(registry.ObjectInfo{Group: 9, Object: 5}); err != nil || !fresh { + t.Fatalf("ClaimDelivered after the malformed one = (%v, %v), want (true, nil)", fresh, err) + } + + // A duplicate records no Subgroup state: the §9.1 check may still reject + // it (the caller's), so Subgroup 3 keeps no priority from it. + if fresh, err := e.ClaimDelivered( + registry.ObjectInfo{Group: 9, Object: 5, Subgroup: 3, Priority: 9}, + ); err != nil || + fresh { + t.Fatalf("a duplicate = (%v, %v), want (false, nil)", fresh, err) + } + if _, err := e.ClaimDelivered(registry.ObjectInfo{Group: 9, Object: 6, Subgroup: 3, Priority: 1}); err != nil { + t.Fatalf("Subgroup 3 kept the duplicate's priority: %v", err) + } +} + +// TestTrackEntry_FinalObjects pins the §2.4.2 conditions that need earlier +// Objects: a Subgroup's Publisher Priority changing (1), an Object past a +// Subgroup's end (2) or two different ends (3), and an Object past the Group's +// (4) or the Track's (5) end, detected in either order. An end is the first +// missing ID: a status at M ends at M (§11.2.1.1), a FIN or END_OF_GROUP bit +// after N at N+1 (§11.4.2, §11.3.1), so the two agree when M = N+1 (§9.1). +func TestTrackEntry_FinalObjects(t *testing.T) { + t.Parallel() + type step func(*registry.TrackEntry) error + // claim does what the relay does: a duplicate, once the §9.1 check + // passed, is recorded too. + claim := func(o registry.ObjectInfo) step { + return func(e *registry.TrackEntry) error { + fresh, err := e.ClaimDelivered(o) + if err == nil && !fresh { + err = e.RecordDuplicate(o) + } + return err + } + } + sub := func(group, object, subgroup uint64) step { + return claim(registry.ObjectInfo{Group: group, Object: object, Subgroup: subgroup}) + } + prio := func(object, subgroup uint64, p uint8) step { + return claim(registry.ObjectInfo{Group: 1, Object: object, Subgroup: subgroup, Priority: p}) + } + status := func(group, object, s uint64) step { + return claim(registry.ObjectInfo{Group: group, Object: object, Status: s}) + } + statusIn := func(object, subgroup, s uint64) step { + return claim(registry.ObjectInfo{Group: 1, Object: object, Subgroup: subgroup, Status: s}) + } + fin := func(subgroup, last uint64, endOfGroup bool) step { + return func(e *registry.TrackEntry) error { + return e.SubgroupEnded(registry.ObjectInfo{Group: 1, Object: last, Subgroup: subgroup}, endOfGroup) + } + } + finOnStatus := func(subgroup, last uint64) step { + return func(e *registry.TrackEntry) error { + return e.SubgroupEnded(registry.ObjectInfo{ + Group: 1, Object: last, Subgroup: subgroup, Status: message.ObjectStatusEndOfGroup, + }, false) + } + } + eog, eot := message.ObjectStatusEndOfGroup, message.ObjectStatusEndOfTrack + for _, tc := range []struct { + name string + steps []step // all but the last are well-formed + malformed bool // the last makes the track malformed + }{ + // 1: Publisher Priority within a Subgroup + {"same priority in a Subgroup", []step{prio(0, 0, 1), prio(1, 0, 1)}, false}, + {"priority changes in a Subgroup", []step{prio(0, 0, 1), prio(1, 0, 2)}, true}, + {"other priority in another Subgroup", []step{prio(0, 0, 1), prio(1, 1, 2)}, false}, + {"datagrams have no Subgroup", []step{ + claim(registry.ObjectInfo{Group: 1, Object: 0, Datagram: true, Priority: 1}), + claim(registry.ObjectInfo{Group: 1, Object: 1, Datagram: true, Priority: 2}), + }, false}, + // 2, 3: the final Object of a Subgroup is the last before a FIN + {"Object past a Subgroup's FIN", []step{sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, false), sub(1, 2, 0)}, true}, + {"FIN below a received Object", []step{sub(1, 0, 0), sub(1, 2, 0), fin(0, 1, false)}, true}, + {"another Subgroup after a FIN", []step{sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, false), sub(1, 5, 1)}, false}, + {"two FINs, same final", []step{sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, false), fin(0, 1, false)}, false}, + {"two FINs, different finals", []step{sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, false), fin(0, 0, false)}, true}, + // 4: the final Object of a Group + {"Object past END_OF_GROUP", []step{status(1, 2, eog), sub(1, 3, 1)}, true}, + {"Object below END_OF_GROUP", []step{status(1, 5, eog), sub(1, 3, 1)}, false}, + {"END_OF_GROUP below a received Object", []step{sub(1, 3, 1), status(1, 2, eog)}, true}, + {"next Group after END_OF_GROUP", []step{status(1, 2, eog), sub(2, 5, 0)}, false}, + {"Object past a datagram's END_OF_GROUP", []step{ + claim(registry.ObjectInfo{Group: 1, Object: 2, Datagram: true, EndOfGroup: true}), + claim(registry.ObjectInfo{Group: 1, Object: 3, Datagram: true}), + }, true}, + {"Object past an END_OF_GROUP Subgroup's FIN", []step{sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), sub(1, 2, 1)}, true}, + {"Object past a plain Subgroup's FIN", []step{sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, false), sub(1, 2, 1)}, false}, + // 5: the final Object of the Track + {"Group past END_OF_TRACK", []step{status(1, 2, eot), sub(2, 0, 0)}, true}, + {"Object past END_OF_TRACK in its Group", []step{status(1, 2, eot), sub(1, 3, 1)}, true}, + {"Object before END_OF_TRACK", []step{status(1, 2, eot), sub(1, 1, 1)}, false}, + {"END_OF_TRACK below a received Object", []step{sub(2, 0, 0), status(1, 2, eot)}, true}, + {"a copy of END_OF_TRACK", []step{status(1, 2, eot), status(1, 2, eot)}, false}, + + // A status at M and a FIN or bit after M-1 are the same end. + {"END_OF_GROUP status after an END_OF_GROUP FIN", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), status(1, 2, eog), + }, false}, + {"END_OF_GROUP FIN after an END_OF_GROUP status", []step{ + status(1, 2, eog), sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), + }, false}, + {"END_OF_GROUP status one past a datagram's END_OF_GROUP", []step{ + claim(registry.ObjectInfo{Group: 1, Object: 1, Datagram: true, EndOfGroup: true}), + status(1, 2, eog), + }, false}, + {"END_OF_GROUP status two past a datagram's END_OF_GROUP", []step{ + claim(registry.ObjectInfo{Group: 1, Object: 1, Datagram: true, EndOfGroup: true}), + status(1, 3, eog), + }, true}, + {"END_OF_GROUP FIN before an END_OF_TRACK", []step{ + status(1, 2, eot), sub(1, 0, 1), sub(1, 1, 1), fin(1, 1, true), + }, false}, + {"a Subgroup ending on a status, and after the Object before it", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, false), finOnStatus(0, 2), + }, false}, + // A status end at M and a FIN end at M+1 are Object M going from + // existing to not existing (§9.1). + {"a Subgroup ending on a status, and after the Object at it", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, false), finOnStatus(0, 1), + }, false}, + {"a Subgroup ending after an Object, and on a status at it", []step{ + sub(1, 0, 0), sub(1, 1, 0), statusIn(1, 0, eog), finOnStatus(0, 1), fin(0, 1, false), + }, false}, + {"an END_OF_GROUP FIN after an Object, then END_OF_GROUP at it", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), statusIn(1, 0, eog), + }, false}, + {"an Object past both", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), statusIn(1, 0, eog), sub(1, 2, 1), + }, true}, + {"a Subgroup ending two apart", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, false), finOnStatus(0, 0), + }, true}, + // A Normal Object at a status's own ID is the late Object of §2.1. + {"Object at END_OF_GROUP's ID", []step{status(1, 2, eog), sub(1, 2, 1)}, false}, + {"Object at END_OF_TRACK's ID", []step{status(1, 2, eot), sub(1, 2, 1)}, false}, + {"Object at the ID after an END_OF_GROUP FIN", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), sub(1, 2, 1), + }, true}, + {"Object, then END_OF_GROUP at its ID", []step{sub(1, 2, 1), statusIn(2, 1, eog)}, false}, + + // A status arriving as a duplicate of a Normal Object still ends it. + {"Object past an END_OF_GROUP that was a duplicate", []step{ + sub(1, 2, 1), statusIn(2, 1, eog), sub(1, 3, 2), + }, true}, + {"Object past an END_OF_TRACK that was a duplicate", []step{ + sub(1, 2, 1), statusIn(2, 1, eot), sub(2, 0, 0), + }, true}, + // A FIN or bit end wins over a status at the same ID, in every order. + {"Object and END_OF_GROUP at 2, then a FIN after 1", []step{ + sub(1, 2, 1), statusIn(2, 1, eog), sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), + }, true}, + {"a FIN after 1 and END_OF_GROUP at 2, then an Object at 2", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), statusIn(2, 1, eog), sub(1, 2, 1), + }, true}, + {"END_OF_GROUP at 2 and an Object there, then a FIN after 1", []step{ + statusIn(2, 1, eog), sub(1, 2, 1), sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), + }, true}, + // Status Objects past an end, and END_OF_TRACK ending its Group. + {"END_OF_TRACK past an END_OF_GROUP", []step{status(1, 3, eog), status(1, 5, eot)}, true}, + {"END_OF_GROUP past an END_OF_TRACK", []step{status(1, 5, eot), status(1, 7, eog)}, true}, + {"END_OF_GROUP in a Group past END_OF_TRACK", []step{status(1, 5, eot), status(2, 0, eog)}, true}, + {"END_OF_TRACK past an END_OF_GROUP FIN", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), status(1, 5, eot), + }, true}, + {"END_OF_GROUP past a Subgroup's FIN", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, false), statusIn(3, 0, eog), + }, true}, + {"END_OF_TRACK below an END_OF_GROUP in a later Group", []step{status(5, 0, eog), status(3, 2, eot)}, true}, + {"a Subgroup's FIN below an END_OF_GROUP in it", []step{ + sub(1, 1, 0), statusIn(5, 0, eog), fin(0, 1, false), + }, true}, + // A status Object at a hard end, then a status end one below it. + {"END_OF_GROUP at an END_OF_GROUP FIN's end, then one below", []step{ + sub(1, 2, 0), fin(0, 2, true), statusIn(3, 1, eog), statusIn(2, 2, eog), + }, true}, + {"END_OF_GROUP at a datagram's END_OF_GROUP end, then one below", []step{ + claim(registry.ObjectInfo{Group: 1, Object: 2, Datagram: true, EndOfGroup: true}), + status(1, 3, eog), statusIn(2, 2, eog), + }, true}, + // A Subgroup whose every Object was dropped has no priority yet. + {"FIN of a Subgroup with nothing recorded, then its Object", []step{ + prio(0, 0, 5), + func(e *registry.TrackEntry) error { + return e.SubgroupEnded(registry.ObjectInfo{Group: 1, Object: 4, Subgroup: 7, Priority: 5}, false) + }, + prio(1, 7, 5), + }, false}, + {"END_OF_TRACK at an END_OF_GROUP FIN's end", []step{ + sub(1, 0, 0), sub(1, 1, 0), fin(0, 1, true), status(1, 2, eot), + }, false}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + e := registry.NewTrackRegistry().GetOrCreate(newTestTrackName("finals")) + last := len(tc.steps) - 1 + for i, s := range tc.steps { + err := s(e) + if i < last && err != nil { + t.Fatalf("step %d: %v", i, err) + } + if i == last && (err != nil) != tc.malformed { + t.Fatalf("last step: err = %v, want malformed = %v", err, tc.malformed) + } + if err != nil && !errors.Is(err, session.ErrMalformedTrack) { + t.Fatalf("err = %v, want it to wrap session.ErrMalformedTrack", err) + } + } + }) + } +} + +// TestTrackEntry_RecordDuplicateChecksAgain: a duplicate's end is checked +// again when recorded, against what was claimed after its ClaimDelivered. +func TestTrackEntry_RecordDuplicateChecksAgain(t *testing.T) { + t.Parallel() + e := registry.NewTrackRegistry().GetOrCreate(newTestTrackName("race")) + claim := func(o registry.ObjectInfo, wantFresh bool) { + t.Helper() + if fresh, err := e.ClaimDelivered(o); err != nil || fresh != wantFresh { + t.Fatalf("ClaimDelivered(%+v) = (%v, %v), want (%v, nil)", o, fresh, err, wantFresh) + } + } + claim(registry.ObjectInfo{Group: 1, Object: 2}, true) + eog := registry.ObjectInfo{Group: 1, Object: 2, Status: message.ObjectStatusEndOfGroup} + claim(eog, false) // a duplicate; the caller's §9.1 check runs now + claim(registry.ObjectInfo{Group: 1, Object: 5}, true) // meanwhile, another upstream + if err := e.RecordDuplicate(eog); !errors.Is(err, session.ErrMalformedTrack) { + t.Fatalf("RecordDuplicate = %v, want the Group ending below Object 5 to be malformed", err) + } +} + +// TestTrackEntry_LateEndPurgesCache: an end placed below Normal Objects +// already received makes them past it, and Object(s) triggering Malformed +// Track status MUST NOT be cached (§2.4.2). +func TestTrackEntry_LateEndPurgesCache(t *testing.T) { + t.Parallel() + e := registry.NewTrackRegistry().GetOrCreate(newTestTrackName("purge")) + for _, id := range []uint64{0, 1, 2} { + e.Cache.Put(&cache.CachedObject{GroupID: 1, ObjectID: id, Payload: []byte("x")}) + if _, err := e.ClaimDelivered(registry.ObjectInfo{Group: 1, Object: id}); err != nil { + t.Fatalf("ClaimDelivered %d: %v", id, err) + } + } + if err := e.SubgroupEnded(registry.ObjectInfo{Group: 1, Object: 1}, false); err == nil { + t.Fatal("a FIN after 1 with Object 2 received is not malformed") + } + if _, ok := e.Cache.Get(1, 2); ok { + t.Error("Object 2, past the Subgroup's end, is still cached") + } + if _, ok := e.Cache.Get(1, 1); !ok { + t.Error("Object 1 was removed from the cache") + } + + // Two conflicting ends: the Objects past the lower one go too. + e = registry.NewTrackRegistry().GetOrCreate(newTestTrackName("purge-conflict")) + for _, id := range []uint64{0, 1} { + e.Cache.Put(&cache.CachedObject{GroupID: 1, ObjectID: id, Payload: []byte("x")}) + if _, err := e.ClaimDelivered(registry.ObjectInfo{Group: 1, Object: id}); err != nil { + t.Fatalf("ClaimDelivered %d: %v", id, err) + } + } + if err := e.SubgroupEnded(registry.ObjectInfo{Group: 1, Object: 1}, false); err != nil { + t.Fatalf("a FIN after 1: %v", err) + } + if err := e.SubgroupEnded(registry.ObjectInfo{Group: 1, Object: 0}, false); err == nil { + t.Fatal("a second FIN, after 0, is not malformed") + } + if _, ok := e.Cache.Get(1, 1); ok { + t.Error("Object 1, past the lower end, is still cached") + } + + // END_OF_TRACK: Objects past it, status Objects included, go too. + eot, eog := message.ObjectStatusEndOfTrack, message.ObjectStatusEndOfGroup + for _, tc := range []struct { + name string + cached []registry.ObjectInfo + late registry.ObjectInfo + gone []message.Location + kept []message.Location + }{ + { + "a second END_OF_TRACK", + []registry.ObjectInfo{{Group: 1, Object: 0}, {Group: 1, Object: 3, Status: eot}}, + registry.ObjectInfo{Group: 1, Object: 1, Status: eot}, + []message.Location{{Group: 1, Object: 3}}, + []message.Location{{Group: 1, Object: 0}}, + }, + { + "END_OF_GROUP in a later Group", + []registry.ObjectInfo{{Group: 2, Object: 0, Status: eog}}, + registry.ObjectInfo{Group: 1, Object: 0, Status: eot}, + []message.Location{{Group: 2, Object: 0}}, + nil, + }, + { + "END_OF_GROUP in its Group, and a later Group", + []registry.ObjectInfo{ + {Group: 1, Object: 0}, {Group: 1, Object: 5, Subgroup: 1, Status: eog}, {Group: 2, Object: 1}, + }, + registry.ObjectInfo{Group: 1, Object: 2, Subgroup: 2, Status: eot}, + []message.Location{{Group: 1, Object: 5}, {Group: 2, Object: 1}}, + []message.Location{{Group: 1, Object: 0}}, + }, + } { + e := registry.NewTrackRegistry().GetOrCreate(newTestTrackName("purge-eot")) + for _, o := range tc.cached { + e.Cache.Put( + &cache.CachedObject{GroupID: o.Group, ObjectID: o.Object, SubgroupID: o.Subgroup, Status: o.Status}, + ) + if _, err := e.ClaimDelivered(o); err != nil { + t.Fatalf("%s: ClaimDelivered %+v: %v", tc.name, o, err) + } + } + if _, err := e.ClaimDelivered(tc.late); err == nil { + t.Fatalf("%s: END_OF_TRACK below an Object received is not malformed", tc.name) + } + for _, l := range tc.gone { + if _, ok := e.Cache.Get(l.Group, l.Object); ok { + t.Errorf("%s: Object %d of Group %d, past END_OF_TRACK, is still cached", tc.name, l.Object, l.Group) + } + } + for _, l := range tc.kept { + if _, ok := e.Cache.Get(l.Group, l.Object); !ok { + t.Errorf("%s: Object %d of Group %d, before END_OF_TRACK, was removed", tc.name, l.Object, l.Group) + } + } + } } // TestTrackRegistry_GetMissingReturnsFalse confirms the unknown-key path of diff --git a/pkg/relay/malformed_track_test.go b/pkg/relay/malformed_track_test.go index 89a36707..609a29f3 100644 --- a/pkg/relay/malformed_track_test.go +++ b/pkg/relay/malformed_track_test.go @@ -79,6 +79,46 @@ func sendGapObjects(t *testing.T, pubSess *session.Session, alias uint64, objs [ } } +// testStream is a subgroup stream for [sendStreams]: its Objects, each with a +// Status, and whether it ends with a FIN or stays open. +type testStream struct { + group, subgroup uint64 + priority uint8 // inline when set + endOfGroup bool // the header's END_OF_GROUP bit + objects []uint64 + status uint64 // of the last Object + open bool +} + +// sendStreams sends each of streams from pubSess on alias. The relay may read +// them in either order; each case using it is malformed in both. +func sendStreams(t *testing.T, pubSess *session.Session, alias uint64, streams []testStream) { + t.Helper() + for _, s := range streams { + sg, err := pubSess.OpenSubgroup(message.SubgroupHeader{ + SubgroupIDMode: message.SubgroupIDExplicit, TrackAlias: alias, GroupID: s.group, + SubgroupID: s.subgroup, InlinePriority: s.priority != 0, PublisherPriority: s.priority, + EndOfGroup: s.endOfGroup, + }) + if err != nil { + t.Errorf("OpenSubgroup: %v", err) + return + } + for i, id := range s.objects { + o := &message.SubgroupObject{Payload: []byte("x")} + if i == len(s.objects)-1 && s.status != 0 { + o = &message.SubgroupObject{ObjectStatus: s.status} + } + _ = sg.WriteObjectAt(id, o) + } + if s.open { + t.Cleanup(func() { _ = sg.Close() }) + continue + } + _ = sg.Close() + } +} + // priorGap is Object Properties carrying a Prior Group or Object ID Gap. func priorGap(typ message.PropertyType, n uint64) []byte { return message.AppendTrackProperties([]wire.KVPair{{Type: typ, IntVal: n}}) @@ -124,6 +164,50 @@ func TestRelay_MalformedObjectEndsTrack(t *testing.T) { {5, 1, 1, priorGap(message.PropertyPriorGroupIDGap, 1)}, }) }}, + // §2.4.2 item 1: a Subgroup's Publisher Priority changes. + {"priority changes in a Subgroup", func(t *testing.T, pubSess *session.Session, alias uint64) { + sendStreams(t, pubSess, alias, []testStream{ + {group: 1, priority: 1, objects: []uint64{0}, open: true}, + {group: 1, priority: 2, objects: []uint64{1}, open: true}, + }) + }}, + // §2.4.2 item 2: an Object past the last one before a Subgroup's FIN. + {"Object past a Subgroup's FIN", func(t *testing.T, pubSess *session.Session, alias uint64) { + sendStreams(t, pubSess, alias, []testStream{ + {group: 1, objects: []uint64{0, 1}}, + {group: 1, objects: []uint64{2}, open: true}, + }) + }}, + // §2.4.2 item 4: an Object past an END_OF_GROUP. + {"Object past END_OF_GROUP", func(t *testing.T, pubSess *session.Session, alias uint64) { + sendStreams(t, pubSess, alias, []testStream{ + {group: 1, objects: []uint64{2}, status: message.ObjectStatusEndOfGroup}, + {group: 1, subgroup: 1, objects: []uint64{3}, open: true}, + }) + }}, + {"Object past an END_OF_GROUP Subgroup's FIN", func(t *testing.T, pubSess *session.Session, alias uint64) { + sendStreams(t, pubSess, alias, []testStream{ + {group: 1, endOfGroup: true, objects: []uint64{0, 1}}, + {group: 1, subgroup: 1, objects: []uint64{2}, open: true}, + }) + }}, + {"datagram past a datagram's END_OF_GROUP", func(t *testing.T, pubSess *session.Session, alias uint64) { + for _, d := range []*message.ObjectDatagram{ + {Type: message.DatagramEndOfGroupBit, TrackAlias: alias, GroupID: 1, ObjectID: 2, ObjectPayload: []byte("x")}, + {TrackAlias: alias, GroupID: 1, ObjectID: 3, ObjectPayload: []byte("x")}, + } { + if err := pubSess.SendDatagram(d); err != nil { + t.Errorf("SendDatagram: %v", err) + } + } + }}, + // §2.4.2 item 5: an Object past the END_OF_TRACK. + {"Object past END_OF_TRACK", func(t *testing.T, pubSess *session.Session, alias uint64) { + sendStreams(t, pubSess, alias, []testStream{ + {group: 1, objects: []uint64{2}, status: message.ObjectStatusEndOfTrack}, + {group: 2, objects: []uint64{0}, open: true}, + }) + }}, {"datagrams with different group gaps in a Group", func(t *testing.T, pubSess *session.Session, alias uint64) { for i, gap := range []uint64{2, 1} { if err := pubSess.SendDatagram(&message.ObjectDatagram{ @@ -459,3 +543,46 @@ func refusedFetchResetsStream(t *testing.T, upstreamProps, objProps []byte) { t.Fatal("the FETCH stream completed; want it reset over what the upstream sent") } } + +// TestRelay_EndSignalsAgree: an END_OF_GROUP status at 2 and a FIN after Object +// 1 on a stream with END_OF_GROUP set end the Group at the same place +// (§11.2.1.1, §11.4.2, §9.1), so the track goes on. +func TestRelay_EndSignalsAgree(t *testing.T) { + t.Parallel() + pubSess, teardown := connectRelay(t, relay.Config{}) + defer teardown() + publishVideoTrack(t, pubSess, "cam1", 7) + subSess := dialAnotherClient(t, pubSess) + subscribeCam1(t, subSess) + events := make(chan objEvent, 16) + go readSubgroups(t.Context(), subSess, events) + + sendStreams(t, pubSess, 7, []testStream{ + {group: 1, endOfGroup: true, objects: []uint64{0, 1}}, + {group: 1, subgroup: 1, objects: []uint64{2}, status: message.ObjectStatusEndOfGroup}, + }) + // Both downstream streams FIN only after the relay recorded both ends. + for fins := 0; fins < 2; { + select { + case ev := <-events: + if ev.err == nil { + continue + } + if !errors.Is(ev.err, io.EOF) { + t.Fatalf("a downstream stream ended with %v, want a FIN", ev.err) + } + fins++ + case <-time.After(2 * time.Second): + t.Fatalf("%d/2 downstream streams FIN'd", fins) + } + } + sendStreams(t, pubSess, 7, []testStream{{group: 2, objects: []uint64{0}, open: true}}) + select { + case ev := <-events: + if ev.err != nil || ev.absID != 0 { + t.Fatalf("got %+v, want Object 0 of Group 2: the track ended", ev) + } + case <-time.After(2 * time.Second): + t.Fatal("Group 2 not forwarded: the track ended") + } +}