diff --git a/STATUS.md b/STATUS.md index 7e91b850..ff64fde5 100644 --- a/STATUS.md +++ b/STATUS.md @@ -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. | diff --git a/pkg/moqt/message/object_properties.go b/pkg/moqt/message/object_properties.go index 666eeffe..5ce173f9 100644 --- a/pkg/moqt/message/object_properties.go +++ b/pkg/moqt/message/object_properties.go @@ -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 +} diff --git a/pkg/moqt/message/properties_test.go b/pkg/moqt/message/properties_test.go index 16e5ddd7..5b3f8219 100644 --- a/pkg/moqt/message/properties_test.go +++ b/pkg/moqt/message/properties_test.go @@ -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)), diff --git a/pkg/relay/forward_state_test.go b/pkg/relay/forward_state_test.go index d5d8bd5d..e4485c10 100644 --- a/pkg/relay/forward_state_test.go +++ b/pkg/relay/forward_state_test.go @@ -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() @@ -240,7 +242,7 @@ 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() }() @@ -248,18 +250,25 @@ func TestRelay_ForwardStateOmissionResetsReopenedStream(t *testing.T) { 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( @@ -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 diff --git a/pkg/relay/handler_fanout.go b/pkg/relay/handler_fanout.go index c8bd73ab..65648356 100644 --- a/pkg/relay/handler_fanout.go +++ b/pkg/relay/handler_fanout.go @@ -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 @@ -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). @@ -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. @@ -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() @@ -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 @@ -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). @@ -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 @@ -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 { @@ -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) @@ -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 @@ -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() @@ -727,6 +772,7 @@ func (w *subgroupWriter) run() { // the open above, which can block. if expired(fwd) { w.dropExpired(hasWritten) + dropped = true continue } @@ -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() } @@ -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) { diff --git a/pkg/relay/handler_fanout_firstobject_test.go b/pkg/relay/handler_fanout_firstobject_test.go index 4836e753..778981dc 100644 --- a/pkg/relay/handler_fanout_firstobject_test.go +++ b/pkg/relay/handler_fanout_firstobject_test.go @@ -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) { diff --git a/pkg/relay/handler_fanout_next_test.go b/pkg/relay/handler_fanout_next_test.go new file mode 100644 index 00000000..ca83f998 --- /dev/null +++ b/pkg/relay/handler_fanout_next_test.go @@ -0,0 +1,153 @@ +package relay + +import ( + "context" + "errors" + "io" + "slices" + "testing" + "time" + + "github.com/floatdrop/moq-go/pkg/moqt/message" + "github.com/floatdrop/moq-go/pkg/moqt/session" + "github.com/floatdrop/moq-go/pkg/moqt/session/sessiontest" +) + +// forwarded is one Object as the subscriber received it, or with end set, the +// end of the stream it came on. +type forwarded struct { + stream int // 1-based index of the subgroup stream it came on + id uint64 + replay bool // the stream's header had FIRST_OBJECT clear (§11.4.2) + end error // io.EOF for a FIN +} + +// readForwarded emits every Object and stream end of every subgroup stream srv +// accepts. +func readForwarded(ctx context.Context, srv *session.Session) <-chan forwarded { + out := make(chan forwarded, 16) + go func() { + for stream := 1; ; stream++ { + ds, err := srv.AcceptDataStream(ctx) + if err != nil { + return + } + sg, ok := ds.(*session.IncomingSubgroupStream) + if !ok { + return + } + for { + if _, err := sg.ReadObject(); err != nil { + out <- forwarded{stream: stream, end: err} + break + } + out <- forwarded{stream: stream, id: sg.ObjectID(), replay: sg.Header.ReplayingSubgroup} + } + } + }() + return out +} + +// awaitForwarded waits for the next Object or stream end the subscriber +// receives. +func awaitForwarded(t *testing.T, in <-chan forwarded) forwarded { + t.Helper() + select { + case f := <-in: + return f + case <-time.After(2 * time.Second): + t.Fatal("nothing forwarded") + return forwarded{} + } +} + +// TestSubgroupWriter_RelayDropBreaksRun: an Object the relay dropped was not +// sent, so the upstream's next Object is not "the next Object" after the one +// before it (§11.4.3: "with no other Objects in between"). It goes on a new +// stream, without FIRST_OBJECT (§11.4.2), and both streams end with a reset, +// since an Object was omitted. +func TestSubgroupWriter_RelayDropBreaksRun(t *testing.T) { + t.Parallel() + src := new(session.IncomingSubgroupStream) // identity only + for _, tc := range []struct { + name string + // Objects 0, 1 and 2 are queued before the writer starts; 2 overflows + // a queue of 2, and expires with a maxAge. + queue int + maxAge time.Duration + want []forwarded // Objects then stream ends, in order + }{ + {"nothing dropped", 3, 0, []forwarded{ + {stream: 1, id: 0}, {stream: 1, id: 1}, {stream: 1, id: 2}, {stream: 1, id: 4}, + {stream: 1, end: io.EOF}, + }}, + {"queue overflow", 2, 0, []forwarded{ + {stream: 1, id: 0}, {stream: 1, id: 1}, {stream: 1, end: errReset}, + {stream: 2, id: 4, replay: true}, {stream: 2, end: errReset}, + }}, + {"expired", 3, time.Nanosecond, []forwarded{ + {stream: 1, id: 0}, {stream: 1, id: 1}, {stream: 1, end: errReset}, + {stream: 2, id: 4, replay: true}, {stream: 2, end: errReset}, + }}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + cli, srv := sessiontest.NewSessionPair(t) + got := readForwarded(t.Context(), srv) + w := newWedgeableWriter(t, cli) + w.inbox = make(chan fwdObject, tc.queue) + offer := func(seq, id uint64, maxAge time.Duration) { + take, follows := w.admit(inboundPos{src: src, seq: seq}, w.hdr, id, nil) + if take { + w.publish(fwdObject{ + obj: &message.SubgroupObject{Payload: []byte("x")}, + absID: id, + first: id == 0, + follows: follows, + maxCacheAge: maxAge, + }) + } + } + + // Upstream Objects 0, 1, 2 and 4, read in turn. + offer(1, 0, 0) + offer(2, 1, 0) + offer(3, 2, tc.maxAge) + go w.run() + var seen []forwarded + for { + f := awaitForwarded(t, got) + seen = append(seen, f) + if f.end == nil && f.id == 1 { + break + } + } + offer(4, 4, 0) + for { + f := awaitForwarded(t, got) + seen = append(seen, f) + if f.end == nil && f.id == 4 { + break + } + } + w.close(false, 0) + joinOrFatal(t, w) + seen = append(seen, awaitForwarded(t, got)) + + for i := range seen { + switch end := seen[i].end; { + case errors.Is(end, io.EOF): + seen[i].end = io.EOF + case end != nil: + seen[i].end = errReset + } + } + if !slices.Equal(seen, tc.want) { + t.Fatalf("forwarded %v,\nwant %v", seen, tc.want) + } + }) + } +} + +// errReset stands for any stream reset in [TestSubgroupWriter_RelayDropBreaksRun]. +var errReset = errors.New("reset") diff --git a/pkg/relay/handler_fanout_test.go b/pkg/relay/handler_fanout_test.go index d6f54c51..14a76cf0 100644 --- a/pkg/relay/handler_fanout_test.go +++ b/pkg/relay/handler_fanout_test.go @@ -608,142 +608,6 @@ func TestSubscribe_InstallsPriorityAndGroupOrder(t *testing.T) { // covered by the unit test on DownstreamSub setters. } -// TestFanout_GapInForwardedObjectIDsOpensNewStream: a gap in the Object IDs of -// one inbound subgroup resets the outbound stream and opens a new one for the -// next Object (§11.4.3). -func TestFanout_GapInForwardedObjectIDsOpensNewStream(t *testing.T) { - t.Parallel() - - pubSess, teardown := connectRelay(t, relay.Config{}) - defer teardown() - - const publisherAlias = uint64(7) - pubReqStream, err := pubSess.Publish(t.Context(), &message.Publish{ - Namespace: ns("video"), - Name: []byte("cam1"), - TrackAlias: publisherAlias, - }) - if err != nil { - t.Fatalf("Publish: %v", err) - } - defer pubReqStream.Close() - - subSess := dialAnotherClient(t, pubSess) - subReq, err := subSess.Subscribe(t.Context(), &message.Subscribe{ - Namespace: ns("video"), - Name: []byte("cam1"), - }) - if err != nil { - t.Fatalf("Subscribe: %v", err) - } - defer subReq.Close() - - // The subscriber-side reader collects every stream the relay opens - // for this subgroup, recording the (firstAbsID, objectCount, endErr) - // tuple for each. The relay should produce exactly two streams: one - // for absID=0 (then reset), and one for absID=2 (then FIN). - type streamSummary struct { - firstAbsID uint64 - count int - endErr error - } - streams := make(chan streamSummary, 4) - go func() { - for { - ds, err := subSess.AcceptDataStream(t.Context()) - if err != nil { - close(streams) - return - } - sg, ok := ds.(*session.IncomingSubgroupStream) - if !ok { - close(streams) - return - } - obj, err := sg.ReadObject() - if err != nil { - streams <- streamSummary{endErr: err} - continue - } - first := obj.ObjectIDDelta - count := 1 - for { - _, err := sg.ReadObject() - if err != nil { - streams <- streamSummary{firstAbsID: first, count: count, endErr: err} - break - } - count++ - } - } - }() - - pubSubgroup, err := pubSess.OpenSubgroup(message.SubgroupHeader{ - SubgroupIDMode: message.SubgroupIDExplicit, - TrackAlias: publisherAlias, - GroupID: 0, - SubgroupID: 0, - }) - if err != nil { - t.Fatalf("OpenSubgroup: %v", err) - } - - // Object with absolute ID 0. - if err := pubSubgroup.WriteObject(&message.SubgroupObject{ - ObjectIDDelta: 0, - Payload: []byte("first"), - }); err != nil { - t.Fatalf("WriteObject #0: %v", err) - } - // Object with absolute ID 2 — delta is (2 - 0 - 1) = 1. Skips ID 1. - if err := pubSubgroup.WriteObject(&message.SubgroupObject{ - ObjectIDDelta: 1, - Payload: []byte("third"), - }); err != nil { - t.Fatalf("WriteObject #2: %v", err) - } - if err := pubSubgroup.Close(); err != nil { - t.Fatalf("pubSubgroup.Close: %v", err) - } - - // Collect the first two summaries — that's all the relay should send. - var got []streamSummary - deadline := time.After(2 * time.Second) - for len(got) < 2 { - select { - case s, ok := <-streams: - if !ok { - t.Fatalf("subscriber stream channel closed after %d summaries: %v", len(got), got) - } - got = append(got, s) - case <-deadline: - t.Fatalf("subscriber did not receive two streams within deadline (got %d: %v)", len(got), got) - } - } - - // First stream: pre-gap object (absID=0), then reset. - if got[0].firstAbsID != 0 { - t.Errorf("stream 1 firstAbsID = %d, want 0", got[0].firstAbsID) - } - if got[0].count != 1 { - t.Errorf("stream 1 count = %d, want 1", got[0].count) - } - if errors.Is(got[0].endErr, io.EOF) { - t.Errorf("stream 1 ended with io.EOF, want a reset") - } - - // Second stream: post-gap object (absID=2), then clean FIN. - if got[1].firstAbsID != 2 { - t.Errorf("stream 2 firstAbsID = %d, want 2", got[1].firstAbsID) - } - if got[1].count != 1 { - t.Errorf("stream 2 count = %d, want 1", got[1].count) - } - if !errors.Is(got[1].endErr, io.EOF) { - t.Errorf("stream 2 ended with %v, want io.EOF (clean FIN)", got[1].endErr) - } -} - // TestFanout_InboundResetCancelsDownstream: a reset inbound subgroup stream // resets the downstream one rather than FINning it (§11.4.3). func TestFanout_InboundResetCancelsDownstream(t *testing.T) { diff --git a/pkg/relay/metrics.go b/pkg/relay/metrics.go index fa5394bb..d6f45ac3 100644 --- a/pkg/relay/metrics.go +++ b/pkg/relay/metrics.go @@ -71,12 +71,11 @@ type TrackRef struct { type ResetCause uint8 const ( - // ResetCauseGap is a §11.4.3 reopen: the next object to forward was not - // consecutive with the last one written, so the current outbound stream - // was reset and a fresh one opened. The relay MUST NOT forward a - // non-consecutive object on an existing subgroup stream, so this is - // correct behaviour — but it is also the direct consequence of an - // earlier drop or filter narrowing, and a subscriber sees the hole. + // ResetCauseGap is a §11.4.3 reopen: the object to forward was not known + // to be "the next Object" after the last one written, so the current + // outbound stream was reset and a fresh one opened. That is correct + // behaviour, but it usually follows an earlier drop, or Objects arriving + // from several upstreams, and a subscriber sees the hole. ResetCauseGap ResetCause = iota // ResetCauseDeliveryTimeout is §8: an object sat unsent past the diff --git a/pkg/relay/next_object_test.go b/pkg/relay/next_object_test.go new file mode 100644 index 00000000..a837f57e --- /dev/null +++ b/pkg/relay/next_object_test.go @@ -0,0 +1,255 @@ +package relay_test + +import ( + "errors" + "io" + "slices" + "testing" + "time" + + "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" +) + +// §11.4.3: the relay forwards an Object on an existing subgroup stream only if +// it is "the next Object", and otherwise resets the stream and opens another. + +// collectStreams reads events until objects Objects have arrived and every +// stream they came on has ended. +func collectStreams(t *testing.T, events <-chan objEvent, objects int) []objEvent { + t.Helper() + var ( + seen []objEvent + got int + streams int + ended int + ) + deadline := time.After(3 * time.Second) + for got < objects || ended < streams { + select { + case ev := <-events: + if ev.stream == 0 { + t.Fatalf("AcceptDataStream: %v", ev.err) + } + seen = append(seen, ev) + streams = max(streams, ev.stream) + if ev.err != nil { + ended++ + } else { + got++ + } + case <-deadline: + t.Fatalf("got %d/%d Objects, %d/%d streams ended: %v", got, objects, ended, streams, layout(seen)) + } + } + return seen +} + +// finned reports whether stream ended with a FIN in events. +func finned(events []objEvent, stream int) bool { + for _, ev := range events { + if ev.stream == stream && ev.err != nil { + return errors.Is(ev.err, io.EOF) + } + } + return false +} + +// TestFanout_NextObject_GapOnOneUpstreamStaysOnStream: Objects read in turn +// from one upstream stream are each the next Object (§11.4.3, second rule), +// so a publisher skipping IDs costs no reset. +func TestFanout_NextObject_GapOnOneUpstreamStaysOnStream(t *testing.T) { + t.Parallel() + pubSess, teardown := connectRelay(t, relay.Config{}) + defer teardown() + pub := publishVideoTrack(t, pubSess, "cam1", 1) + subSess := dialAnotherClient(t, pubSess) + subscribeCam1(t, subSess) + events := make(chan objEvent, 16) + go readSubgroups(t.Context(), subSess, events) + + sg, err := pub.OpenSubgroup(message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit}) + if err != nil { + t.Fatalf("OpenSubgroup: %v", err) + } + for _, id := range []uint64{0, 2, 5} { + if err := sg.WriteObjectAt(id, &message.SubgroupObject{Payload: []byte("x")}); err != nil { + t.Fatalf("WriteObjectAt %d: %v", id, err) + } + } + if err := sg.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + + seen := collectStreams(t, events, 3) + if got, want := layout(seen), [][]uint64{{0, 2, 5}}; !slices.EqualFunc(got, want, slices.Equal) { + t.Fatalf("streams = %v, want %v", got, want) + } + if !finned(seen, 1) { + t.Error("stream ended with a reset, want a FIN: no Object was omitted") + } +} + +// TestFanout_NextObject_FilteredObjectsBetween: Objects between two forwarded +// ones that did not pass the subscriber's filters keep the second the next +// Object (§11.4.3). +func TestFanout_NextObject_FilteredObjectsBetween(t *testing.T) { + t.Parallel() + pubSess, teardown := connectRelay(t, relay.Config{}) + defer teardown() + pub := publishVideoTrack(t, pubSess, "cam1", 1) + subSess := dialAnotherClient(t, pubSess) + subscribeCam1(t, subSess, message.RangeFilterParam(&message.RangeFilter{ + Type: message.ParamObjectIDFilter, + Ranges: []message.Range{{Start: 0, End: 0}, {Start: 2, End: 3}}, + })) + events := make(chan objEvent, 16) + go readSubgroups(t.Context(), subSess, events) + + sg, err := pub.OpenSubgroup(message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit}) + if err != nil { + t.Fatalf("OpenSubgroup: %v", err) + } + for range 4 { + if err := sg.WriteObject(&message.SubgroupObject{Payload: []byte("x")}); err != nil { + t.Fatalf("WriteObject: %v", err) + } + } + if err := sg.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + + seen := collectStreams(t, events, 3) + if got, want := layout(seen), [][]uint64{{0, 2, 3}}; !slices.EqualFunc(got, want, slices.Equal) { + t.Fatalf("streams = %v, want %v", got, want) + } + if finned(seen, 1) { + t.Error("stream ended with a FIN, want a reset: Object 1 was omitted") + } +} + +// twoPublishers PUBLISHes video/cam1 from two sessions, subscribes a third, and +// opens Subgroup 0 of Group 0, with Object Properties, from each publisher. +func twoPublishers(t *testing.T) (a, b *session.OutgoingSubgroupStream, events <-chan objEvent) { + t.Helper() + pubA, teardown := connectRelay(t, relay.Config{}) + t.Cleanup(teardown) + pubB := dialAnotherClient(t, pubA) + subSess := dialAnotherClient(t, pubA) + aPub := publishVideoTrack(t, pubA, "cam1", 1) + bPub := publishVideoTrack(t, pubB, "cam1", 2) + subscribeCam1(t, subSess) + ch := make(chan objEvent, 16) + go readSubgroups(t.Context(), subSess, ch) + + // Properties on both: the merged stream takes the first one's header. + hdr := message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit, Properties: true} + a, err := aPub.OpenSubgroup(hdr) + if err != nil { + t.Fatalf("A OpenSubgroup: %v", err) + } + b, err = bPub.OpenSubgroup(hdr) + if err != nil { + t.Fatalf("B OpenSubgroup: %v", err) + } + return a, b, ch +} + +// writeAndAwait writes the Object at id on sg and waits for the subscriber to +// receive it, so the next write is ordered after it at the relay. +func writeAndAwait( + t *testing.T, + sg *session.OutgoingSubgroupStream, + id uint64, + props []byte, + events <-chan objEvent, + streams *[]objEvent, +) { + t.Helper() + if err := sg.WriteObjectAt(id, &message.SubgroupObject{Properties: props, Payload: []byte("x")}); err != nil { + t.Fatalf("WriteObjectAt %d: %v", id, err) + } + for { + select { + case ev := <-events: + *streams = append(*streams, ev) + if ev.err == nil && ev.absID == id { + return + } + case <-time.After(2 * time.Second): + t.Fatalf("Object %d not forwarded", id) + } + } +} + +// layout is the Object IDs of each stream in events. +func layout(events []objEvent) [][]uint64 { + var out [][]uint64 + for _, ev := range events { + if ev.err != nil { + continue + } + for len(out) < ev.stream { + out = append(out, nil) + } + out[ev.stream-1] = append(out[ev.stream-1], ev.absID) + } + return out +} + +func priorObjectIDGap(gap uint64) []byte { + return message.AppendTrackProperties([]wire.KVPair{{Type: message.PropertyPriorObjectIDGap, IntVal: gap}}) +} + +// TestFanout_NextObject_AcrossUpstreams: with two upstreams feeding one +// Subgroup (§9.3), an Object from the other upstream after a gap is the next +// Object only if its Prior Object ID Gap (§12.9) covers the gap (§11.4.3, +// third rule). +func TestFanout_NextObject_AcrossUpstreams(t *testing.T) { + t.Parallel() + for _, tc := range []struct { + name string + props []byte + want [][]uint64 + }{ + {"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. + {"gap covering the previous Object", priorObjectIDGap(3), [][]uint64{{0}, {3}}}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + a, b, events := twoPublishers(t) + var seen []objEvent + writeAndAwait(t, a, 0, nil, events, &seen) + writeAndAwait(t, b, 3, tc.props, events, &seen) + if got := layout(seen); !slices.EqualFunc(got, tc.want, slices.Equal) { + t.Fatalf("streams = %v, want %v", got, tc.want) + } + }) + } +} + +// TestFanout_NextObject_DuplicateBetweenBreaksRun: an upstream's Object that +// lost the §2.1 dedup claim was not sent on the downstream stream, so the +// upstream's following Object is not the next Object after its earlier one +// (§11.4.3: "with no other Objects in between"). +func TestFanout_NextObject_DuplicateBetweenBreaksRun(t *testing.T) { + t.Parallel() + a, b, events := twoPublishers(t) + var seen []objEvent + writeAndAwait(t, a, 1, nil, events, &seen) + writeAndAwait(t, b, 0, nil, events, &seen) + // B's 1 duplicates A's and is dropped. + if err := b.WriteObjectAt(1, &message.SubgroupObject{Payload: []byte("x")}); err != nil { + t.Fatalf("WriteObjectAt 1: %v", err) + } + writeAndAwait(t, b, 3, nil, events, &seen) + if got, want := layout(seen), [][]uint64{{1}, {0}, {3}}; !slices.EqualFunc(got, want, slices.Equal) { + t.Fatalf("streams = %v, want %v", got, want) + } +} diff --git a/pkg/relay/publish_done_test.go b/pkg/relay/publish_done_test.go index 91497c2b..257404e9 100644 --- a/pkg/relay/publish_done_test.go +++ b/pkg/relay/publish_done_test.go @@ -264,20 +264,28 @@ func TestPublishDone_AfterGapReopen(t *testing.T) { pubSess, teardown := connectRelay(t, relay.Config{}) defer teardown() pub := publishVideoTrack(t, pubSess, "cam1", 1) + other := publishVideoTrack(t, dialAnotherClient(t, pubSess), "cam1", 2) subSess := dialAnotherClient(t, pubSess) subReq := subscribeCam1(t, subSess) - sg, err := pub.OpenSubgroup(message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit}) + hdr := message.SubgroupHeader{SubgroupIDMode: message.SubgroupIDExplicit} + sg, err := pub.OpenSubgroup(hdr) if err != nil { t.Fatalf("OpenSubgroup: %v", err) } + otherSg, err := other.OpenSubgroup(hdr) + if err != nil { + t.Fatalf("other OpenSubgroup: %v", err) + } wrote := make(chan struct{}) go func() { defer close(wrote) _ = sg.WriteObject(&message.SubgroupObject{Payload: []byte("0")}) - // Object 2 after object 0: the relay may not carry it on the same - // stream (§11.4.3), so it resets that one and opens another. - _ = sg.WriteObject(&message.SubgroupObject{ObjectIDDelta: 1, Payload: []byte("2")}) + // Objects 0 and 2 from two upstreams: in either order, neither is the + // next Object after the other (§11.4.3), so the relay resets the + // first stream and opens another. + _ = otherSg.WriteObjectAt(2, &message.SubgroupObject{Payload: []byte("2")}) + _ = otherSg.Close() }() ended := make(chan struct{}, 2) for range 2 { @@ -299,6 +307,10 @@ func TestPublishDone_AfterGapReopen(t *testing.T) { }() } + // PUBLISH_DONE goes out once both upstreams are done. + if err := other.Done(moqt.PublishDoneTrackEnded, "done"); err != nil { + t.Fatalf("other Done: %v", err) + } if err := pub.Done(moqt.PublishDoneTrackEnded, "done"); err != nil { t.Fatalf("Done: %v", err) }