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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 12 additions & 8 deletions STATUS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand Down Expand Up @@ -433,13 +433,17 @@ 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
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.
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.
- **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
Expand Down
21 changes: 21 additions & 0 deletions pkg/moqt/message/object_properties.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
27 changes: 27 additions & 0 deletions pkg/moqt/message/properties_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package message

import (
"bytes"
"testing"
"time"

Expand Down Expand Up @@ -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)),
Expand Down
177 changes: 177 additions & 0 deletions pkg/relay/duplicate_consistency_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,177 @@
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},
// 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()
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)
})
}
}
12 changes: 12 additions & 0 deletions pkg/relay/handler_datagram.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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
}

Expand Down
60 changes: 60 additions & 0 deletions pkg/relay/handler_duplicate.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
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). 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 {
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:
return nil
case !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)
}
Loading
Loading