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
39 changes: 26 additions & 13 deletions STATUS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down
32 changes: 24 additions & 8 deletions pkg/moqt/message/object_properties.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
28 changes: 28 additions & 0 deletions pkg/moqt/message/properties_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)),
Expand Down
186 changes: 186 additions & 0 deletions pkg/relay/gap_properties_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
})
}
}
10 changes: 8 additions & 2 deletions pkg/relay/handler_datagram.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down
Loading
Loading