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
2 changes: 1 addition & 1 deletion STATUS.md
Original file line number Diff line number Diff line change
Expand Up @@ -219,7 +219,7 @@ By package, bottom-up along the dependency stack:
| 11.4 | Streams (subgroup / fetch) | DONE | Typed in/out subgroup + fetch streams. |
| 11.4.1 | Stream cancellation | DONE | Bidi request-stream termination ends the request (handlers unregister on stream end); the relay sends PUBLISH_DONE on graceful subscription termination rather than abrupt reset. |
| 11.4.2 | Subgroup header + delta object IDs | DONE | All subgroup-ID modes; `ReadDecoded` resolves deltas. A subgroup stream whose Track Alias is not bound yet is held unread (§11.4.2 MAY buffer): `session.Demux` parks it, count-bounded; the relay waits up to 1 s (`IncomingSubgroupStream.AwaitInboundTrack`, at most 32 streams per session, the rest reset with EXCESSIVE_LOAD). The §11.4.2 MUST to give control streams connection flow control first is not met by the bundled transports, so enough early data can delay the SUBSCRIBE_OK until the relay's wait runs out and the streams are reset. |
| 11.4.3 | Closing subgroup streams | DONE | Relay forwards only the next object on a stream (gap → reset+reopen), FINs on clean inbound EOF, resets on inbound reset, resets with MALFORMED_TRACK after a terminal EndOfGroup/EndOfTrack object (§2.4.2), marks reliable boundaries for RESET_STREAM_AT (`SetReliableBoundary`, transport-gated on `EnableStreamResetPartialDelivery`), and resets (not FINs) in-flight subgroups whose group falls out of range after a narrowing REQUEST_UPDATE. A subscription that skipped any Object of the Subgroup other than one before its Start Location (a filter, Forward State 0, a Start raised past Objects already sent, an inbox overflow with EXCESSIVE_LOAD, an expiry) gets resets, never a FIN, on that Subgroup's streams. Objects published before a subscription joined are treated as before its Start. |
| 11.4.3 | Closing subgroup streams | DONE | Relay forwards only the next object on a stream, otherwise reset+reopen: the next object is one ID greater, read next from the same upstream stream (only filtered-out objects between), or covered by its Prior Object ID Gap; an object the relay dropped, or one from another upstream, breaks the run. FINs on clean inbound EOF, resets on inbound reset, resets with MALFORMED_TRACK after a terminal EndOfGroup/EndOfTrack object (§2.4.2), marks reliable boundaries for RESET_STREAM_AT (`SetReliableBoundary`, transport-gated on `EnableStreamResetPartialDelivery`), and resets (not FINs) in-flight subgroups whose group falls out of range after a narrowing REQUEST_UPDATE. A subscription that skipped any Object of the Subgroup other than one before its Start Location (a filter, Forward State 0, a Start raised past Objects already sent, an inbox overflow with EXCESSIVE_LOAD, an expiry) gets resets, never a FIN, on that Subgroup's streams. Objects published before a subscription joined are treated as before its Start. |
| 11.4.4 | Fetch header | DONE | Serialization Flags of 128 or more that are not an End of Range close the session. |
| 11.4.4.1 | Fetch flags | DONE | All subgroup modes + delta/priority/properties/status flags. A first Object that references a prior Object's fields closes the session. |
| 11.4.4.2 | End of range | DONE | Non-existent (0x8C) / unknown (0x10C) handled. An Object after a leading marker that references a prior Subgroup ID or Priority closes the session. |
Expand Down
14 changes: 14 additions & 0 deletions pkg/moqt/message/object_properties.go
Original file line number Diff line number Diff line change
Expand Up @@ -76,3 +76,17 @@ 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.
//
// Must not allocate: per-Object path.
func PriorObjectIDGap(raw []byte) (uint64, bool) {
var c objectPropertiesCheck
if c.walk(raw, false) != nil || c.objectGaps != 1 {
return 0, false
}
return c.objectGap, true
}
30 changes: 30 additions & 0 deletions pkg/moqt/message/properties_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -285,6 +285,36 @@ func TestCheckObjectProperties(t *testing.T) {
}
}

// TestPriorObjectIDGap: the Prior Object ID Gap (§12.9) is found in either list
// (§12.7), and Properties that make the track malformed carry none.
func TestPriorObjectIDGap(t *testing.T) {
for _, tc := range []struct {
name string
raw []byte
gap uint64
hasGap bool
}{
{"empty", nil, 0, false},
{"other properties", AppendTrackProperties([]wire.KVPair{kv(0x40, 1)}), 0, false},
{"mutable", AppendTrackProperties([]wire.KVPair{kv(0x40, 1), kv(PropertyPriorObjectIDGap, 3)}), 3, true},
{"zero", AppendTrackProperties([]wire.KVPair{kv(PropertyPriorObjectIDGap, 0)}), 0, true},
{"inside Immutable", AppendTrackProperties([]wire.KVPair{immutable(kv(PropertyPriorObjectIDGap, 2))}), 2, true},
{"two instances", AppendTrackProperties([]wire.KVPair{
kv(PropertyPriorObjectIDGap, 1), immutable(kv(PropertyPriorObjectIDGap, 1)),
}), 0, false},
{"unparseable", []byte{0x02}, 0, false},
} {
gap, ok := PriorObjectIDGap(tc.raw)
if gap != tc.gap || ok != tc.hasGap {
t.Errorf("%s: PriorObjectIDGap = (%d, %v), want (%d, %v)", tc.name, gap, ok, tc.gap, tc.hasGap)
}
}
raw := AppendTrackProperties([]wire.KVPair{immutable(kv(PropertyPriorObjectIDGap, 2))})
if n := testing.AllocsPerRun(10, func() { PriorObjectIDGap(raw) }); n != 0 {
t.Errorf("PriorObjectIDGap 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
38 changes: 26 additions & 12 deletions pkg/relay/forward_state_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -216,9 +216,11 @@ func TestRelay_SkipBeforeStartKeepsFIN(t *testing.T) {
}
}

// TestRelay_ForwardStateOmissionResetsReopenedStream: after a pause omitted an
// Object, the stream reopened on resume ends with a reset too.
func TestRelay_ForwardStateOmissionResetsReopenedStream(t *testing.T) {
// TestRelay_ForwardStateOmissionKeepsStream: an Object omitted while paused did
// not pass the subscriber's filters (§5.1.5: "The Forward parameter is also a
// type of filter"), so the Object after the resume is the next Object and
// stays on the stream (§11.4.3), which still ends with a reset.
func TestRelay_ForwardStateOmissionKeepsStream(t *testing.T) {
t.Parallel()
pubSess, teardown := connectRelay(t, relay.Config{})
defer teardown()
Expand All @@ -240,26 +242,33 @@ func TestRelay_ForwardStateOmissionResetsReopenedStream(t *testing.T) {
return
}
<-resumed
_ = sg.WriteObject(&message.SubgroupObject{Payload: []byte("2")}) // gap: a new stream
_ = sg.WriteObject(&message.SubgroupObject{Payload: []byte("2")})
_ = sg.Close()
}()

ds, err := subSess.AcceptDataStream(t.Context())
if err != nil {
t.Fatalf("AcceptDataStream: %v", err)
}
first, ok := ds.(*session.IncomingSubgroupStream)
in, ok := ds.(*session.IncomingSubgroupStream)
if !ok {
t.Fatalf("AcceptDataStream = %T, want a subgroup stream", ds)
}
if _, err := first.ReadObject(); err != nil {
t.Fatalf("ReadObject: %v", err)
type result struct {
ids []uint64
end error
}
done := make(chan result, 1)
go func() {
var r result
for {
if _, err := first.ReadObject(); err != nil {
o, err := in.ReadDecoded()
if err != nil {
r.end = err
done <- r
return
}
r.ids = append(r.ids, o.ObjectID)
}
}()
if _, err := subSess.UpdateRequest(
Expand All @@ -280,11 +289,16 @@ func TestRelay_ForwardStateOmissionResetsReopenedStream(t *testing.T) {
}
close(resumed)

ids, end := readUntilEnd(t, subSess)
if !slices.Equal(ids, []uint64{2}) {
t.Fatalf("the reopened stream carried Objects %v, want [2]", ids)
var r result
select {
case r = <-done:
case <-time.After(2 * time.Second):
t.Fatal("the subgroup stream did not end")
}
if !slices.Equal(r.ids, []uint64{0, 2}) {
t.Fatalf("the stream carried Objects %v, want [0 2]", r.ids)
}
requireReset(t, end, "Object 1 was omitted while paused")
requireReset(t, r.end, "Object 1 was omitted while paused")
}

// TestRelay_StartRaisedToLaterGroupResetsPromptly: a Start raised past the
Expand Down
94 changes: 82 additions & 12 deletions pkg/relay/handler_fanout.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,19 @@ type fwdObject struct {
// first marks the subgroup's true first object (§11.4.2 FIRST_OBJECT);
// only an outbound stream beginning with it sets the bit.
first bool

// follows reports that the Object handed to the writer before this one
// was read just before it from the same inbound stream, with only
// Objects the subscriber's filters rejected between (see
// [subgroupWriter.admit]).
follows bool
}

// inboundPos is an Object's place on its inbound subgroup stream: the stream,
// and how many Objects had been read from it, this one included.
type inboundPos struct {
src *session.IncomingSubgroupStream
seq uint64
}

// subgroupWriterSet is the payload of a [registry.SharedSubgroup]: one
Expand Down Expand Up @@ -265,6 +278,9 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming

var (
firstObj = true
// pos counts every Object read, dedup losers included: each is an
// Object between its neighbours (§11.4.3).
pos = inboundPos{src: stream}
// terminalSeen: an EndOfGroup/EndOfTrack was read on this inbound
// stream; any later object makes the track malformed (§11.4.3,
// §2.4.2).
Expand Down Expand Up @@ -315,6 +331,7 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming

isTrueFirst := firstObj && !hdr.ReplayingSubgroup
firstObj = false
pos.seq++
objectID := stream.ObjectID() // resolved by ReadObject (§11.4.2)

// Tracked whether or not this copy wins the dedup claim below.
Expand Down Expand Up @@ -368,8 +385,14 @@ func (h *sessionHandler) runFanout(ctx context.Context, stream *session.Incoming
if w == nil {
continue
}
if w.admit(hdr, objectID, obj.Properties) {
w.publish(fwdObject{obj: obj, absID: objectID, first: isTrueFirst, maxCacheAge: liveMaxAge})
if take, follows := w.admit(pos, hdr, objectID, obj.Properties); take {
w.publish(fwdObject{
obj: obj,
absID: objectID,
first: isTrueFirst,
maxCacheAge: liveMaxAge,
follows: follows,
})
}
}
sg.Mu.Unlock()
Expand Down Expand Up @@ -435,7 +458,8 @@ func (h *sessionHandler) openWriterForSub(
// subgroupWriter is the per-subscriber writer goroutine: it drains an inbox
// onto outbound subgroup streams on the subscriber's session.
//
// - A gap in Object IDs resets the stream and opens a fresh one (§11.4.3).
// - An Object that is not the next Object resets the stream and opens a
// fresh one (§11.4.3, see [isNextObject]).
// - A clean inbound EOF FINs the stream; an inbound error resets it.
// - A full inbox drops the object. An object that waited longer than
// maxLag resets with TOO_FAR_BEHIND and terminates the subscription
Expand Down Expand Up @@ -485,15 +509,27 @@ type subgroupWriter struct {
// straggler from another upstream. Only touched by admit, under sg.Mu.
lastAdmitted uint64
hasAdmitted bool
// lastPos is where the last Object admit let through was read, moved
// past the filtered Objects read after it; zero once the relay drops
// one. Only touched by admit and publish, under sg.Mu.
lastPos inboundPos
}

// admit decides whether w takes the Object at objectID of the subgroup hdr
// names, closing w when it will take none again.
func (w *subgroupWriter) admit(hdr message.SubgroupHeader, objectID uint64, props []byte) bool {
// names, read at pos, closing w when it will take none again. follows is
// [fwdObject.follows] for a taken Object.
func (w *subgroupWriter) admit(
pos inboundPos,
hdr message.SubgroupHeader,
objectID uint64,
props []byte,
) (take, follows bool) {
follows = w.lastPos.src == pos.src && w.lastPos.seq+1 == pos.seq
switch w.sub.ForwardDecision(hdr.GroupID, objectID, hdr.SubgroupID, hdr.PublisherPriority, props) {
case registry.Forward:
w.lastAdmitted, w.hasAdmitted = objectID, true
return true
w.lastPos = pos
return true, follows
case registry.SkipObject, registry.SkipPaused:
// The stream stays open for later Objects, but the Subgroup is now
// incomplete (§11.4.3).
Expand All @@ -509,7 +545,13 @@ func (w *subgroupWriter) admit(hdr message.SubgroupHeader, objectID uint64, prop
w.markIncomplete(moqt.StreamResetCancelled)
}
}
return false
// §11.4.3: an Object that "did not pass the subscriber's filters" does
// not separate the ones either side of it. Forward State is one of them
// (§5.1.5).
if follows {
w.lastPos = pos
}
return false, false
}

// markIncomplete sets subgroupWriter.incomplete; the first code recorded
Expand Down Expand Up @@ -541,7 +583,8 @@ func (w *subgroupWriter) resetCode() moqt.StreamResetCode {

// publish enqueues fwd without blocking, stamping its enqueue time for the
// lag check. On overflow the object is dropped, and past maxDropsBeforeReset
// the writer is closed in reset mode. It is a no-op after close.
// the writer is closed in reset mode. It is a no-op after close. Callers hold
// sg.Mu.
func (w *subgroupWriter) publish(fwd fwdObject) {
w.dropsMu.Lock()
if w.closed {
Expand All @@ -556,6 +599,7 @@ func (w *subgroupWriter) publish(fwd fwdObject) {
w.metrics.ObjectForwarded(w.ref, w.hdr.SubgroupID)
default:
w.metrics.ObjectDropped(w.ref, w.hdr.SubgroupID)
w.lastPos = inboundPos{} // the next Object does not follow a sent one
w.dropsMu.Lock()
w.drops++
w.markIncompleteLocked(moqt.StreamResetExcessiveLoad)
Expand Down Expand Up @@ -634,6 +678,9 @@ func (w *subgroupWriter) run() {
prevID uint64
hasWritten bool
writeFailed bool
// dropped: an Object was dropped since the last one written, so the
// next is not known to follow it.
dropped bool
)

// reopen resets the current outbound stream (if any) and opens a fresh
Expand Down Expand Up @@ -712,10 +759,8 @@ func (w *subgroupWriter) run() {
}
}

// §11.4.3: only "the next Object" may go on an existing stream. Of
// the draft's three ways to tell, this relay uses only "one greater
// than the previous Object" (a choice) and reopens on any other gap.
if hasWritten && fwd.absID != prevID+1 {
// §11.4.3: only "the next Object" may go on an existing stream.
if hasWritten && !isNextObject(fwd, prevID, dropped) {
w.metrics.SubgroupStreamReset(w.ref, w.hdr.SubgroupID, ResetCauseGap)
if !reopen(fwd.first) {
failWrites()
Expand All @@ -727,6 +772,7 @@ func (w *subgroupWriter) run() {
// the open above, which can block.
if expired(fwd) {
w.dropExpired(hasWritten)
dropped = true
continue
}

Expand Down Expand Up @@ -761,6 +807,7 @@ func (w *subgroupWriter) run() {
}
prevID = fwd.absID
hasWritten = true
dropped = false
// §11.4.3: a later reset still delivers what was written.
w.out.MarkReliable()
}
Expand Down Expand Up @@ -823,6 +870,29 @@ func (w *subgroupWriter) run() {
w.closeOut(true, 0)
}

// isNextObject reports whether fwd is "the next Object" (§11.4.3) on a stream
// whose last Object is prevID. Of the draft's ways to tell, the relay uses:
// the Object ID is one greater; fwd follows the last Object on its inbound
// 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.
func isNextObject(fwd fwdObject, prevID uint64, dropped bool) bool {
if fwd.absID == prevID+1 {
return true
}
if dropped || fwd.absID <= prevID {
return false
}
if fwd.follows {
return true
}
gap, ok := message.PriorObjectIDGap(fwd.obj.Properties)
return ok && fwd.absID-gap == prevID+1
}

// openCounted opens a subgroup stream, counting it for the §10.12 Stream
// Count; it fails once the subscription has terminated.
func (w *subgroupWriter) openCounted(hdr message.SubgroupHeader) (*session.OutgoingSubgroupStream, error) {
Expand Down
34 changes: 0 additions & 34 deletions pkg/relay/handler_fanout_firstobject_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -165,40 +165,6 @@ func TestFanout_FirstObjectBitClearedForFilteredHead(t *testing.T) {
}
}

// TestFanout_FirstObjectBitAcrossGapReopen: after a gap reopen (§11.4.3) only
// the original stream carries FIRST_OBJECT.
func TestFanout_FirstObjectBitAcrossGapReopen(t *testing.T) {
t.Parallel()
pub, sub := firstObjectTopology(t, nil)
// Objects 0, 1, then a jump to 5: the relay must not forward the
// non-consecutive object on the same outbound stream.
writeSubgroupObjects(t, pub, message.SubgroupHeader{
SubgroupIDMode: message.SubgroupIDImplicitZero, TrackAlias: 7, GroupID: 0,
}, []uint64{0, 1, 5})

caps := captureSubgroups(t, sub, 2)
var origin, reopened *subgroupCapture
for i := range caps {
if len(caps[i].Objects) > 0 && caps[i].Objects[0] == 0 {
origin = &caps[i]
} else {
reopened = &caps[i]
}
}
if origin == nil || reopened == nil {
t.Fatalf("expected an origin and a reopened stream, got %+v", caps)
}
if origin.Header.ReplayingSubgroup {
t.Error("origin stream begins with the subgroup's first object: FIRST_OBJECT must be set")
}
if !reopened.Header.ReplayingSubgroup {
t.Error("gap-reopened stream is mid-subgroup: FIRST_OBJECT must be clear")
}
if len(reopened.Objects) != 1 || reopened.Objects[0] != 5 {
t.Errorf("reopened stream objects = %v, want [5]", reopened.Objects)
}
}

// TestFanout_FirstObjectBitNotInvented: an inbound replay stream (FIRST_OBJECT
// clear) is not forwarded with the bit set.
func TestFanout_FirstObjectBitNotInvented(t *testing.T) {
Expand Down
Loading
Loading