Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
0bc25d5
fix(relay): keep INCLUDE_PROPERTIES=0 across a REQUEST_UPDATE (§10.9,…
floatdrop Sep 27, 2026
6480ab3
fix(relay): answer TRACK_STATUS for a track SUBSCRIBE would accept (§…
floatdrop Sep 27, 2026
9ac663f
fix(relay): close with GOAWAY_TIMEOUT only after a GOAWAY ran out (§3.5)
floatdrop Sep 27, 2026
6e99b0b
fix(relay): carry the first Object's delivery-timeout override onto r…
floatdrop Sep 27, 2026
9518daa
fix(relay): reset a merged Subgroup a clean survivor does not cover (…
floatdrop Sep 27, 2026
b1128f1
fix(relay): set FETCH_OK End Of Track at the END_OF_TRACK Object (§10…
floatdrop Sep 27, 2026
ada7018
fix(relay): reset a cancelled FETCH's request and data streams (§5.2)
floatdrop Sep 27, 2026
4de6761
fix(relay): bound Objects stitched from an upstream FETCH by its MAX_…
floatdrop Sep 27, 2026
2b72d55
fix(relay): cancel FETCHes and reset fetch streams of a malformed tra…
floatdrop Sep 27, 2026
18b4708
fix(relay): verify a REQUEST_UPDATE's tokens on SUBSCRIBE, FETCH and …
floatdrop Sep 27, 2026
5350f87
fix(session,relay): refuse malformed Track Properties with INTERNAL_E…
floatdrop Sep 27, 2026
5e853ad
fix(relay): keep a PUBLISH_SKIPPED across a prefix update that moves …
floatdrop Sep 27, 2026
82bb1b9
fix(relay): do not forward a SUBSCRIBE_TRACKS holder the track its ow…
floatdrop Sep 27, 2026
9d31287
bench: refresh the baseline for the draft-20 compliance fixes
floatdrop Sep 27, 2026
e3a70d4
fix(session,relay): initiate no requests on a session with a GOAWAY i…
floatdrop Sep 27, 2026
8e8085b
fix(relay): end a fill at the Largest Object the subscriber was told …
floatdrop Sep 27, 2026
023d79e
fix(relay): check a forward's registered downstream under its claim, …
floatdrop Sep 27, 2026
23c67e9
fix(relay): judge a merged Subgroup's FIN from the ledger's lowest Ob…
floatdrop Sep 27, 2026
90cba9a
test(relay): register the session before Stop in the close-code test
floatdrop Sep 27, 2026
a5b043d
style(relay): wrap the close-code test's TRACK_STATUS call (golines)
floatdrop Sep 27, 2026
58e40d4
fix(relay): let a client SUBSCRIBE and FETCH a track it publishes (§5.1)
floatdrop Sep 27, 2026
2288a1f
fix(relay): reset a refused FETCH's request stream with its data stre…
floatdrop Sep 27, 2026
a97626e
feat(relay): hold a SUBSCRIBE for a publisher under RENDEZVOUS_TIMEOU…
floatdrop Sep 27, 2026
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
94 changes: 46 additions & 48 deletions STATUS.md

Large diffs are not rendered by default.

400 changes: 205 additions & 195 deletions benchmarks/baseline-go1.27.txt

Large diffs are not rendered by default.

9 changes: 9 additions & 0 deletions pkg/moqt/session/goaway.go
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,15 @@ func (s *Session) SendGoaway(timeout time.Duration, newURI string) error {
return s.sendControl(msg)
}

// GoawaySent reports whether [Session.SendGoaway] was called: the session is
// draining, and §10.4 says its sender "SHOULD avoid initiating requests unless
// required by migration".
func (s *Session) GoawaySent() bool {
s.mu.Lock()
defer s.mu.Unlock()
return s.goawaySent
}

// handleGoaway records a received GOAWAY and notifies any waiter on
// GoawayReceived. §10.4: a second GOAWAY on the same control stream MUST
// terminate the session with PROTOCOL_VIOLATION.
Expand Down
2 changes: 1 addition & 1 deletion pkg/moqt/session/request.go
Original file line number Diff line number Diff line change
Expand Up @@ -944,7 +944,7 @@ func (r *Request) AcceptSubscribe(ok *message.SubscribeOK) (*Publication, error)
// Track Properties that fail validation (see
// [WithKnownMandatoryTrackProperties]) are rejected with REQUEST_ERROR —
// UNSUPPORTED_EXTENSION for an unknown Mandatory Track Property (§2.5.1),
// MALFORMED_TRACK for ones that do not parse — and the error returned. A
// INTERNAL_ERROR for ones that do not parse — and the error returned. A
// session-fatal value (§12.5, §12.6) closed the session in AcceptRequest. An
// alias collision closes the session with DUPLICATE_TRACK_ALIAS and returns
// *ErrDuplicateTrackAlias (§11.1).
Expand Down
11 changes: 6 additions & 5 deletions pkg/moqt/session/track_properties.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,9 +35,8 @@ func (e *ErrUnsupportedMandatoryTrackProperty) Error() string {

// ErrMalformedTrackProperties is wrapped by the error [ValidateTrackProperties]
// returns when raw Track Properties do not parse: a Key-Value-Pair that
// "cannot be parsed" makes the track malformed (§12.7, §2.4.2). Answering
// with MALFORMED_TRACK, which §10.6 defines only for FETCH, is this package's
// choice.
// "cannot be parsed" makes the track malformed (§12.7, §2.4.2). A request is
// refused with INTERNAL_ERROR (see [TrackPropertiesRejectCode]).
var ErrMalformedTrackProperties = errors.New("moqt/session: malformed track properties")

// ErrTrackPropertiesNotAllowed is returned, and nothing sent, when asked to
Expand Down Expand Up @@ -124,10 +123,12 @@ func (s *Session) checkTrackPropertyValues(raw []byte, context string) error {
}

// TrackPropertiesRejectCode is the REQUEST_ERROR code for a Track Properties
// validation error: UNSUPPORTED_EXTENSION (§2.5.1) or MALFORMED_TRACK.
// validation error on a PUBLISH or SUBSCRIBE: UNSUPPORTED_EXTENSION for an
// unknown Mandatory Track Property (§2.5.1), else INTERNAL_ERROR.
// MALFORMED_TRACK is defined only "In response to a FETCH" (§10.6.2).
func TrackPropertiesRejectCode(err error) moqt.RequestErrorCode {
if _, ok := errors.AsType[*ErrUnsupportedMandatoryTrackProperty](err); ok {
return moqt.RequestUnsupportedExtension
}
return moqt.RequestMalformedTrack
return moqt.RequestInternalError
}
5 changes: 3 additions & 2 deletions pkg/moqt/session/track_status_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -190,7 +190,8 @@ func TestTrackStatusFollowupClosesSession(t *testing.T) {

// TestAcceptPublishTrackPropertiesRejected: a PUBLISH with an unknown
// Mandatory Track Property is refused with UNSUPPORTED_EXTENSION (§2.5.1), and
// unparseable Track Properties with MALFORMED_TRACK.
// unparseable Track Properties with INTERNAL_ERROR (§10.6.2 defines
// MALFORMED_TRACK only for FETCH).
func TestAcceptPublishTrackPropertiesRejected(t *testing.T) {
for _, tc := range []struct {
name string
Expand All @@ -200,7 +201,7 @@ func TestAcceptPublishTrackPropertiesRejected(t *testing.T) {
{"unknown mandatory", message.AppendTrackProperties([]wire.KVPair{
{Type: message.MandatoryTrackPropertyMin, IntVal: 1},
}), moqt.RequestUnsupportedExtension},
{"malformed", []byte{0x01}, moqt.RequestMalformedTrack},
{"malformed", []byte{0x01}, moqt.RequestInternalError},
} {
t.Run(tc.name, func(t *testing.T) {
client, server := openPair(t,
Expand Down
21 changes: 17 additions & 4 deletions pkg/relay/cache/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,11 @@ type CachedObject struct {
MaxCacheDuration time.Duration
HasMaxCacheDuration bool

// Stitched marks an Object read from an upstream FETCH for one response
// rather than stored: ReceivedAt is when it was read, and only a positive
// MaxCacheDuration bounds it, not the relay's TTL (see [ObjectCache.Expired]).
Stitched bool

// EndOfUnknownRange marks this element as a §11.4.4.2 End of Unknown
// Range (0x10C) FETCH marker rather than a stored object: every Location
// from the previous element in the response stream (exclusive) through
Expand Down Expand Up @@ -258,12 +263,20 @@ func (c *ObjectCache) notExpiredLocked(obj *CachedObject) bool {
return c.maxAge <= 0 || age <= c.maxAge
}

// Expired reports whether obj, taken from this cache, may no longer be
// served (§12.3: "MUST NOT start forwarding"). Elements the cache did not
// store (range markers, Objects stitched from upstream) never expire.
// Expired reports whether obj, in a FETCH or fill response on this cache's
// track, may no longer be served (§12.3: "MUST NOT start forwarding"): one
// taken from this cache within its own MAX_CACHE_DURATION and the relay's TTL
// (see notExpiredLocked). Range markers never expire. An
// Object stitched from an upstream FETCH expires only past its own
// MAX_CACHE_DURATION ("any individual Object received through this
// subscription or fetch"); a present 0 sets no limit on it, since it is passed
// through rather than served from the cache (interpretation).
func (c *ObjectCache) Expired(obj *CachedObject) bool {
if obj.ReceivedAt.IsZero() {
switch {
case obj.ReceivedAt.IsZero():
return false
case obj.Stitched:
return obj.MaxCacheDuration > 0 && time.Since(obj.ReceivedAt) > obj.MaxCacheDuration
}
c.mu.RLock()
defer c.mu.RUnlock()
Expand Down
81 changes: 80 additions & 1 deletion pkg/relay/cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,8 @@ import (

// The relay's per-track cache as seen by FETCH: size-based eviction, and
// MAX_CACHE_DURATION (§12.3), after which the relay must not start forwarding
// an Object, from the cache or from a live subscriber's queue.
// an Object, from the cache, from a live subscriber's queue, or from an upstream
// FETCH it passes through.

// TestFetch_CacheEvictionUnderLoad: past MaxCacheSize the cache evicts the
// oldest Objects, so a FETCH of the early range returns fewer Objects than it
Expand Down Expand Up @@ -334,3 +335,81 @@ func TestRelay_MaxCacheDurationExpiresDuringFetch(t *testing.T) {
)
}
}

// TestRelay_MaxCacheDurationBoundsStitchedObject: an Object received through
// an upstream FETCH is bound by that FETCH's MAX_CACHE_DURATION (§12.3: "any
// individual Object received through this subscription or fetch"). One the
// upstream sent, then held its stream open past the duration, is not
// forwarded but marked unknown. A present 0, which the cache reads as "never
// serve from the cache", sets no limit on an Object passed through.
func TestRelay_MaxCacheDurationBoundsStitchedObject(t *testing.T) {
t.Parallel()
const hold = 150 * time.Millisecond
for _, tc := range []struct {
name string
millis uint64
want fetchElem
}{
{"expired", 50, unknownAt(0, 1)},
{"zero", 0, obj(0, 1)},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
upSess, teardown := connectRelay(t, relay.Config{})
t.Cleanup(teardown)
if _, err := upSess.PublishNamespace(
t.Context(),
&message.PublishNamespace{Namespace: ns("video")},
); err != nil {
t.Fatalf("PublishNamespace: %v", err)
}
go func() {
for {
req, err := upSess.AcceptRequest(t.Context())
if err != nil {
return
}
switch m := req.First.(type) {
case *message.Subscribe:
if req.Reply(&message.SubscribeOK{TrackAlias: 42}) != nil {
return
}
// The live stream misses Object 1.
publishCam1Group(t, upSess, 42, true, cam1Object{0, 0, nil}, cam1Object{0, 2, nil})
case *message.Fetch:
if req.Reply(&message.FetchOK{
EndLocation: fetchOKEnd(m),
TrackProperties: message.AppendTrackProperties(
trackProp(message.PropertyMaxCacheDuration, tc.millis)),
}) != nil {
return
}
out, err := upSess.OpenFetchStream(message.FetchHeader{RequestID: m.RequestID})
if err != nil {
return
}
_ = out.WriteObject(&message.FetchObject{
SerializationFlags: message.FetchFlagGroupIDDelta | message.FetchFlagObjectIDDelta |
message.FetchFlagPriority | uint64(message.FetchSubgroupIDExplicit),
GroupIDDelta: 0, ObjectIDDelta: 1, ObjectPayload: []byte("x"),
})
time.Sleep(hold)
_ = out.Close()
}
}
}()
live := dialAnotherClient(t, upSess)
subscribeCam1(t, live)
go drainAll(t.Context(), live)
fc := dialAnotherClient(t, upSess)
waitRelayLargest(t, fc, ns("video"), []byte("cam1"), 0, 2)

got := fetchCam1Range(t, fc, message.Location{}, message.Location{Group: 0, Object: 2},
message.GroupOrderAscending)
want := []fetchElem{obj(0, 0), tc.want, obj(0, 2)}
if !slices.Equal(got, want) {
t.Fatalf("FETCH elements %v, want %v", got, want)
}
})
}
}
77 changes: 77 additions & 0 deletions pkg/relay/drain_requests_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
package relay_test

import (
"context"
"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/relay"
)

// TestRelay_NoRequestsToPeerSentGoaway: once the relay has sent a publisher
// GOAWAY it avoids initiating requests to it (§10.4: "the sender SHOULD avoid
// initiating requests unless required by migration"), so a SUBSCRIBE that
// needs that publisher is refused with GOING_AWAY rather than sent upstream.
func TestRelay_NoRequestsToPeerSentGoaway(t *testing.T) {
t.Parallel()
l := newPipeListener()
r := relay.New(l, relay.Config{GoawayTimeout: 5 * time.Second})
go func() { _ = r.Start(t.Context()) }()
dial := func() *session.Session {
conn, err := l.Dial()
if err != nil {
t.Fatalf("Dial: %v", err)
}
s, err := session.Client(t.Context(), conn)
if err != nil {
t.Fatalf("session.Client: %v", err)
}
return s
}
pubSess, subSess := dial(), dial()
stopped := make(chan struct{})
defer func() {
_ = pubSess.Close(moqt.SessionNoError, "done")
_ = subSess.Close(moqt.SessionNoError, "done")
<-stopped
}()

video := ns("video")
if _, err := pubSess.PublishNamespace(t.Context(), &message.PublishNamespace{Namespace: video}); err != nil {
t.Fatalf("PublishNamespace: %v", err)
}
subscribes := make(chan *session.Request, 4)
go func() {
for {
req, err := pubSess.AcceptRequest(t.Context())
if err != nil {
return
}
subscribes <- req
_ = req.Reply(&message.SubscribeOK{TrackAlias: 1})
}
}()

go func() {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
_ = r.Stop(ctx)
close(stopped)
}()
select {
case <-pubSess.GoawayReceived():
case <-time.After(2 * time.Second):
t.Fatal("the publisher never got the relay's GOAWAY")
}

_, err := subSess.Subscribe(t.Context(), &message.Subscribe{Namespace: video, Name: []byte("cam1")})
select {
case req := <-subscribes:
t.Fatalf("the relay sent %s to a publisher it had sent GOAWAY", req.First.Type())
default:
}
requireRejectedWithCode(t, err, moqt.RequestGoingAway)
}
28 changes: 28 additions & 0 deletions pkg/relay/export_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,3 +20,31 @@ func SetTestHookEarlyStreamWaiting(hook func(alias uint64)) (restore func()) {
testHookEarlyStreamWaiting.Store(&hook)
return func() { testHookEarlyStreamWaiting.Store(prev) }
}

// SetTestHookBeforeDownstreamRegistered installs hook, to be called once a
// SUBSCRIBE has an upstream for its track and before its downstream is
// registered, and returns a function restoring the previous value. See
// [testHookBeforeDownstreamRegistered].
func SetTestHookBeforeDownstreamRegistered(hook func(track.FullTrackName)) (restore func()) {
prev := testHookBeforeDownstreamRegistered.Load()
testHookBeforeDownstreamRegistered.Store(&hook)
return func() { testHookBeforeDownstreamRegistered.Store(prev) }
}

// SetTestHookBeforeFill installs hook, to be called as a fill is about to be
// evaluated, and returns a function restoring the previous value. See
// [testHookBeforeFill].
func SetTestHookBeforeFill(hook func(track.FullTrackName)) (restore func()) {
prev := testHookBeforeFill.Load()
testHookBeforeFill.Store(&hook)
return func() { testHookBeforeFill.Store(prev) }
}

// SetTestHookBeforeForwardClaim installs hook, to be called as a forward of a
// track to a SUBSCRIBE_TRACKS holder is about to claim it, and returns a
// function restoring the previous value. See [testHookBeforeForwardClaim].
func SetTestHookBeforeForwardClaim(hook func(track.FullTrackName)) (restore func()) {
prev := testHookBeforeForwardClaim.Load()
testHookBeforeForwardClaim.Store(&hook)
return func() { testHookBeforeForwardClaim.Store(prev) }
}
84 changes: 84 additions & 0 deletions pkg/relay/fetch_holes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -394,3 +394,87 @@ func TestFetch_FillTimeoutBoundsUpstreamRead(t *testing.T) {
t.Fatalf("FETCH elements %v, want %v", got, want)
}
}

// TestFetch_CancelResetsStreams: a requester that cancels a FETCH while the
// relay still waits on the upstream for a hole has the data stream and the
// request stream reset with CANCELLED (§5.2: "It MUST reset the bidi request
// stream and unidirectional data stream associated with the FETCH"), not
// served once FILL_TIMEOUT runs out.
func TestFetch_CancelResetsStreams(t *testing.T) {
t.Parallel()
l := newPipeListener()
resets := make(chan streamReset, 16)
l.resetsFor = resetsOn(3, resets) // the upstream is 1, the live subscriber 2
upSess, teardown := connectRelayOn(t, relay.Config{}, l)
t.Cleanup(teardown)
if _, err := upSess.PublishNamespace(t.Context(), &message.PublishNamespace{Namespace: ns("video")}); err != nil {
t.Fatalf("PublishNamespace: %v", err)
}
go func() {
for {
req, err := upSess.AcceptRequest(t.Context())
if err != nil {
return
}
switch m := req.First.(type) {
case *message.Subscribe:
if req.Reply(&message.SubscribeOK{TrackAlias: 42}) != nil {
return
}
// The live stream misses Object 1.
publishCam1Group(t, upSess, 42, true, cam1Object{0, 0, nil}, cam1Object{0, 2, nil})
case *message.Fetch:
if req.Reply(&message.FetchOK{EndLocation: fetchOKEnd(m)}) != nil {
return
}
out, err := upSess.OpenFetchStream(message.FetchHeader{RequestID: m.RequestID})
if err != nil {
return
}
// Nothing more: the stream stays open.
t.Cleanup(func() { out.Cancel(moqt.StreamResetCancelled) })
}
}
}()
live := dialAnotherClient(t, upSess)
subscribeCam1(t, live)
go drainAll(t.Context(), live)
fc := dialAnotherClient(t, upSess)
waitRelayLargest(t, fc, ns("video"), []byte("cam1"), 0, 2)

fr, err := fc.Fetch(t.Context(), &message.Fetch{
Namespace: ns("video"), Name: []byte("cam1"),
Parameters: message.Parameters{
fetchRangeFilter(message.Location{}, message.Location{Group: 0, Object: 2}),
message.FillTimeoutParam(5 * time.Second),
},
})
if err != nil {
t.Fatalf("Fetch: %v", err)
}
if _, err := fc.AcceptDataStream(t.Context()); err != nil {
t.Fatalf("AcceptDataStream: %v", err)
}
_ = fr.Close()

var uni, bidi bool
deadline := time.After(500 * time.Millisecond)
for !uni || !bidi {
select {
case r := <-resets:
if r.code != moqt.StreamResetCancelled {
t.Fatalf("the relay reset a stream (%v) with %v, want CANCELLED", r.stream, r.code)
}
switch r.stream {
case fetchStreamReset:
uni = true
case requestStreamReset:
bidi = true
case fetchStreamStop:
}
case <-deadline:
t.Fatalf("within 500ms of the cancel: data stream reset %t, request stream reset %t; want both",
uni, bidi)
}
}
}
Loading
Loading