diff --git a/STATUS.md b/STATUS.md index 71d81bf1..371a7e68 100644 --- a/STATUS.md +++ b/STATUS.md @@ -236,8 +236,8 @@ By package, bottom-up along the dependency stack: | 12.5 | DEFAULT_PUBLISHER_GROUP_ORDER | 0x22 | DONE | Validated. | | 12.6 | DYNAMIC_GROUPS | 0x30 | DONE | Property defined & scope-validated (flow: see §5.1.6.1). | | 12.7 | Immutable properties | 0x0B | DONE | Relays cache & forward verbatim, never add. Property lookups search its contents too (`message.ExpandImmutable`), the mutable value winning: delivery timeouts, MAX_CACHE_DURATION, DYNAMIC_GROUPS, Mandatory Track Property screening, and property Range Filters. | -| 12.8 | Prior group ID gap | 0x3C | PARTIAL| Object-scope; encoder in `msf/groupid.go`. More than one, or one past the Group ID, makes the track malformed (`message.CheckObjectProperties`); the rules that need earlier Objects are not checked — see Limitations. | -| 12.9 | Prior object ID gap | 0x3E | PARTIAL| Object-scope. More than one, or one past the Object ID, makes the track malformed; the rules that need earlier Objects are not checked — see Limitations. | +| 12.8 | Prior group ID gap | 0x3C | PARTIAL| Object-scope; encoder in `msf/groupid.go`. More than one, or one past the Group ID, makes the track malformed (`message.CheckObjectProperties`), and the relay also ends the track for two values in one Group. Against the last 32 Groups of any upstream (`registry.TrackEntry.ClaimDelivered`), the relay neither forwards nor caches an Object in a Group announced absent; a gap covering a received Group is accepted (§2.1, §9.1, see Limitations). Not in upstream FETCH responses, nor in a session that is not a relay's. | +| 12.9 | Prior object ID gap | 0x3E | PARTIAL| Object-scope. More than one, or one past the Object ID, makes the track malformed. The relay neither forwards nor caches an Object announced absent, and accepts a gap covering a received Object, as for §12.8. | ## §13 Security considerations @@ -428,19 +428,32 @@ 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: only per-Object conditions (§2.4.2, §12.8, §12.9)** — - the session reports Object Properties that make a track malformed +- **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. Not detected: the gap rules that - need earlier Objects (a gap covering an Object already received, an Object - inside a gap already communicated, differing Prior Group ID Gaps in a - Group), and §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. + publisher cancelled, the Object not cached. The relay also ends it for two + Prior Group ID Gap values in one Group. Not detected: §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. +- **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 + previously received" as malformed-track conditions, but §2.1 says the first + "is not a protocol error and the Track is not malformed", and lets an Object + go from existing to not existing. The relay follows §2.1: it neither + forwards nor caches such an Object (§9.1 SHOULD NOT) and accepts the covering + gap, keeping any cached copy of the Object it covers (§9.1 makes updating + the cache a MAY). §9.1's specific SHOULD NOT is taken over §9.4's general + "MUST NOT reorder or drop objects received on a multi-object stream". + Checked on live subgroup and datagram Objects against the last 32 Groups; + not in upstream FETCH responses, nor in a session that is not a relay's. The + gap properties are forwarded unchanged, so a subscriber reading §12.9 + literally may still end the track itself. - **LOC properties inside Immutable Properties (§12.7)** — `loc.Properties.Parse` reads only the mutable list, so a LOC Timestamp and the like placed inside Immutable Properties is not found. Filling the fields from there would make diff --git a/pkg/moqt/message/object_properties.go b/pkg/moqt/message/object_properties.go index 5ce173f9..1cfe9433 100644 --- a/pkg/moqt/message/object_properties.go +++ b/pkg/moqt/message/object_properties.go @@ -77,16 +77,32 @@ 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. +// PriorGaps are an Object's Prior Group ID Gap (§12.8) and Prior Object ID +// Gap (§12.9), each if present. +type PriorGaps struct { + Group, Object uint64 + HasGroup, HasObject bool +} + +// ObjectPriorGaps returns the Prior Group and Object ID Gaps in an Object's raw +// Properties, searching Immutable Properties too (§12.7). Properties that +// [CheckObjectProperties] rejects for a reason it can see without the Object's +// IDs carry none. // // Must not allocate: per-Object path. -func PriorObjectIDGap(raw []byte) (uint64, bool) { +func ObjectPriorGaps(raw []byte) PriorGaps { var c objectPropertiesCheck - if c.walk(raw, false) != nil || c.objectGaps != 1 { - return 0, false + if c.walk(raw, false) != nil { + return PriorGaps{} + } + return PriorGaps{ + Group: c.groupGap, Object: c.objectGap, + HasGroup: c.groupGaps == 1, HasObject: c.objectGaps == 1, } - return c.objectGap, true +} + +// PriorObjectIDGap is [ObjectPriorGaps] reduced to the Prior Object ID Gap. +func PriorObjectIDGap(raw []byte) (uint64, bool) { + g := ObjectPriorGaps(raw) + return g.Object, g.HasObject } diff --git a/pkg/moqt/message/properties_test.go b/pkg/moqt/message/properties_test.go index 5b3f8219..2a02d18c 100644 --- a/pkg/moqt/message/properties_test.go +++ b/pkg/moqt/message/properties_test.go @@ -315,6 +315,34 @@ func TestPriorObjectIDGap(t *testing.T) { } } +// TestObjectPriorGaps: both gaps are read in one walk, from either list (§12.7). +func TestObjectPriorGaps(t *testing.T) { + for _, tc := range []struct { + name string + raw []byte + want PriorGaps + }{ + {"empty", nil, PriorGaps{}}, + {"both", AppendTrackProperties([]wire.KVPair{ + kv(PropertyPriorGroupIDGap, 2), kv(PropertyPriorObjectIDGap, 3), + }), PriorGaps{Group: 2, Object: 3, HasGroup: true, HasObject: true}}, + {"group inside Immutable", AppendTrackProperties([]wire.KVPair{ + immutable(kv(PropertyPriorGroupIDGap, 0)), + }), PriorGaps{HasGroup: true}}, + {"malformed", AppendTrackProperties([]wire.KVPair{ + kv(PropertyPriorGroupIDGap, 1), kv(PropertyPriorGroupIDGap, 1), + }), PriorGaps{}}, + } { + if got := ObjectPriorGaps(tc.raw); got != tc.want { + t.Errorf("%s: ObjectPriorGaps = %+v, want %+v", tc.name, got, tc.want) + } + } + raw := AppendTrackProperties([]wire.KVPair{kv(PropertyPriorGroupIDGap, 2)}) + if n := testing.AllocsPerRun(10, func() { ObjectPriorGaps(raw) }); n != 0 { + t.Errorf("ObjectPriorGaps allocates %v times, want 0", n) + } +} + func BenchmarkCheckObjectProperties(b *testing.B) { raw := AppendTrackProperties([]wire.KVPair{ kv(0x40, 1), kv(PropertyPriorObjectIDGap, 1), immutable(kv(0x42, 7)), diff --git a/pkg/relay/gap_properties_test.go b/pkg/relay/gap_properties_test.go new file mode 100644 index 00000000..42912d58 --- /dev/null +++ b/pkg/relay/gap_properties_test.go @@ -0,0 +1,186 @@ +package relay_test + +import ( + "context" + "testing" + "time" + + "github.com/floatdrop/moq-go/pkg/moqt/message" + "github.com/floatdrop/moq-go/pkg/moqt/session" + "github.com/floatdrop/moq-go/pkg/relay" +) + +// got is an Object's location as the subscriber received it. +type got struct{ group, id uint64 } + +// receiveDatagrams emits the location of each datagram sess receives. +func receiveDatagrams(ctx context.Context, sess *session.Session, out chan<- got) { + for { + d, err := sess.ReceiveDatagram(ctx) + if err != nil { + return + } + out <- got{d.GroupID, d.ObjectID} + } +} + +// receiveSubgroupObjects emits the location of each Object on every subgroup +// stream sess accepts. +func receiveSubgroupObjects(ctx context.Context, sess *session.Session, out chan<- got) { + for { + ds, err := sess.AcceptDataStream(ctx) + if err != nil { + return + } + sg, ok := ds.(*session.IncomingSubgroupStream) + if !ok { + return + } + go func() { + for { + o, err := sg.ReadDecoded() + if err != nil { + return + } + out <- got{o.GroupID, o.ObjectID} + } + }() + } +} + +// TestRelay_ObjectInsideAnnouncedGapDropped: an Object inside a Prior Object or +// Group ID Gap announced earlier is known not to exist (§2.1), so the relay +// neither forwards nor caches it (§9.1), and the track goes on. +func TestRelay_ObjectInsideAnnouncedGapDropped(t *testing.T) { + t.Parallel() + type object struct { + group, id uint64 + props []byte + } + for _, tc := range []struct { + name string + datagram bool + // first announces the gap; dropped is inside it; next follows, on the + // same downstream stream as dropped for subgroups. + first, dropped, next object + }{ + { + name: "subgroup, object gap", + first: object{1, 3, priorGap(message.PropertyPriorObjectIDGap, 2)}, + dropped: object{1, 2, nil}, next: object{1, 4, nil}, + }, + { + name: "datagram, object gap", datagram: true, + first: object{1, 3, priorGap(message.PropertyPriorObjectIDGap, 2)}, + dropped: object{1, 2, nil}, next: object{1, 4, nil}, + }, + { + name: "datagram, group gap", datagram: true, + first: object{5, 0, priorGap(message.PropertyPriorGroupIDGap, 2)}, + dropped: object{4, 0, nil}, next: object{5, 1, nil}, + }, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + pubSess, teardown := connectRelay(t, relay.Config{}) + defer teardown() + pub := publishVideoTrack(t, pubSess, "cam1", 1) + subSess := newCam1Subscriber(t, pubSess) + + received := make(chan got, 8) + if tc.datagram { + go receiveDatagrams(t.Context(), subSess, received) + } else { + go receiveSubgroupObjects(t.Context(), subSess, received) + } + await := func() got { + t.Helper() + select { + case g := <-received: + return g + case <-time.After(2 * time.Second): + t.Fatal("no Object forwarded") + return got{} + } + } + send := func(sg *session.OutgoingSubgroupStream, o object) { + t.Helper() + if tc.datagram { + d := &message.ObjectDatagram{ + TrackAlias: 1, + GroupID: o.group, + ObjectID: o.id, + Properties: o.props, + ObjectPayload: []byte("x"), + } + if len(o.props) > 0 { + d.Type = message.DatagramPropertiesBit + } + if err := pubSess.SendDatagram(d); err != nil { + t.Fatalf("SendDatagram: %v", err) + } + return + } + if err := sg.WriteObjectAt( + o.id, + &message.SubgroupObject{Properties: o.props, Payload: []byte("x")}, + ); err != nil { + t.Fatalf("WriteObjectAt: %v", err) + } + } + open := func(group, subgroup uint64) *session.OutgoingSubgroupStream { + t.Helper() + if tc.datagram { + return nil + } + sg, err := pub.OpenSubgroup(message.SubgroupHeader{ + SubgroupIDMode: message.SubgroupIDExplicit, GroupID: group, SubgroupID: subgroup, Properties: true, + }) + if err != nil { + t.Fatalf("OpenSubgroup: %v", err) + } + t.Cleanup(func() { _ = sg.Close() }) + return sg + } + + send(open(tc.first.group, 0), tc.first) + if g := await(); g != (got{tc.first.group, tc.first.id}) { + t.Fatalf("forwarded %+v, want the first Object", g) + } + // Datagrams are read in order; so are Objects on one stream. + later := open(tc.dropped.group, 1) + send(later, tc.dropped) + send(later, tc.next) + if g := await(); g != (got{tc.next.group, tc.next.id}) { + t.Fatalf("forwarded %+v, want %+v: the Object inside the gap is not forwarded", g, tc.next) + } + + // subSess's reader would take the FETCH response stream. + fetcher := dialAnotherClient(t, pubSess) + fetchReq, err := fetcher.Fetch(t.Context(), &message.Fetch{ + Namespace: ns("video"), Name: []byte("cam1"), + Parameters: message.Parameters{fetchRangeFilter( + message.Location{}, + message.Location{Group: tc.next.group, Object: tc.next.id}, + )}, + }) + if err != nil { + t.Fatalf("Fetch: %v", err) + } + defer fetchReq.Close() + servedNext := false + for _, e := range collectFetchElems(t, fetcher, message.GroupOrderAscending, 2*time.Second) { + if e.Marker { + continue + } + if e.Group == tc.dropped.group && e.Object == tc.dropped.id { + t.Fatalf("FETCH served Object %d of Group %d, inside the announced gap", e.Object, e.Group) + } + servedNext = servedNext || e.Group == tc.next.group && e.Object == tc.next.id + } + if !servedNext { + t.Fatal("FETCH did not serve the Object after the gap") + } + }) + } +} diff --git a/pkg/relay/handler_datagram.go b/pkg/relay/handler_datagram.go index c27e79c2..91f1aab9 100644 --- a/pkg/relay/handler_datagram.go +++ b/pkg/relay/handler_datagram.go @@ -60,8 +60,14 @@ func (h *sessionHandler) handleDatagram(ctx context.Context, d *message.ObjectDa d.PublisherPriority = in.DefaultPublisherPriority } - // §2.1: the first copy of {GroupID, ObjectID} wins. - if !entry.ClaimDelivered(d.GroupID, d.ObjectID) { + // §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)) + if err != nil { + h.endMalformedTrack(ctx, entry, h.sess, err) + return + } + if !fresh { return } diff --git a/pkg/relay/handler_fanout.go b/pkg/relay/handler_fanout.go index 35d23114..a7d0ea75 100644 --- a/pkg/relay/handler_fanout.go +++ b/pkg/relay/handler_fanout.go @@ -189,7 +189,7 @@ func (h *sessionHandler) resolveInboundTrack( // §9.3: inbound streams carrying the same (GroupID, SubgroupID) share one // outbound writer per subscriber (§2.2: a Subgroup is not split across // streams), and [registry.TrackEntry.ClaimDelivered] drops duplicate objects -// (§2.1). Each writer is a [subgroupWriter] goroutine behind a bounded queue. +// (§9.3) and ones an announced gap says do not exist (§2.1, §9.1). Each writer is a [subgroupWriter] goroutine behind a bounded queue. // The inbound FIN-vs-reset reaches the outbound streams only when the last // contributor leaves. func (h *sessionHandler) runFanout(ctx context.Context, stream *session.IncomingSubgroupStream) { @@ -305,6 +305,13 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming // §2.4.2). terminalSeen bool ) + // malformed ends the track for an Object read from stream (§2.4.2); the + // Object is neither cached nor forwarded. + malformed := func(err error) { + stream.Cancel(moqt.StreamResetMalformedTrack) + inboundReset, inboundResetCode = true, moqt.StreamResetMalformedTrack + h.endMalformedTrack(ctx, entry, h.sess, err) + } for { obj, err := pending, error(nil) @@ -322,10 +329,7 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming return } if errors.Is(err, session.ErrMalformedTrack) { - stream.Cancel(moqt.StreamResetMalformedTrack) - inboundReset = true - inboundResetCode = moqt.StreamResetMalformedTrack - h.endMalformedTrack(ctx, entry, h.sess, err) + malformed(err) return } h.log.LogAttrs(ctx, slog.LevelDebug, "fanout: inbound ReadObject failed", @@ -340,11 +344,7 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming h.log.LogAttrs(ctx, slog.LevelDebug, "fanout: object after EndOfGroup/EndOfTrack — malformed track", slog.Uint64("group", hdr.GroupID), slog.Uint64("subgroup", hdr.SubgroupID)) - stream.Cancel(moqt.StreamResetMalformedTrack) - inboundReset = true - inboundResetCode = moqt.StreamResetMalformedTrack - h.endMalformedTrack(ctx, entry, h.sess, - fmt.Errorf("object after END_OF_GROUP / END_OF_TRACK in Group %d", hdr.GroupID)) + malformed(fmt.Errorf("object after END_OF_GROUP / END_OF_TRACK in Group %d", hdr.GroupID)) return } @@ -356,9 +356,15 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming // Tracked whether or not this copy wins the dedup claim below. terminal := obj.IsTerminal() - // §2.1: the first upstream to deliver {GroupID, ObjectID} forwards it. - // Outside sg.Mu, so dedup losers never touch the writer set. - if !entry.ClaimDelivered(hdr.GroupID, objectID) { + // §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)) + if err != nil { + malformed(err) + return + } + if !fresh { if terminal { terminalSeen = true } @@ -929,9 +935,8 @@ func (w *subgroupWriter) reopenCause( // 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). 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. +// covers prevID says prevID no longer exists (§2.1), which shows nothing about +// the order. func isNextObject(fwd fwdObject, prevID uint64, dropped bool) bool { if fwd.absID == prevID+1 { return true diff --git a/pkg/relay/internal/registry/track_entry.go b/pkg/relay/internal/registry/track_entry.go index a1716195..b014d137 100644 --- a/pkg/relay/internal/registry/track_entry.go +++ b/pkg/relay/internal/registry/track_entry.go @@ -1,6 +1,7 @@ package registry import ( + "fmt" "maps" "slices" "sync" @@ -26,7 +27,7 @@ import ( // overlapping sessions while migrating WiFi → cellular, // - redundant origins (N-redundant encoders) used for live-media // reliability — the relay deduplicates objects by {GroupID, ObjectID} -// (§2.1) via [TrackEntry.ClaimDelivered] so each object is forwarded +// (§9.3) via [TrackEntry.ClaimDelivered] so each object is forwarded // downstream exactly once. // // Concurrency: @@ -138,7 +139,12 @@ type TrackEntry struct { // group is treated as already-delivered (a peer lagging by that many groups // is beyond any useful reorder window). deliveredMax/HasMax track the largest // group seen, for the pruning window. - delivered map[uint64]map[uint64]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 deliveredMax uint64 deliveredHasMax bool @@ -221,51 +227,123 @@ func (e *TrackEntry) CopySubgroups() []*SharedSubgroup { // few groups behind) while keeping per-track dedup memory bounded. const deliveredGroupWindow = 32 -// ClaimDelivered is the §2.1 dedup gate across multiple upstream publishers. It -// records (group, object) as forwarded and reports whether the caller is the -// first to do so (true → forward it) or it was already forwarded by a peer +// ClaimDelivered is the dedup gate across multiple upstream publishers (§9.3). +// It records (group, object) as forwarded and reports whether the caller is +// the first to do so (true → forward it) or it was already forwarded by a peer // upstream (false → drop it). The ledger persists on the entry (not on a // per-Subgroup structure) and is independent of the size-bounded Object Cache, // so redundant streams that do not temporally overlap, or peers lagging by more // than the cache capacity, still dedup correctly. Memory is bounded to the most // recent [deliveredGroupWindow] groups. -func (e *TrackEntry) ClaimDelivered(group, object uint64) bool { +// +// gaps are the Object's Prior Group and Object ID Gaps (§12.8, §12.9), which +// the ledger records for the whole track. An Object inside a gap announced +// earlier is known not to exist, and that is permanent (§2.1): false, since a +// caching relay "SHOULD NOT cache or forward" it (§9.1). A gap covering an +// Object already received is accepted: an Object may go from existing to not +// existing (§2.1). +// +// Interpretation: §12.8 and §12.9 list both cases as making the track +// malformed, but §2.1 says the first "is not a protocol error and the Track is +// not malformed"; the relay follows §2.1 and §9.1 for both. It also takes +// §9.1's specific SHOULD NOT forward over §9.4's general "MUST NOT reorder or +// 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) { e.deliveredMu.Lock() defer e.deliveredMu.Unlock() - if e.delivered == nil { - e.delivered = make(map[uint64]map[uint64]struct{}) + // An object from a group already aged out of the window is treated as + // already delivered — a peer lagging that far behind is past any useful + // reorder window, and re-forwarding it would be a large out-of-order break. + if e.deliveredHasMax && group <= e.deliveredMax && e.deliveredMax-group >= deliveredGroupWindow { + return false, nil + } + g := e.delivered[group] + if gaps.HasGroup && g != nil && g.hasGroupGap && g.groupGap != gaps.Group { + 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 e.announcedAbsentLocked(g, group, object) { + return false, nil } + if e.delivered == nil { + e.delivered = make(map[uint64]*deliveredGroup) + } // Advance the window when a newer group appears, pruning groups that have // fallen out of it. if !e.deliveredHasMax || group > e.deliveredMax { e.deliveredMax = group e.deliveredHasMax = true - for g := range e.delivered { - if e.deliveredMax-g >= deliveredGroupWindow { - delete(e.delivered, g) - } - } + maps.DeleteFunc(e.delivered, func(id uint64, _ *deliveredGroup) bool { + return e.deliveredMax-id >= deliveredGroupWindow + }) + e.groupGaps = slices.DeleteFunc(e.groupGaps, func(r idRange) bool { + return e.deliveredMax-r.hi >= deliveredGroupWindow + }) } - // An object from a group already aged out of the window is treated as - // already delivered — a peer lagging that far behind is past any useful - // reorder window, and re-forwarding it would be a large out-of-order break. - if group <= e.deliveredMax && e.deliveredMax-group >= deliveredGroupWindow { - return false + if g == nil { + g = &deliveredGroup{objects: make(map[uint64]struct{})} + e.delivered[group] = g } + e.recordGapsLocked(g, group, object, gaps) + if _, ok := g.objects[object]; ok { + return false, nil + } + g.objects[object] = struct{}{} + return true, nil +} + +// 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{} + // 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. + groupGap uint64 + hasGroupGap bool +} + +// idRange is an inclusive range of Group or Object IDs. +type idRange struct{ lo, hi uint64 } + +func (r idRange) contains(id uint64) bool { return r.lo <= id && id <= r.hi } + +// gapRange is the IDs a Prior Group or Object ID Gap of n on id announces +// absent (§12.8, §12.9); none for n == 0. The session layer has already +// rejected an n above id ([message.CheckObjectProperties]). +func gapRange(id, n uint64) (idRange, bool) { + return idRange{lo: id - n, hi: id - 1}, n > 0 +} - set := e.delivered[group] - if set == nil { - set = make(map[uint64]struct{}) - e.delivered[group] = set +// announcedAbsentLocked reports whether the Object at (group, object), g +// being its Group's ledger entry, is inside a gap announced earlier. The +// Group's scan grows with its distinct Object ID gaps, bounded only by the +// window; publishers rarely send many, so the cost is accepted. +func (e *TrackEntry) announcedAbsentLocked(g *deliveredGroup, group, object uint64) bool { + if slices.ContainsFunc(e.groupGaps, func(r idRange) bool { return r.contains(group) }) { + return true } - if _, ok := set[object]; ok { - return false + return g != nil && slices.ContainsFunc(g.objectGaps, func(r idRange) bool { return r.contains(object) }) +} + +// recordGapsLocked records the gaps an Object at (group, object) announced. +func (e *TrackEntry) recordGapsLocked(g *deliveredGroup, group, object uint64, gaps message.PriorGaps) { + if r, ok := gapRange(object, gaps.Object); gaps.HasObject && ok && !slices.Contains(g.objectGaps, r) { + g.objectGaps = append(g.objectGaps, r) + } + if gaps.HasGroup && !g.hasGroupGap { + g.groupGap, g.hasGroupGap = gaps.Group, true + if r, ok := gapRange(group, gaps.Group); ok { + e.groupGaps = append(e.groupGaps, r) + } } - set[object] = struct{}{} - return true } // ReleaseSubgroup drops one contributor from (group, subgroup) and reports diff --git a/pkg/relay/internal/registry/track_test.go b/pkg/relay/internal/registry/track_test.go index cf3a84d7..783d5937 100644 --- a/pkg/relay/internal/registry/track_test.go +++ b/pkg/relay/internal/registry/track_test.go @@ -1,6 +1,7 @@ package registry_test import ( + "errors" "sync" "sync/atomic" "testing" @@ -31,40 +32,138 @@ func TestTrackEntry_ClaimDelivered(t *testing.T) { r := registry.NewTrackRegistry() e := r.GetOrCreate(newTestTrackName("dedup")) + claim := func(group, object uint64) bool { + t.Helper() + fresh, err := e.ClaimDelivered(group, object, message.PriorGaps{}) + if err != nil { + t.Fatalf("ClaimDelivered(%d, %d): %v", group, object, err) + } + return fresh + } // First sighting wins; an exact repeat loses. - if !e.ClaimDelivered(0, 5) { + if !claim(0, 5) { t.Fatal("first ClaimDelivered(0,5) should win") } - if e.ClaimDelivered(0, 5) { + if claim(0, 5) { t.Fatal("repeat ClaimDelivered(0,5) should lose") } // A gap-fill in the same group (object 5 already seen, 2 not) is independent. - if !e.ClaimDelivered(0, 2) { + if !claim(0, 2) { t.Fatal("ClaimDelivered(0,2) should win — distinct object in a seen group") } // A different group is independent. - if !e.ClaimDelivered(1, 5) { + if !claim(1, 5) { t.Fatal("ClaimDelivered(1,5) should win — distinct group") } // Advance the group far enough that group 0 ages out of the window; a late // straggler from group 0 must then be treated as already delivered. - if !e.ClaimDelivered(1000, 0) { + if !claim(1000, 0) { t.Fatal("ClaimDelivered(1000,0) should win") } - if e.ClaimDelivered(0, 9) { + if claim(0, 9) { t.Fatal("ClaimDelivered(0,9) should lose — group 0 has aged out of the dedup window") } // The current group still dedups normally after the window advanced. - if !e.ClaimDelivered(1000, 1) { + if !claim(1000, 1) { t.Fatal("ClaimDelivered(1000,1) should win in the current group") } - if e.ClaimDelivered(1000, 1) { + if claim(1000, 1) { t.Fatal("repeat ClaimDelivered(1000,1) should lose") } } +// TestTrackEntry_ClaimDeliveredGapProperties pins how the ledger treats the +// Prior Group and Object ID Gaps (§12.8, §12.9) of earlier Objects. An Object +// inside an announced gap is known not to exist (§2.1) and dropped (§9.1); a gap +// covering a received Object is accepted (§2.1); two Prior Group ID Gap values +// in one Group make the track malformed (§12.8). +func TestTrackEntry_ClaimDeliveredGapProperties(t *testing.T) { + t.Parallel() + type claim struct { + group, object uint64 + gaps message.PriorGaps + } + const ( + forwarded = iota + dropped + malformed + ) + var none message.PriorGaps + objectGap := func(n uint64) message.PriorGaps { return message.PriorGaps{Object: n, HasObject: true} } + groupGap := func(n uint64) message.PriorGaps { return message.PriorGaps{Group: n, HasGroup: true} } + for _, tc := range []struct { + name string + claims []claim // all but the last are forwarded + want int // the last one's outcome + }{ + // §12.9 + {"object gap over missing IDs", []claim{{1, 0, none}, {1, 3, objectGap(2)}}, forwarded}, + {"object gap covering a received Object", []claim{{1, 1, none}, {1, 3, objectGap(2)}}, forwarded}, + {"Object inside an announced object gap", []claim{{1, 3, objectGap(2)}, {1, 2, none}}, dropped}, + {"Object below an announced object gap", []claim{{1, 3, objectGap(2)}, {1, 0, none}}, forwarded}, + {"same ID in another Group", []claim{{1, 3, objectGap(2)}, {2, 2, none}}, forwarded}, + {"a copy with the same object gap", []claim{{1, 3, objectGap(2)}, {1, 3, objectGap(2)}}, dropped}, + // §12.8 + {"group gap over missing Groups", []claim{{1, 0, none}, {4, 0, groupGap(2)}}, forwarded}, + {"group gap covering a received Group", []claim{{2, 5, none}, {4, 0, groupGap(2)}}, forwarded}, + {"Group inside an announced group gap", []claim{{4, 0, groupGap(2)}, {3, 0, none}}, dropped}, + {"Group covered after it arrived", []claim{{2, 5, none}, {4, 0, groupGap(2)}, {2, 6, none}}, dropped}, + {"same group gap twice in a Group", []claim{{4, 0, groupGap(2)}, {4, 1, groupGap(2)}}, forwarded}, + {"different group gaps in a Group", []claim{{4, 0, groupGap(2)}, {4, 1, groupGap(1)}}, malformed}, + {"group gap on one Object of a Group only", []claim{{4, 0, groupGap(2)}, {4, 1, none}}, forwarded}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + 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) + got := forwarded + switch { + case err != nil: + got = malformed + if !errors.Is(err, session.ErrMalformedTrack) { + t.Fatalf("claim %d %+v: %v, want it to wrap session.ErrMalformedTrack", i, c, err) + } + case !fresh: + got = dropped + } + want := forwarded + if i == last { + want = tc.want + } + if got != want { + t.Fatalf("claim %d %+v: outcome %d (err %v), want %d", i, c, got, err, want) + } + } + }) + } +} + +// TestTrackEntry_ClaimDeliveredMalformedLeavesNoTrace: a malformed claim +// records nothing: neither its Object nor its Prior Group ID Gap. +func TestTrackEntry_ClaimDeliveredMalformedLeavesNoTrace(t *testing.T) { + t.Parallel() + 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) + 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 { + 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) +} + // TestTrackRegistry_GetMissingReturnsFalse confirms the unknown-key path of // Get is a clean miss rather than a zero entry. func TestTrackRegistry_GetMissingReturnsFalse(t *testing.T) { diff --git a/pkg/relay/malformed_track_test.go b/pkg/relay/malformed_track_test.go index 8cf3602f..89a36707 100644 --- a/pkg/relay/malformed_track_test.go +++ b/pkg/relay/malformed_track_test.go @@ -11,6 +11,7 @@ import ( "github.com/floatdrop/moq-go/pkg/moqt" "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" ) @@ -49,6 +50,40 @@ func requireUpstreamCancelled(t *testing.T, pub *session.Publication) { } } +// gapObject is an Object for [sendGapObjects]. +type gapObject struct { + group, subgroup, object uint64 + props []byte +} + +// sendGapObjects sends each of objs on its own subgroup stream. The relay may +// read the streams in either order; each case using it is malformed in both. +func sendGapObjects(t *testing.T, pubSess *session.Session, alias uint64, objs []gapObject) { + t.Helper() + for _, o := range objs { + sg, err := pubSess.OpenSubgroup(message.SubgroupHeader{ + SubgroupIDMode: message.SubgroupIDExplicit, TrackAlias: alias, + GroupID: o.group, SubgroupID: o.subgroup, Properties: true, + }) + if err != nil { + t.Errorf("OpenSubgroup: %v", err) + return + } + if err := sg.WriteObjectAt( + o.object, + &message.SubgroupObject{Properties: o.props, Payload: []byte("x")}, + ); err != nil { + t.Errorf("WriteObjectAt: %v", err) + } + _ = 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}}) +} + // TestRelay_MalformedObjectEndsTrack: each kind of malformed Object ends the // downstream subscription with MALFORMED_TRACK and cancels the upstream one // (§2.4.2). @@ -82,6 +117,23 @@ func TestRelay_MalformedObjectEndsTrack(t *testing.T) { _ = sg.WriteObject(&message.SubgroupObject{Payload: []byte("x")}) _ = sg.Close() }}, + // §12.8: two Prior Group ID Gap values in one Group. + {"different group gaps in a Group", func(t *testing.T, pubSess *session.Session, alias uint64) { + sendGapObjects(t, pubSess, alias, []gapObject{ + {5, 0, 0, priorGap(message.PropertyPriorGroupIDGap, 2)}, + {5, 1, 1, priorGap(message.PropertyPriorGroupIDGap, 1)}, + }) + }}, + {"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{ + Type: message.DatagramPropertiesBit, TrackAlias: alias, GroupID: 5, ObjectID: uint64(i), + Properties: priorGap(message.PropertyPriorGroupIDGap, gap), ObjectPayload: []byte("x"), + }); err != nil { + t.Errorf("SendDatagram: %v", err) + } + } + }}, {"datagram Object Properties", func(t *testing.T, pubSess *session.Session, alias uint64) { if err := pubSess.SendDatagram(&message.ObjectDatagram{ Type: message.DatagramPropertiesBit, TrackAlias: alias, GroupID: 1, diff --git a/pkg/relay/next_object_test.go b/pkg/relay/next_object_test.go index a837f57e..de3717b3 100644 --- a/pkg/relay/next_object_test.go +++ b/pkg/relay/next_object_test.go @@ -217,8 +217,8 @@ 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. + // A gap covering an Object already sent says it no longer exists + // (§2.1), which shows nothing about the order. {"gap covering the previous Object", priorObjectIDGap(3), [][]uint64{{0}, {3}}}, } { t.Run(tc.name, func(t *testing.T) {