From d607d53bb35b666fdb1eb80afea1da580584d465 Mon Sep 17 00:00:00 2001 From: Vsevolod Strukchinsky Date: Sat, 26 Sep 2026 13:21:01 +0500 Subject: [PATCH 1/2] =?UTF-8?q?fix(relay):=20a=20duplicate=20that=20differ?= =?UTF-8?q?s=20from=20the=20cached=20copy=20makes=20the=20track=20malforme?= =?UTF-8?q?d=20(=C2=A79.1)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit §9.1: "An endpoint that receives a duplicate Object with a different Forwarding Preference, Subgroup ID, Priority or Payload MUST treat the track as Malformed" (also §2.4.2 items 6 and 7). The relay dropped a copy that lost the dedup claim without looking at it. It now compares the copy with the first one in the cache, on the subgroup and datagram paths, and ends the track (PUBLISH_DONE MALFORMED_TRACK) on a difference in those fields or in Immutable Properties (§12.7, byte for byte; one copy having them and the other not counts). Mutable Properties may differ (§9.1). Payloads are compared only when both copies are Normal: Normal may become End of Group or Track (§9.1), and the reverse is a late Object (§2.1). Decisions made with the user; the limitation is documented: a copy is not compared once the first one left the cache. New: message.ImmutableProperties (zero-alloc). runFanout sets terminalSeen once, which is equivalent and keeps it under gocyclo. Tests whose "duplicates" differed in Payload or Subgroup are corrected: §9.1 now makes them malformed. TestRelay_DuplicateConsistency's seven malformed cases were verified red with the check disabled; its guards fail when all Properties, or Payloads regardless of Status, are compared. Co-Authored-By: Claude Opus 5.5 (1M context) --- STATUS.md | 10 +- pkg/moqt/message/object_properties.go | 21 +++ pkg/moqt/message/properties_test.go | 27 ++++ pkg/relay/duplicate_consistency_test.go | 165 ++++++++++++++++++++++ pkg/relay/handler_datagram.go | 12 ++ pkg/relay/handler_duplicate.go | 55 ++++++++ pkg/relay/handler_fanout.go | 23 +-- pkg/relay/handler_fanout_multipub_test.go | 4 +- pkg/relay/handler_fanout_test.go | 16 +-- 9 files changed, 311 insertions(+), 22 deletions(-) create mode 100644 pkg/relay/duplicate_consistency_test.go create mode 100644 pkg/relay/handler_duplicate.go diff --git a/STATUS.md b/STATUS.md index 371a7e68..f48cda97 100644 --- a/STATUS.md +++ b/STATUS.md @@ -142,7 +142,7 @@ By package, bottom-up along the dependency stack: | § | Feature | Status | Notes | |-------|--------------------------------------|--------|-------| -| 9.1 | Caching relays | DONE | LRU+TTL object cache (`cache/cache.go`); updates limited to non-existence/properties. | +| 9.1 | Caching relays | DONE | LRU+TTL object cache (`cache/cache.go`); updates limited to non-existence/properties. A duplicate of a cached Object with a different Forwarding Preference, Subgroup ID, Priority or Payload, or different Immutable Properties (§2.4.2, §12.7), ends the track as malformed (`relay/handler_duplicate.go`); not checked once the first copy left the cache — see Limitations. | | 9.2 | Forward handling | DONE | FORWARD flag honoured; Forward=0 pauses delivery. Upstream Forward is set to 1 only when a downstream subscriber forwards, else the relay pauses it (Forward=0) and resumes on the first forwarding subscriber. | | 9.3 | Multiple publishers | DONE | Per-track upstreams; dedup by `{GroupID, ObjectID}`. Upstreams of one Subgroup share one downstream stream per subscriber, with the first one's SUBGROUP_HEADER; a later one's Object Properties reopen it with PROPERTIES set, so none are dropped (§2.5). Like a §11.4.3 gap reopen, the reset keeps already-written Objects only where RESET_STREAM_AT is in use; always setting PROPERTIES would avoid it at a byte per Object. | | 9.4 | Subscriber interactions | DONE | Upstream subscription established before SUBSCRIBE_OK; aggregation. | @@ -433,8 +433,12 @@ Known protocol gaps, roughly ordered by how load-bearing they are: (`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. Not detected: §2.4.2's list other - than an Object after END_OF_GROUP on the same stream. A downstream FETCH + Prior Group ID Gap values in one Group, and for a duplicate that differs from + the cached first copy (§9.1; Payloads 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 diff --git a/pkg/moqt/message/object_properties.go b/pkg/moqt/message/object_properties.go index 1cfe9433..42f7c418 100644 --- a/pkg/moqt/message/object_properties.go +++ b/pkg/moqt/message/object_properties.go @@ -106,3 +106,24 @@ func PriorObjectIDGap(raw []byte) (uint64, bool) { g := ObjectPriorGaps(raw) return g.Object, g.HasObject } + +// ImmutableProperties returns the value of the Immutable Properties (§12.7) in +// an Object's raw Properties, serialized as received, and whether there is one. +// Properties that do not parse carry none. +// +// Must not allocate. +func ImmutableProperties(raw []byte) ([]byte, bool) { + r := wire.NewReader(raw) + var prev uint64 + for !r.Empty() { + kv, next, err := r.KVPairView(prev) + if err != nil { + return nil, false + } + if kv.Type == PropertyImmutableProperties { + return kv.ByteVal, true + } + prev = next + } + return nil, false +} diff --git a/pkg/moqt/message/properties_test.go b/pkg/moqt/message/properties_test.go index 2a02d18c..847ec682 100644 --- a/pkg/moqt/message/properties_test.go +++ b/pkg/moqt/message/properties_test.go @@ -1,6 +1,7 @@ package message import ( + "bytes" "testing" "time" @@ -343,6 +344,32 @@ func TestObjectPriorGaps(t *testing.T) { } } +// TestImmutableProperties: the Immutable Properties value is returned as sent, +// its serialization included (§12.7). +func TestImmutableProperties(t *testing.T) { + inner := AppendTrackProperties([]wire.KVPair{kv(0x40, 1)}) + for _, tc := range []struct { + name string + raw []byte + want []byte + ok bool + }{ + {"empty", nil, nil, false}, + {"mutable only", AppendTrackProperties([]wire.KVPair{kv(0x40, 1)}), nil, false}, + {"present", AppendTrackProperties([]wire.KVPair{kv(0x42, 2), immutable(kv(0x40, 1))}), inner, true}, + {"unparseable", []byte{0x02}, nil, false}, + } { + got, ok := ImmutableProperties(tc.raw) + if ok != tc.ok || !bytes.Equal(got, tc.want) { + t.Errorf("%s: ImmutableProperties = (%x, %v), want (%x, %v)", tc.name, got, ok, tc.want, tc.ok) + } + } + raw := AppendTrackProperties([]wire.KVPair{immutable(kv(0x40, 1))}) + if n := testing.AllocsPerRun(10, func() { ImmutableProperties(raw) }); n != 0 { + t.Errorf("ImmutableProperties 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/duplicate_consistency_test.go b/pkg/relay/duplicate_consistency_test.go new file mode 100644 index 00000000..13c6cacf --- /dev/null +++ b/pkg/relay/duplicate_consistency_test.go @@ -0,0 +1,165 @@ +package relay_test + +import ( + "testing" + "time" + + "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" +) + +// objCopy is one publisher's copy of Object {1, 0}. +type objCopy struct { + datagram bool + subgroup uint64 + priority uint8 // inline when set; else the track default + status uint64 + props []byte + payload string +} + +// sendCopy sends c as Object {1, 0} from sess on alias and, if withNext, a +// further Object: {1, 1} on the same subgroup stream, or datagram {2, 0}; the +// relay reads them in order. Only the copy's send must succeed: a malformed one +// gets the stream cancelled. +func sendCopy(t *testing.T, sess *session.Session, alias uint64, c objCopy, withNext bool) { + t.Helper() + if c.datagram { + ds := []*message.ObjectDatagram{ + {GroupID: 1, ObjectStatus: c.status, Properties: c.props, ObjectPayload: []byte(c.payload)}, + {GroupID: 2, ObjectPayload: []byte("next")}, + } + if !withNext { + ds = ds[:1] + } + for _, d := range ds { + d.TrackAlias = alias + d.Type = message.DatagramDefaultPriorityBit + if len(d.Properties) > 0 { + d.Type |= message.DatagramPropertiesBit + } + if d.ObjectStatus != 0 { + d.Type |= message.DatagramStatusBit + } + if err := sess.SendDatagram(d); err != nil { + t.Errorf("SendDatagram: %v", err) + } + } + return + } + sg, err := sess.OpenSubgroup(message.SubgroupHeader{ + SubgroupIDMode: message.SubgroupIDExplicit, TrackAlias: alias, GroupID: 1, SubgroupID: c.subgroup, + Properties: true, InlinePriority: c.priority != 0, PublisherPriority: c.priority, + }) + if err != nil { + t.Errorf("OpenSubgroup: %v", err) + return + } + t.Cleanup(func() { _ = sg.Close() }) + if err := sg.WriteObject(&message.SubgroupObject{ + ObjectStatus: c.status, Properties: c.props, Payload: []byte(c.payload), + }); err != nil { + t.Errorf("WriteObject: %v", err) + } + if withNext { + _ = sg.WriteObject(&message.SubgroupObject{Payload: []byte("next")}) + } +} + +func immutableProps(v uint64) []byte { + return message.AppendTrackProperties([]wire.KVPair{{ + Type: message.PropertyImmutableProperties, + ByteVal: message.AppendTrackProperties([]wire.KVPair{{Type: 0x40, IntVal: v}}), + }}) +} + +func mutableProps(v uint64) []byte { + return message.AppendTrackProperties([]wire.KVPair{{Type: 0x40, IntVal: v}}) +} + +// TestRelay_DuplicateConsistency: a duplicate of a cached Object with a +// different Forwarding Preference, Subgroup ID, Priority or Payload (§9.1), or +// different Immutable Properties (§2.4.2, §12.7), makes the track malformed. +// Mutable Properties may differ (§9.1), and Normal may become End of Group +// (§9.1: existing to not existing). +func TestRelay_DuplicateConsistency(t *testing.T) { + t.Parallel() + sub := objCopy{payload: "a"} + dg := objCopy{datagram: true, payload: "a"} + with := func(c objCopy, f func(*objCopy)) objCopy { f(&c); return c } + for _, tc := range []struct { + name string + first, second objCopy + malformed bool + }{ + {"payload", sub, with(sub, func(c *objCopy) { c.payload = "b" }), true}, + {"priority", sub, with(sub, func(c *objCopy) { c.priority = 7 }), true}, + {"Subgroup ID", sub, with(sub, func(c *objCopy) { c.subgroup = 1 }), true}, + {"forwarding preference", sub, dg, true}, + { + "Immutable Properties", + with(sub, func(c *objCopy) { c.props = immutableProps(1) }), + with(sub, func(c *objCopy) { c.props = immutableProps(2) }), true, + }, + {"Immutable Properties removed", with(sub, func(c *objCopy) { c.props = immutableProps(1) }), sub, true}, + {"datagram payload", dg, with(dg, func(c *objCopy) { c.payload = "b" }), true}, + + {"identical", sub, sub, false}, + { + "mutable Properties", + with(sub, func(c *objCopy) { c.props = mutableProps(1) }), + with(sub, func(c *objCopy) { c.props = mutableProps(2) }), false, + }, + {"identical datagram", dg, dg, false}, + {"Normal, then End of Group", dg, with(dg, func(c *objCopy) { + c.status, c.payload = message.ObjectStatusEndOfGroup, "" + }), false}, + } { + t.Run(tc.name, func(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) + go receiveSubgroupObjects(t.Context(), subSess, received) + await := func(want got) { + t.Helper() + for { + select { + case g := <-received: + if g == want { + return + } + case <-time.After(2 * time.Second): + t.Fatalf("Object %+v not forwarded", want) + } + } + } + + sendCopy(t, pubA, 1, tc.first, false) + await(got{1, 0}) // forwarded, so cached + sendCopy(t, pubB, 2, tc.second, true) + if tc.malformed { + if pd := awaitPublishDone(t, subReq); pd.StatusCode != moqt.PublishDoneMalformedTrack { + t.Fatalf("PUBLISH_DONE %#x, want MALFORMED_TRACK", uint64(pd.StatusCode)) + } + return + } + // B's next Object is read after its copy: forwarded, the track + // survived the copy. + next := got{1, 1} + if tc.second.datagram { + next = got{2, 0} + } + await(next) + }) + } +} diff --git a/pkg/relay/handler_datagram.go b/pkg/relay/handler_datagram.go index 91f1aab9..6ff3a6cd 100644 --- a/pkg/relay/handler_datagram.go +++ b/pkg/relay/handler_datagram.go @@ -7,6 +7,7 @@ import ( "github.com/floatdrop/moq-go/pkg/moqt/message" "github.com/floatdrop/moq-go/pkg/moqt/session" + "github.com/floatdrop/moq-go/pkg/relay/cache" "github.com/floatdrop/moq-go/pkg/relay/internal/registry" ) @@ -68,6 +69,17 @@ func (h *sessionHandler) handleDatagram(ctx context.Context, d *message.ObjectDa return } if !fresh { + if err := checkDuplicate(entry.Cache, &cache.CachedObject{ + GroupID: d.GroupID, + ObjectID: d.ObjectID, + PublisherPriority: d.PublisherPriority, + ForwardingPref: cache.ForwardingDatagram, + Status: d.ObjectStatus, + Properties: d.Properties, + Payload: d.ObjectPayload, + }); err != nil { + h.endMalformedTrack(ctx, entry, h.sess, err) + } return } diff --git a/pkg/relay/handler_duplicate.go b/pkg/relay/handler_duplicate.go new file mode 100644 index 00000000..90fc131d --- /dev/null +++ b/pkg/relay/handler_duplicate.go @@ -0,0 +1,55 @@ +package relay + +import ( + "bytes" + "fmt" + + "github.com/floatdrop/moq-go/pkg/moqt/message" + "github.com/floatdrop/moq-go/pkg/moqt/session" + "github.com/floatdrop/moq-go/pkg/relay/cache" +) + +// checkDuplicate compares dup, a copy that lost the §9.3 dedup claim, with the +// first copy in c. A different Forwarding Preference, Subgroup ID, Priority or +// Payload (§9.1), or different Immutable Properties (§2.4.2, §12.7; one copy +// having them and the other not counts), makes the track malformed: the error +// wraps [session.ErrMalformedTrack]. Mutable Properties may differ (§9.1). +// +// Normal becoming End of Group or End of Track is the existing-to-not-existing +// change §9.1 allows, and the reverse a late Object (§2.1), so Payloads are +// compared only when both copies are Normal. +// +// 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 { + first, ok := c.Get(dup.GroupID, dup.ObjectID) + if !ok || first.IsRangeMarker() { + return nil + } + var field string + switch { + case first.ForwardingPref != dup.ForwardingPref: + field = "Forwarding Preference" + case first.ForwardingPref == cache.ForwardingSubgroup && first.SubgroupID != dup.SubgroupID: + field = "Subgroup ID" + case first.PublisherPriority != dup.PublisherPriority: + field = "Priority" + case first.Status == message.ObjectStatusNormal && dup.Status == message.ObjectStatusNormal && + !bytes.Equal(first.Payload, dup.Payload): + field = "Payload" + case !sameImmutableProperties(first.Properties, dup.Properties): + field = "Immutable Properties" + default: + return nil + } + 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) +} + +// sameImmutableProperties reports whether a and b, raw Object Properties, +// carry the same Immutable Properties, byte for byte, or neither has any. +func sameImmutableProperties(a, b []byte) bool { + av, aok := message.ImmutableProperties(a) + bv, bok := message.ImmutableProperties(b) + return aok == bok && bytes.Equal(av, bv) +} diff --git a/pkg/relay/handler_fanout.go b/pkg/relay/handler_fanout.go index a7d0ea75..bb200b02 100644 --- a/pkg/relay/handler_fanout.go +++ b/pkg/relay/handler_fanout.go @@ -353,8 +353,9 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming pos.seq++ objectID := stream.ObjectID() // resolved by ReadObject (§11.4.2) - // Tracked whether or not this copy wins the dedup claim below. - terminal := obj.IsTerminal() + // Whether or not this copy wins the dedup claim below; the next + // iteration acts on it, so it can be set now. + terminalSeen = obj.IsTerminal() // §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 @@ -365,8 +366,18 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming return } if !fresh { - if terminal { - terminalSeen = true + if err := checkDuplicate(entry.Cache, &cache.CachedObject{ + GroupID: hdr.GroupID, + ObjectID: objectID, + SubgroupID: hdr.SubgroupID, + PublisherPriority: hdr.PublisherPriority, + ForwardingPref: cache.ForwardingSubgroup, + Status: obj.ObjectStatus, + Properties: obj.Properties, + Payload: obj.Payload, + }); err != nil { + malformed(err) + return } continue // redundant copy already forwarded by a peer upstream. } @@ -423,10 +434,6 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming } } sg.Mu.Unlock() - - if terminal { - terminalSeen = true - } } } diff --git a/pkg/relay/handler_fanout_multipub_test.go b/pkg/relay/handler_fanout_multipub_test.go index 499c6e1a..7c997ddc 100644 --- a/pkg/relay/handler_fanout_multipub_test.go +++ b/pkg/relay/handler_fanout_multipub_test.go @@ -123,7 +123,7 @@ func TestFanout_MultiPublisher_DeduplicatesObjects(t *testing.T) { for i := range 3 { if err := bSg.WriteObject(&message.SubgroupObject{ ObjectIDDelta: 0, - Payload: []byte{byte('a' + i)}, + Payload: []byte{byte('A' + i)}, // the same Object: §9.1 forbids another Payload }); err != nil { t.Fatalf("B WriteObject #%d: %v", i, err) } @@ -208,7 +208,7 @@ func TestFanout_MultiPublisher_DedupSurvivesCacheEviction(t *testing.T) { for i := range n { if err := bSg.WriteObject(&message.SubgroupObject{ ObjectIDDelta: 0, - Payload: []byte{byte('a' + i)}, + Payload: []byte{byte('A' + i)}, // the same Object: §9.1 forbids another Payload }); err != nil { t.Fatalf("B WriteObject #%d: %v", i, err) } diff --git a/pkg/relay/handler_fanout_test.go b/pkg/relay/handler_fanout_test.go index 14a76cf0..655c2639 100644 --- a/pkg/relay/handler_fanout_test.go +++ b/pkg/relay/handler_fanout_test.go @@ -732,9 +732,10 @@ func TestFanout_UpdatesTrackEntryLargestObject(t *testing.T) { defer subReq.Close() go drainAll(t.Context(), subSess) - // Phase 1: publish three objects (absIDs 0, 1, 2) on group 4. After - // this the relay's TrackEntry.LargestObject must be {Group: 4, - // Object: 2}. + // Phase 1: publish absIDs 0 and 2 on group 4. After this the relay's + // TrackEntry.LargestObject must be {Group: 4, Object: 2}. Object 1 is + // left for phase 3: sending it twice with another Payload or Subgroup + // would make the track malformed (§9.1). sg, err := pubSess.OpenSubgroup(message.SubgroupHeader{ SubgroupIDMode: message.SubgroupIDExplicit, TrackAlias: publisherAlias, @@ -744,12 +745,9 @@ func TestFanout_UpdatesTrackEntryLargestObject(t *testing.T) { if err != nil { t.Fatalf("OpenSubgroup: %v", err) } - for i := range 3 { - if err := sg.WriteObject(&message.SubgroupObject{ - ObjectIDDelta: 0, - Payload: []byte{byte('A' + i)}, - }); err != nil { - t.Fatalf("WriteObject #%d: %v", i, err) + for _, id := range []uint64{0, 2} { + if err := sg.WriteObjectAt(id, &message.SubgroupObject{Payload: []byte{byte('A' + id)}}); err != nil { + t.Fatalf("WriteObjectAt %d: %v", id, err) } } if err := sg.Close(); err != nil { From e16792a0a789a6363d517d97ce4fa290841c80b0 Mon Sep 17 00:00:00 2001 From: Vsevolod Strukchinsky Date: Sat, 26 Sep 2026 13:24:02 +0500 Subject: [PATCH 2/2] fix(relay): compare a duplicate's Immutable Properties only when both copies are Normal MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review follow-up. Only a Normal Object has a Payload or Properties (§11.2.1.1, §11.2.1.2), but Immutable Properties were compared whatever the Status, so a legal duplicate ended the track: a Normal copy with Immutable Properties becoming End of Group (§9.1), or the reverse, a late Object (§2.1). Forwarding Preference, Subgroup ID and Priority are still compared across a Status change. Also documents that a copy an announced gap says does not exist is still compared with a cached first copy: §2.1 excuses its arrival, not a different content. Both new guard cases of TestRelay_DuplicateConsistency were verified red before the fix. Co-Authored-By: Claude Opus 5.5 (1M context) --- STATUS.md | 20 ++++++++++---------- pkg/relay/duplicate_consistency_test.go | 12 ++++++++++++ pkg/relay/handler_duplicate.go | 11 ++++++++--- 3 files changed, 30 insertions(+), 13 deletions(-) diff --git a/STATUS.md b/STATUS.md index f48cda97..fba65b55 100644 --- a/STATUS.md +++ b/STATUS.md @@ -434,16 +434,16 @@ Known protocol gaps, roughly ordered by how load-bearing they are: 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 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. + 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. - **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 13c6cacf..ece03b6c 100644 --- a/pkg/relay/duplicate_consistency_test.go +++ b/pkg/relay/duplicate_consistency_test.go @@ -117,6 +117,18 @@ func TestRelay_DuplicateConsistency(t *testing.T) { {"Normal, then End of Group", dg, with(dg, func(c *objCopy) { c.status, c.payload = message.ObjectStatusEndOfGroup, "" }), false}, + // Only a Normal Object carries Properties (§11.2.1.2). + {"Normal with Immutable Properties, then End of Group", with(dg, func(c *objCopy) { + c.props = immutableProps(1) + }), with(dg, func(c *objCopy) { + c.status, c.payload = message.ObjectStatusEndOfGroup, "" + }), false}, + // The late Object of §2.1. + {"End of Group, then Normal with Immutable Properties", with(dg, func(c *objCopy) { + c.status, c.payload = message.ObjectStatusEndOfGroup, "" + }), with(dg, func(c *objCopy) { + c.props = immutableProps(1) + }), false}, } { t.Run(tc.name, func(t *testing.T) { t.Parallel() diff --git a/pkg/relay/handler_duplicate.go b/pkg/relay/handler_duplicate.go index 90fc131d..3d238224 100644 --- a/pkg/relay/handler_duplicate.go +++ b/pkg/relay/handler_duplicate.go @@ -16,9 +16,13 @@ import ( // wraps [session.ErrMalformedTrack]. Mutable Properties may differ (§9.1). // // Normal becoming End of Group or End of Track is the existing-to-not-existing -// change §9.1 allows, and the reverse a late Object (§2.1), so Payloads are +// change §9.1 allows, and the reverse a late Object (§2.1). Only a Normal +// Object has a Payload or Properties (§11.2.1.1, §11.2.1.2), so those are // compared only when both copies are Normal. // +// 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. +// // 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 { @@ -34,8 +38,9 @@ func checkDuplicate(c *cache.ObjectCache, dup *cache.CachedObject) error { field = "Subgroup ID" case first.PublisherPriority != dup.PublisherPriority: field = "Priority" - case first.Status == message.ObjectStatusNormal && dup.Status == message.ObjectStatusNormal && - !bytes.Equal(first.Payload, dup.Payload): + case first.Status != message.ObjectStatusNormal || dup.Status != message.ObjectStatusNormal: + return nil + case !bytes.Equal(first.Payload, dup.Payload): field = "Payload" case !sameImmutableProperties(first.Properties, dup.Properties): field = "Immutable Properties"