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
14 changes: 6 additions & 8 deletions STATUS.md
Original file line number Diff line number Diff line change
Expand Up @@ -199,7 +199,7 @@ By package, bottom-up along the dependency stack:
| 10.12 | PUBLISH_DONE | 0x0B | DONE | Sent once every stream of the subscription has closed and no datagram send is in progress, with the exact Stream Count; written on its own goroutine, so subscribers do not wait on each other. When a track's last upstream ends, its PUBLISH_DONE code reaches subscribers if it is about the track (TRACK_ENDED, MALFORMED_TRACK); codes about the relay's own upstream subscription become INTERNAL_ERROR. Session `Publication.Done` resets the subgroups still open with CANCELLED, refuses later opens and writes (`ErrPublicationEnded`), and counts every subgroup opened, however the opens race it. |
| 10.13 | FETCH | 0x16 | DONE | Standalone, the only kind in draft-20. From the cache, a Location is non-existent only on a signal: a Prior Group or Object ID Gap, a Group's or the Track's end, or an upstream's FETCH. Other uncached Locations are FETCHed from a fetch-capable upstream in one span, within FILL_TIMEOUT, or else marked End of Unknown (or Timed-Out) Range. |
| 10.14 | FETCH_OK | 0x18 | DONE | The relay sets End Of Track when the End Location is the Object an END_OF_TRACK status made the Track's final one; it does not learn a Track's end from an upstream FETCH_OK's End Of Track. An End Location before the FETCH's Start closes the session. A Start relative to the Largest Object is compared through End ≤ Largest; an End of {0,0} is let through, as it cannot be told apart from "no content yet". |
| 10.15 | TRACK_STATUS | 0x0D | DONE | Reply via REQUEST_OK, then FIN; any follow-up from the requester closes the session. The relay answers from an Established subscription, else forwards TRACK_STATUS to every candidate SUBSCRIBE would try, concurrently, and combines the answers as a SUBSCRIBE_OK would (the entry's Track Properties, else the first answer's; the largest LARGEST_OBJECT), or gives the refusal SUBSCRIBE would. A forwarding TRACK_STATUS counts against `Config.MaxSubscriptionsPerSession` (§13.1), and each upstream round trip is bounded at 5s, counting as TIMEOUT (§13.6). Relay policy: not to the requester's own session while a TRACK_STATUS for the track to it is in flight (a loop, §6.2). |
| 10.15 | TRACK_STATUS | 0x0D | DONE | Reply via REQUEST_OK, then FIN; any follow-up from the requester closes the session. The relay answers from an Established subscription, else forwards TRACK_STATUS to every candidate SUBSCRIBE would try, concurrently, and combines the answers as a SUBSCRIBE_OK would (the entry's Track Properties, else the first answer's; the largest LARGEST_OBJECT), or gives the refusal SUBSCRIBE would. A forwarding TRACK_STATUS counts against `Config.MaxSubscriptionsPerSession` (§13.1), and the round is bounded at 5s, a silent candidate counting as TIMEOUT (§13.6). Concurrent requests for a track share one round, each answered with its own INCLUDE_PROPERTIES; a requester's STOP_SENDING stops only its wait. Relay policy: not to the requester's own session while a TRACK_STATUS for the track to it is in flight (a loop, §6.2). |
| 10.16 | PUBLISH_NAMESPACE | 0x06 | DONE | |
| 10.17 | NAMESPACE | 0x08 | DONE | Per namespace, counted over local and remote sources. |
| 10.18 | NAMESPACE_DONE | 0x0E | DONE | Never before its NAMESPACE. |
Expand Down Expand Up @@ -538,13 +538,11 @@ Relay:
by a relay peer (§6.2 has no loop protection), so it declines the second hop
rather than loop. Self-subscriptions are otherwise "identical" (§5.1).
- Filters are not aggregated upstream (§6.3.1 SHOULD).
- Concurrent TRACK_STATUS requests for one track with no Established
subscription are each forwarded upstream, not coalesced.
- A SUBSCRIBE the relay cannot open to a candidate for want of bidi-stream
credit is answered DOES_NOT_EXIST with no Retry Interval ("SHOULD NOT be
retried"), and outranks another candidate's GOING_AWAY or TIMEOUT, though
the publisher is there; EXCESSIVE_LOAD with a retry would fit better
(§10.6.2).
- A shared TRACK_STATUS round asks the candidates known when it started: a
request joining it later is not answered by a publisher that arrived in
between. A TRACK_STATUS loop through different sessions (the relay's
request reaching it back over another one) is not detected, and its round
waits out its 5s bound.

Open questions for interop: whether an End of Range marker carries an Object
Payload Length (Figure 28 vs §11.4.4.2), and whether EXPIRES may appear in
Expand Down
7 changes: 7 additions & 0 deletions pkg/relay/export_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,3 +59,10 @@ func SetTrackStatusTimeout(d time.Duration) (restore func()) {
prev := trackStatusTimeout.Swap(int64(d))
return func() { trackStatusTimeout.Store(prev) }
}

// SetTestHookTrackStatusJoined installs hook as testHookTrackStatusJoined and
// returns the restore.
func SetTestHookTrackStatusJoined(hook func(track.FullTrackName)) (restore func()) {
prev := testHookTrackStatusJoined.Swap(&hook)
return func() { testHookTrackStatusJoined.Store(prev) }
}
21 changes: 19 additions & 2 deletions pkg/relay/handler_subscribe.go
Original file line number Diff line number Diff line change
Expand Up @@ -569,12 +569,16 @@ func preferCandidateErr(last, err error) error {
}

// candidateRetry is the Retry Interval err is answered with: for a
// GOING_AWAY without one, the most the relay's own jittered one can be (see
// [upstreamRejection]), else the upstream's.
// GOING_AWAY without one, or a request not opened for want of stream credit,
// the most the relay's own jittered one can be (see [upstreamRejection]), else
// the upstream's.
func candidateRetry(err error) uint64 {
if errors.Is(err, errGoingAway) {
return goingAwayRetry
}
if errors.Is(err, session.ErrNoStreamCredit) {
return excessiveLoadRetryMax
}
rej, ok := errors.AsType[*session.RequestRejectedError](err)
if !ok {
return 0
Expand Down Expand Up @@ -979,6 +983,15 @@ func upstreamRejection(err error) *session.RequestRejectedError {
RetryInterval: retryIntervalAfter(goingAwayRetryAfter),
}
}
if errors.Is(err, session.ErrNoStreamCredit) {
// A live publisher the relay could not open a request to for want
// of stream credit: it "cannot process the request at this time"
// (§10.6.2), and may shortly.
return &session.RequestRejectedError{
Code: moqt.RequestExcessiveLoad,
RetryInterval: retryIntervalAfter(excessiveLoadRetry),
}
}
up, ok := errors.AsType[*session.RequestRejectedError](err)
if !ok {
return &session.RequestRejectedError{Code: moqt.RequestDoesNotExist}
Expand Down Expand Up @@ -1014,6 +1027,10 @@ const (
goingAwayRetry = uint64(goingAwayRetryAfter/time.Millisecond) * 3 / 2
)

// excessiveLoadRetryMax is the largest Retry Interval the jitter can make of
// excessiveLoadRetry.
const excessiveLoadRetryMax = uint64(excessiveLoadRetry/time.Millisecond) * 3 / 2

// isTrackPropertiesErr reports whether err is a Track Properties validation
// failure from [session.Session.Subscribe] or [session.Session.Fetch]: an
// unknown Mandatory Track Property, or Track Properties that do not parse.
Expand Down
89 changes: 83 additions & 6 deletions pkg/relay/handler_track_status.go
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ func (h *sessionHandler) handleTrackStatus(ctx context.Context, req *session.Req
h.rejectExcessiveLoad(ctx, req, "subscription")
return
}
oks, err := h.trackStatusUpstream(ctx, req, fullName)
oks, err := h.forwardTrackStatus(ctx, req, fullName)
h.limiter.releaseSub()
if len(oks) == 0 {
h.rejectTrackStatus(ctx, req, err)
Expand Down Expand Up @@ -125,23 +125,101 @@ func (h *sessionHandler) rejectTrackStatus(ctx context.Context, req *session.Req
}
}

// testHookTrackStatusJoined, when set by a test, runs once a request has
// joined, or started, the forwarded TRACK_STATUS round for its track.
var testHookTrackStatusJoined atomic.Pointer[func(track.FullTrackName)]

// trackStatusRounds is a relay's forwarded TRACK_STATUS rounds in flight, by
// track; see [sessionHandler.forwardTrackStatus]. Its ctx ends when the relay
// stops.
type trackStatusRounds struct {
ctx context.Context
end context.CancelFunc

mu sync.Mutex
rounds map[track.Key]*trackStatusRound
}

// trackStatusRound is one round's result, set before done closes.
type trackStatusRound struct {
done chan struct{}
oks []*message.TrackStatusOK
err error
}

func newTrackStatusRounds() *trackStatusRounds {
ctx, end := context.WithCancel(context.Background())
return &trackStatusRounds{ctx: ctx, end: end, rounds: make(map[track.Key]*trackStatusRound)}
}

// forwardTrackStatus answers TRACK_STATUS for fullName from a round of
// [sessionHandler.trackStatusUpstream], shared with the other requests for the
// track that arrive while it is in flight (relay policy, so N requests do not
// make N upstream ones). The round runs on a relay-scoped goroutine, bounded
// by its timeout and by the relay stopping, not by any one requester, who
// stops waiting on its own STOP_SENDING.
//
// A request from a session the relay is asking about the track right now may
// be the relay's own TRACK_STATUS routed back to it on that session (§6.2),
// which the round waits on: it runs a round of its own instead, where the loop
// guard skips that session. A loop through other sessions (relay A's request
// reaching it back over a different session) is not detected; such a round
// waits out its timeout.
func (h *sessionHandler) forwardTrackStatus(
ctx context.Context,
req *session.Request,
fullName track.FullTrackName,
) ([]*message.TrackStatusOK, error) {
waitCtx, cancel := context.WithCancel(ctx)
defer cancel()
defer context.AfterFunc(req.Stream.Context(), cancel)()
key := fullName.Key()
if h.tracks.RequestPending(message.TypeTrackStatus, h.sess, key) {
return h.trackStatusUpstream(waitCtx, fullName)
}

rs := h.statusRounds
rs.mu.Lock()
round := rs.rounds[key]
if round == nil {
round = &trackStatusRound{done: make(chan struct{})}
rs.rounds[key] = round
// Started from this tracked handler, so Stop joins it.
h.relayGo(func() {
round.oks, round.err = h.trackStatusUpstream(rs.ctx, fullName)
rs.mu.Lock()
delete(rs.rounds, key)
rs.mu.Unlock()
close(round.done)
})
}
rs.mu.Unlock()
if hook := testHookTrackStatusJoined.Load(); hook != nil {
(*hook)(fullName)
}
select {
case <-round.done:
return round.oks, round.err
case <-waitCtx.Done():
return nil, waitCtx.Err()
}
}

// trackStatusUpstream forwards TRACK_STATUS for fullName, concurrently, to
// every candidate SUBSCRIBE would try (see [sessionHandler.subscribeUpstream]):
// each local publisher of a covering namespace and each relay Discovery
// resolves. It returns their TRACK_STATUS_OKs in that order and, when there
// are none, the refusal SUBSCRIBE would give ([preferCandidateErr]), nil when there
// was no candidate.
//
// trackStatusUpstreamTimeout bounds the whole forwarding, resolving the
// Discovery candidates included, and it all ends when the requester cancels
// (STOP_SENDING). A draining candidate is sent nothing
// trackStatusUpstreamTimeout bounds the whole round, resolving the Discovery
// candidates included. A draining candidate is sent nothing
// (§10.4). Relay policy, as for SUBSCRIBE and FETCH: the requester's own
// session is skipped while a TRACK_STATUS for the track to it is in flight,
// since this request may be that one routed back, and a second would loop
// (§6.2).
func (h *sessionHandler) trackStatusUpstream(
ctx context.Context,
req *session.Request,
fullName track.FullTrackName,
) ([]*message.TrackStatusOK, error) {
timeout := trackStatusUpstreamTimeout
Expand All @@ -150,7 +228,6 @@ func (h *sessionHandler) trackStatusUpstream(
}
upCtx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
defer context.AfterFunc(req.Stream.Context(), cancel)()

key := fullName.Key()
var lastErr error
Expand Down
26 changes: 17 additions & 9 deletions pkg/relay/relay.go
Original file line number Diff line number Diff line change
Expand Up @@ -299,6 +299,10 @@ type Relay struct {
// the FETCH. Shared across every session handler.
fetch *registry.FetchRouter

// statusRounds shares forwarded TRACK_STATUS rounds across every
// session handler; see [sessionHandler.forwardTrackStatus].
statusRounds *trackStatusRounds

// upstreams dials and pools relay-to-relay sessions for Discovery-driven
// cross-relay upstream SUBSCRIBE. nil when Config.Dialer is unset (the
// single-instance case); session handlers treat a nil pool as "no
Expand Down Expand Up @@ -394,14 +398,15 @@ func New(listener Listener, cfg Config) *Relay {
}

r := &Relay{
listener: listener,
cfg: cfg,
log: log.With("component", "relay"),
tracks: registry.NewTrackRegistry(trackOpts...),
names: registry.NewNamespaceRegistry(nameOpts...),
fetch: registry.NewFetchRouter(),
sessions: make(map[*session.Session]struct{}),
stopCh: make(chan struct{}),
listener: listener,
cfg: cfg,
log: log.With("component", "relay"),
tracks: registry.NewTrackRegistry(trackOpts...),
names: registry.NewNamespaceRegistry(nameOpts...),
fetch: registry.NewFetchRouter(),
statusRounds: newTrackStatusRounds(),
sessions: make(map[*session.Session]struct{}),
stopCh: make(chan struct{}),
}

// When a Dialer is configured, the relay can follow Discovery
Expand Down Expand Up @@ -642,7 +647,7 @@ func (r *Relay) serveSession(ctx context.Context, sess *session.Session, leg Leg

handler := newSessionHandler(
sess, r.log, r.tracks, r.names,
r.cfg.Authorizer, r.cfg.Metrics, leg, r.fetch, r.upstreams,
r.cfg.Authorizer, r.cfg.Metrics, leg, r.fetch, r.statusRounds, r.upstreams,
r.cfg.Discovery, r.cfg.RelayAddr,
r.cfg.SendQueueSize, r.cfg.MaxDropsBeforeReset, r.cfg.MaxFanoutLag,
r.cfg.MaxSubscriptionsPerSession, r.cfg.MaxNamespaceRequestsPerSession,
Expand Down Expand Up @@ -758,6 +763,9 @@ func (r *Relay) Stop(ctx context.Context) error {
if r.upstreams != nil {
r.upstreams.close()
}
// Forwarded TRACK_STATUS rounds are upstream work too: cut short
// here, joined by handlers.Wait below.
r.statusRounds.end()

// 7. Force-close anything still standing, with GOAWAY_TIMEOUT
// only where it is true: "the peer took too long to close the
Expand Down
7 changes: 7 additions & 0 deletions pkg/relay/session_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,9 @@ type sessionHandler struct {
auth Authorizer
metrics Metrics
fetch *registry.FetchRouter
// statusRounds shares forwarded TRACK_STATUS rounds across the relay's
// sessions; see [sessionHandler.forwardTrackStatus].
statusRounds *trackStatusRounds
// leg records whether this session was dialled by the relay
// (LegUpstream) or by the peer (LegLocal). Every [Metrics] call this
// handler makes carries it, so an operator can separate what the
Expand Down Expand Up @@ -90,6 +93,7 @@ func newSessionHandler(
metrics Metrics,
leg Leg,
fetch *registry.FetchRouter,
statusRounds *trackStatusRounds,
upstreams *upstreamPool,
discovery discovery.DiscoveryStore,
relayAddr string,
Expand All @@ -110,6 +114,7 @@ func newSessionHandler(
metrics: metrics,
leg: leg,
fetch: fetch,
statusRounds: statusRounds,
upstreams: upstreams,
discovery: discovery,
relayAddr: relayAddr,
Expand Down Expand Up @@ -425,6 +430,8 @@ const excessiveLoadRetry = time.Second

// retryIntervalAfter is a Retry Interval (§10.6.2) inviting a retry after d
// plus up to 50% jitter, encoded as milliseconds plus one.
//
//nolint:unparam // callers pass separate policy constants, equal only today.
func retryIntervalAfter(d time.Duration) uint64 {
ms := uint64(d / time.Millisecond) //nolint:gosec // G115: callers pass a positive constant duration.
return ms + rand.Uint64N(ms/2) + 1 //nolint:gosec // G404: retry jitter, not a secret.
Expand Down
56 changes: 54 additions & 2 deletions pkg/relay/session_upstream_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -539,10 +539,25 @@ func TestSubscribe_UnsupportedMandatoryPropertyOutranksMalformed(t *testing.T) {
}
}

// requireRetryableExcessiveLoad fails unless err is EXCESSIVE_LOAD inviting a
// retry after about a second (§10.6.2).
func requireRetryableExcessiveLoad(t *testing.T, err error) {
t.Helper()
requireRejectedWithCode(t, err, moqt.RequestExcessiveLoad)
if rej, _ := errors.AsType[*session.RequestRejectedError](
err,
); rej.RetryInterval < 1001 ||
rej.RetryInterval > 1500 {
t.Fatalf("Retry Interval = %d, want a retry after about 1s", rej.RetryInterval)
}
}

// TestRendezvous_NoStreamCreditEndsHold: a SUBSCRIBE the relay cannot even
// open to a live publisher, for want of bidi-stream credit, is not a sign the
// track has no publisher, so a RENDEZVOUS_TIMEOUT hold does not wait out its
// budget on it: the subscriber is answered at once, as without a hold.
// budget on it: the subscriber is answered at once, with EXCESSIVE_LOAD and a
// retry, since the relay "cannot process the request at this time"
// (§10.6.2).
func TestRendezvous_NoStreamCreditEndsHold(t *testing.T) {
t.Parallel()
subSess, teardown := connectRelay(t, relay.Config{})
Expand All @@ -552,12 +567,49 @@ func TestRendezvous_NoStreamCreditEndsHold(t *testing.T) {
done := subscribeRendezvous(t.Context(), subSess, 5*time.Second)
select {
case err := <-done:
requireRejectedWithCode(t, err, moqt.RequestDoesNotExist)
requireRetryableExcessiveLoad(t, err)
case <-time.After(time.Second):
t.Fatal("SUBSCRIBE held against a live publisher the relay had no stream credit for")
}
}

// TestSubscribe_NoStreamCreditOutranksNoPublisherYet: a candidate the relay
// has no stream credit for is a live publisher, so its EXCESSIVE_LOAD wins
// over another candidate's GOING_AWAY, whatever the order; and a forwarded
// TRACK_STATUS gets the same answer.
func TestSubscribe_NoStreamCreditOutranksNoPublisherYet(t *testing.T) {
t.Parallel()
for _, order := range []string{"credit-less first", "credit-less last"} {
t.Run(order, func(t *testing.T) {
t.Parallel()
subSess, teardown := connectRelay(t, relay.Config{})
t.Cleanup(teardown)
creditless := func() {
pub := dialAnotherClientWithLimits(t, subSess, -1, 0)
publishNS(t, pub, "video")
}
if order == "credit-less first" {
creditless()
}
refusingPublishers(t, subSess, refuse(moqt.RequestGoingAway, 2001))
if order == "credit-less last" {
creditless()
}
_, err := subSess.Subscribe(t.Context(), &message.Subscribe{Namespace: ns("video"), Name: []byte("cam1")})
requireRetryableExcessiveLoad(t, err)
})
}
t.Run("TRACK_STATUS", func(t *testing.T) {
t.Parallel()
subSess, teardown := connectRelay(t, relay.Config{})
t.Cleanup(teardown)
pub := dialAnotherClientWithLimits(t, subSess, -1, 0)
publishNS(t, pub, "video")
_, err := trackStatusCam1(t, subSess)
requireRetryableExcessiveLoad(t, err)
})
}

// TestSubscribe_OtherRefusalTieIsOrderFree: two refusals of the "any other"
// kind with the same Retry Interval, one answered as INTERNAL_ERROR (an
// UNAUTHORIZED about the relay's hop) and one passed on (EXCESSIVE_LOAD),
Expand Down
Loading
Loading