From 427de7f9c79e3d1e324e2fddd369fa64e542e741 Mon Sep 17 00:00:00 2001 From: Vsevolod Strukchinsky Date: Sun, 27 Sep 2026 19:20:11 +0500 Subject: [PATCH 1/3] fix(relay): close the sessions Start's ctx ends, so Stop cannot hang on them MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Start documents cancelling its ctx as terminating live sessions, but it only ended each session's handler: the session was unregistered, so Stop no longer closed it, and left open. A relay-scoped reader of an upstream SUBSCRIBE on such a session then never returned, and Stop waited on it without bound. A handler with an inbound subgroup stream open did not even end, so its session stayed open and registered. handleConn now closes the session with NO_ERROR (§3.5: no GOAWAY was sent, so none ran out) as soon as Start's ctx ends, via context.AfterFunc, and again once serveSession returns if the ctx has ended: stop can otherwise win against an AfterFunc that has not started yet. Run is unaffected: it hands Start a ctx that is never cancelled. Pooled upstream sessions run under the pool's context and are closed by Stop as before. TestRelay_StartCtxClosesSessions, idle and with a publisher's subgroup stream open, was verified red before the change. Closing only after serveSession returns fails the open-stream case every time. A bare `defer stop()` failed the idle case intermittently under -race with -count=20. Co-Authored-By: Claude Opus 5.5 (1M context) --- STATUS.md | 6 --- pkg/relay/relay.go | 12 ++++++ pkg/relay/start_ctx_test.go | 80 +++++++++++++++++++++++++++++++++++++ 3 files changed, 92 insertions(+), 6 deletions(-) create mode 100644 pkg/relay/start_ctx_test.go diff --git a/STATUS.md b/STATUS.md index 9a117cb0..0f908c20 100644 --- a/STATUS.md +++ b/STATUS.md @@ -553,12 +553,6 @@ Relay: period and ignores `Stop`'s ctx, so a cancelled `Stop` can still wait up to `GoawayTimeout` for it. - Filters are not aggregated upstream (§6.3.1 SHOULD). -- Cancelling `Start`'s ctx ends each session's handler without closing the - session, and drops it from the set `Stop` closes; `Stop` then waits without - bound on the reader of an upstream SUBSCRIBE on such a session, whose stream - stays open. `Start`'s doc says the cancel "terminates live sessions". - Reproduced by two relays wired through Discovery and stopped in `t.Cleanup` - (after `t.Context` ends) while a cross-relay subscription is live. Documentation: diff --git a/pkg/relay/relay.go b/pkg/relay/relay.go index 04d40963..9710e864 100644 --- a/pkg/relay/relay.go +++ b/pkg/relay/relay.go @@ -595,7 +595,19 @@ func (r *Relay) handleConn(ctx context.Context, conn session.Conn) { return } + // Start documents cancelling its ctx as terminating live sessions, and a + // handler ending does not close its session: left open, a relay-scoped + // reader on one of its request streams, and with it Stop, would wait on + // it. NO_ERROR (§3.5): no GOAWAY was sent, so none ran out. + closeSess := func() { _ = sess.Close(moqt.SessionNoError, "relay: stopped") } + stop := context.AfterFunc(ctx, closeSess) r.serveSession(ctx, sess, LegLocal) + // The handler can end on ctx before AfterFunc has run closeSess, and stop + // then keeps it from running at all. + stop() + if ctx.Err() != nil { + closeSess() + } } // serveSession runs the per-session lifecycle for a Session that has already diff --git a/pkg/relay/start_ctx_test.go b/pkg/relay/start_ctx_test.go new file mode 100644 index 00000000..b4390d3d --- /dev/null +++ b/pkg/relay/start_ctx_test.go @@ -0,0 +1,80 @@ +package relay_test + +import ( + "context" + "testing" + "time" + + "github.com/floatdrop/moq-go/pkg/moqt/message" + "github.com/floatdrop/moq-go/pkg/moqt/session" + "github.com/floatdrop/moq-go/pkg/relay" +) + +// TestRelay_StartCtxClosesSessions: cancelling Start's ctx "terminates live +// sessions immediately" (see [relay.Relay.Start]): each is closed, not just +// left without a handler. Otherwise a publisher's session stays open with the +// relay's upstream SUBSCRIBE on it, and Stop, which no longer counts the +// session, waits without bound on that SUBSCRIBE's reader. +// +// With a subgroup stream from the publisher still open, the handler cannot +// end before the session does, so only closing it at once ends it. +func TestRelay_StartCtxClosesSessions(t *testing.T) { + t.Parallel() + t.Run("idle", func(t *testing.T) { + t.Parallel() + testStartCtxClosesSessions(t, false) + }) + t.Run("subgroup stream open", func(t *testing.T) { + t.Parallel() + testStartCtxClosesSessions(t, true) + }) +} + +func testStartCtxClosesSessions(t *testing.T, streamOpen bool) { + ctx, cancel := context.WithCancel(t.Context()) + tr := startTestRelay(ctx, relay.Config{}) + pub := dialClient(t, tr) + subs := acceptSubscribes(t, pub) + publishNS(t, pub, "video") + subSess := dialClient(t, tr) + subscribeCam1(t, subSess) + accepted := awaitAcceptedSubscribe(t, subs, "video/cam1") + if streamOpen { + sg, err := accepted.pub.OpenSubgroup( + message.SubgroupHeader{GroupID: 0, SubgroupIDMode: message.SubgroupIDExplicit}, + ) + if err != nil { + t.Fatalf("OpenSubgroup: %v", err) + } + if err := sg.WriteObjectAt(0, &message.SubgroupObject{Payload: []byte("a")}); err != nil { + t.Fatalf("WriteObjectAt: %v", err) + } + // Forwarded, so the relay is reading the stream when ctx ends. + ds, err := subSess.AcceptDataStream(t.Context()) + if err != nil { + t.Fatalf("AcceptDataStream: %v", err) + } + if _, err := ds.(*session.IncomingSubgroupStream).ReadObject(); err != nil { + t.Fatalf("ReadObject: %v", err) + } + } + + cancel() + tr.requireStartReturned(t) + select { + case <-pub.Done(): + case <-time.After(2 * time.Second): + t.Fatal("the publisher's session is still open after Start's ctx ended") + } + + stopped := make(chan struct{}) + go func() { + defer close(stopped) + _ = tr.r.Stop(context.Background()) + }() + select { + case <-stopped: + case <-time.After(2 * time.Second): + t.Fatal("Stop did not return after Start's ctx ended") + } +} From 39e4f1359b87292ee306bd776d952fedd7ec7681 Mon Sep 17 00:00:00 2001 From: Vsevolod Strukchinsky Date: Sun, 27 Sep 2026 19:23:16 +0500 Subject: [PATCH 2/3] =?UTF-8?q?fix(relay):=20bound=20a=20straggler=20sessi?= =?UTF-8?q?on's=20drain=20by=20Stop's=20ctx=20(=C2=A73.5)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A session that registers after Stop took its snapshot drains on its own: GOAWAY, grace period, force-close. That wait ignored Stop's ctx, so a cancelled Stop could still wait up to GoawayTimeout for it, although Stop's bulk drain cuts short on the same ctx. beginShutdown now records Stop's ctx in place of the shuttingDown flag, and drainStraggler also ends on it, closing with NO_ERROR: the grace period did not run out, so GOAWAY_TIMEOUT does not apply (§3.5), as in the bulk drain. TestRelay_addSessionDrainsStraggler gains a case with a one-hour grace period and Stop's ctx ending at 100ms; it fails with the ctx case removed. Co-Authored-By: Claude Opus 5.5 (1M context) --- STATUS.md | 3 -- pkg/relay/relay.go | 43 ++++++++++++++++------------ pkg/relay/shutdown_straggler_test.go | 32 ++++++++++++++------- 3 files changed, 45 insertions(+), 33 deletions(-) diff --git a/STATUS.md b/STATUS.md index 0f908c20..521c55fe 100644 --- a/STATUS.md +++ b/STATUS.md @@ -549,9 +549,6 @@ Relay: - A SUBSCRIBE whose only candidate upstream is a draining relay reached through the upstream pool gets DOES_NOT_EXIST, not the GOING_AWAY a draining local publisher yields: the pool skips such a relay before any request. -- A session that registers after `Stop` began is drained on its own grace - period and ignores `Stop`'s ctx, so a cancelled `Stop` can still wait up to - `GoawayTimeout` for it. - Filters are not aggregated upstream (§6.3.1 SHOULD). Documentation: diff --git a/pkg/relay/relay.go b/pkg/relay/relay.go index 9710e864..03561bb9 100644 --- a/pkg/relay/relay.go +++ b/pkg/relay/relay.go @@ -310,13 +310,13 @@ type Relay struct { // sessions tracks every Session that has completed SETUP and not yet // been torn down. Stop iterates it under sessionsMu to broadcast GOAWAY - // and to wait for drain. shuttingDown is set (under sessionsMu, by + // and to wait for drain. stopCtx, Stop's ctx, is set (under sessionsMu, by // beginShutdown) when Stop snapshots the set; addSession reads it under the // same lock to decide whether a newly-registered session is a straggler - // Stop's snapshot missed. - sessionsMu sync.Mutex - sessions map[*session.Session]struct{} - shuttingDown bool + // Stop's snapshot missed, and bounds that straggler's drain by it. + sessionsMu sync.Mutex + sessions map[*session.Session]struct{} + stopCtx context.Context // stopOnce guards Stop so the second caller short-circuits. stopCh is // closed by Stop to signal the accept loop to exit and to release any @@ -707,8 +707,8 @@ func (r *Relay) Stop(ctx context.Context) error { // doing potentially-blocking session work. The atomicity partitions // sessions cleanly: every session is either in this snapshot (its // drain is owned by steps 4–7 below) or registered later (it observes - // shuttingDown in addSession and owns its own drain) — never both. - sessions := r.beginShutdown() + // stopCtx in addSession and owns its own drain) — never both. + sessions := r.beginShutdown(ctx) // 4. Send GOAWAY to each session if a grace period is set. A // zero timeout means "don't bother with GOAWAY"; close @@ -785,20 +785,20 @@ func (r *Relay) Stop(ctx context.Context) error { func (r *Relay) addSession(s *session.Session, leg Leg) { r.sessionsMu.Lock() r.sessions[s] = struct{}{} - shuttingDown := r.shuttingDown + stopCtx := r.stopCtx r.sessionsMu.Unlock() r.cfg.Metrics.SessionOpened(leg) // Straggler cover: if shutdown was already in progress when we registered, - // Stop's snapshot — taken under sessionsMu together with the shuttingDown - // flag (see beginShutdown) — does NOT include this session, so Stop will - // neither GOAWAY nor close it. Own that lifecycle here. When shutdown began - // after we registered, shuttingDown is false and Stop's snapshot covers us; + // Stop's snapshot — taken under sessionsMu together with stopCtx (see + // beginShutdown) — does NOT include this session, so Stop will neither + // GOAWAY nor close it. Own that lifecycle here. When shutdown began after + // we registered, stopCtx is nil and Stop's snapshot covers us; // exactly one owner either way. The drain runs under r.handlers so Stop's // handlers.Wait joins it (safe: this runs inside serveSession, itself a // tracked handler, so the WaitGroup counter is already non-zero). - if shuttingDown { - r.handlers.Go(func() { r.drainStraggler(s) }) + if stopCtx != nil { + r.handlers.Go(func() { r.drainStraggler(stopCtx, s) }) } } @@ -806,8 +806,9 @@ func (r *Relay) addSession(s *session.Session, leg Leg) { // session that registered after Stop snapshotted the live-session set, so // Stop's bulk drain (Stop steps 4–7) does not cover it. It mirrors that bulk // drain for one session: GOAWAY, wait for the peer to drain or the grace period -// to elapse, then force-close. Spawned by addSession only during shutdown. -func (r *Relay) drainStraggler(s *session.Session) { +// to elapse, or Stop's ctx to end, then force-close. Spawned by addSession only +// during shutdown. +func (r *Relay) drainStraggler(stopCtx context.Context, s *session.Session) { goawayExpired := false if r.cfg.GoawayTimeout > 0 { sent := s.SendGoaway(r.cfg.GoawayTimeout, "") == nil @@ -818,6 +819,9 @@ func (r *Relay) drainStraggler(s *session.Session) { goawayExpired = sent case <-s.Done(): return // peer drained within the grace period + case <-stopCtx.Done(): + // Cut short, as Stop's bulk drain is: the grace period did not + // run out, so NO_ERROR (§3.5). } } _ = s.Close(shutdownCloseCode(goawayExpired), "relay shutdown") @@ -844,11 +848,12 @@ func (r *Relay) removeSession(s *session.Session, leg Leg) { // currently-registered sessions, atomically under sessionsMu. The atomicity is // what lets addSession partition sessions into exactly two non-overlapping // groups: those in the returned snapshot (drained by Stop) and those registered -// afterward (which see shuttingDown and drain themselves via drainStraggler). -func (r *Relay) beginShutdown() []*session.Session { +// afterward (which see stopCtx and drain themselves via drainStraggler, bounded +// by ctx, Stop's). +func (r *Relay) beginShutdown(ctx context.Context) []*session.Session { r.sessionsMu.Lock() defer r.sessionsMu.Unlock() - r.shuttingDown = true + r.stopCtx = ctx out := make([]*session.Session, 0, len(r.sessions)) for s := range r.sessions { out = append(out, s) diff --git a/pkg/relay/shutdown_straggler_test.go b/pkg/relay/shutdown_straggler_test.go index 4f111a73..f32b0640 100644 --- a/pkg/relay/shutdown_straggler_test.go +++ b/pkg/relay/shutdown_straggler_test.go @@ -41,16 +41,21 @@ func (c *stragglerCodeConn) CloseWithError(code uint64, reason string) error { // snapshot still goes through GOAWAY, grace and force-close, via // addSession's drainStraggler. It closes with GOAWAY_TIMEOUT only when it was // sent a GOAWAY and the grace period ran out (§3.5), and NO_ERROR when no -// GOAWAY is configured. +// GOAWAY is configured or Stop's ctx ends the drain first, as the bulk drain +// does. func TestRelay_addSessionDrainsStraggler(t *testing.T) { t.Parallel() for _, tc := range []struct { name string grace time.Duration - want moqt.SessionErrorCode + // stopAfter, when set, ends Stop's ctx that long into the drain. + stopAfter time.Duration + want moqt.SessionErrorCode }{ - {"grace period ran out", 150 * time.Millisecond, moqt.SessionGoawayTimeout}, - {"no GOAWAY configured", 0, moqt.SessionNoError}, + {"grace period ran out", 150 * time.Millisecond, 0, moqt.SessionGoawayTimeout}, + {"no GOAWAY configured", 0, 0, moqt.SessionNoError}, + // Stop's ctx bounds a straggler's drain as it does the bulk drain. + {"Stop's ctx ended", time.Hour, 100 * time.Millisecond, moqt.SessionNoError}, } { t.Run(tc.name, func(t *testing.T) { t.Parallel() @@ -76,14 +81,19 @@ func TestRelay_addSessionDrainsStraggler(t *testing.T) { clientSess, serverSess := cl.s, sv.s defer func() { _ = clientSess.Close(0, "") }() - // Simulate Stop having already begun: beginShutdown marks - // shuttingDown and snapshots the (still empty) session set. The - // straggler registers next. - if snap := r.beginShutdown(); len(snap) != 0 { + // Simulate Stop having already begun: beginShutdown records Stop's + // ctx and snapshots the (still empty) session set. The straggler + // registers next. + stopCtx, cancelStop := context.WithCancel(t.Context()) + defer cancelStop() + if tc.stopAfter > 0 { + time.AfterFunc(tc.stopAfter, cancelStop) + } + if snap := r.beginShutdown(stopCtx); len(snap) != 0 { t.Fatalf("beginShutdown snapshot = %d sessions, want 0", len(snap)) } - // addSession must observe shuttingDown and take ownership of the + // addSession must observe Stop's ctx and take ownership of the // drain. r.addSession(serverSess, LegLocal) @@ -97,11 +107,11 @@ func TestRelay_addSessionDrainsStraggler(t *testing.T) { } // ...and, because this client ignores the GOAWAY, force-close at - // the grace boundary so the session terminates. + // the grace boundary, or once Stop's ctx ends. select { case <-serverSess.Done(): case <-time.After(2 * time.Second): - t.Fatal("straggler session was not closed after the grace period") + t.Fatal("straggler session was not force-closed") } select { case code := <-codes: From 942585a6d3bc86cee0d0dd696581bf1c78f89381 Mon Sep 17 00:00:00 2001 From: Vsevolod Strukchinsky Date: Sun, 27 Sep 2026 19:28:26 +0500 Subject: [PATCH 3/3] =?UTF-8?q?fix(relay):=20refuse=20with=20GOING=5FAWAY?= =?UTF-8?q?=20when=20the=20only=20upstream=20relay=20is=20draining=20(?= =?UTF-8?q?=C2=A710.4)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A SUBSCRIBE whose only candidate was a relay reached through the upstream pool that had sent GOAWAY got DOES_NOT_EXIST: the pool skips such a relay before any request (§10.4), so no answer of its own existed. A draining local publisher yields GOING_AWAY ("The endpoint has received a GOAWAY", §10.6.2). resolveUpstreams now reports whether it skipped a draining relay, and subscribeUpstream counts that as a GOING_AWAY answer, ranked with the other "no publisher yet" errors; under RENDEZVOUS_TIMEOUT the hold continues, as before. Interpretation, marked in the code: GOING_AWAY is taken to outrank §10.2.6's DOES_NOT_EXIST for "no publisher is available", since the publisher is known and GOING_AWAY says to retry. The last "no publisher yet" answer still sets the code, so the order of candidates matters there; that, and an upstream relay's own GOING_AWAY answer being passed on as INTERNAL_ERROR, are recorded in STATUS.md. TestCrossRelay_DrainingUpstreamRelayAnswersGoingAway (alone, and with a local publisher refusing too) uses a Discovery store that keeps a draining relay listed, as an eventually-consistent backend may; it fails with the handler change removed. TestResolveUpstreamsSkipsGoingAwayRelay asserts the new result. Co-Authored-By: Claude Opus 5.5 (1M context) --- STATUS.md | 9 ++- pkg/relay/handler_subscribe.go | 9 ++- pkg/relay/pool_goaway_test.go | 102 +++++++++++++++++++++++++++++++ pkg/relay/relay_upstream_pool.go | 21 ++++--- pkg/relay/relay_upstream_test.go | 7 ++- 5 files changed, 133 insertions(+), 15 deletions(-) create mode 100644 pkg/relay/pool_goaway_test.go diff --git a/STATUS.md b/STATUS.md index 521c55fe..364ce8b9 100644 --- a/STATUS.md +++ b/STATUS.md @@ -546,10 +546,13 @@ Relay: flight: the relay cannot tell either from its own request routed back to it 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). -- A SUBSCRIBE whose only candidate upstream is a draining relay reached through - the upstream pool gets DOES_NOT_EXIST, not the GOING_AWAY a draining local - publisher yields: the pool skips such a relay before any request. - Filters are not aggregated upstream (§6.3.1 SHOULD). +- Among SUBSCRIBE candidates that only say the track has no publisher yet + (DOES_NOT_EXIST, TIMEOUT, draining), the last to answer sets the refusal + code, so a draining local publisher and a remote relay's DOES_NOT_EXIST + yield DOES_NOT_EXIST while the reverse yields GOING_AWAY. An upstream relay + that answers GOING_AWAY itself, before its GOAWAY reaches this relay, is + passed on as INTERNAL_ERROR. Documentation: diff --git a/pkg/relay/handler_subscribe.go b/pkg/relay/handler_subscribe.go index 969acc98..902ba4c3 100644 --- a/pkg/relay/handler_subscribe.go +++ b/pkg/relay/handler_subscribe.go @@ -628,7 +628,14 @@ func (h *sessionHandler) subscribeUpstream( establish(pub.Session, "local-publisher", extra) } - remotes := h.upstreams.resolveUpstreams(ctx, fullName.Namespace) + remotes, draining := h.upstreams.resolveUpstreams(ctx, fullName.Namespace) + // A draining relay was sent no request (§10.4); it answers as a draining + // publisher would, ranked with the other candidates' errors. Taken to + // outrank §10.2.6's DOES_NOT_EXIST for "no publisher is available": the + // publisher is known, and GOING_AWAY (§10.6.2) says to retry. + if draining && !isTrackPropertiesErr(lastErr) && awaitsPublisher(lastErr) { + lastErr = errGoingAway + } for _, remote := range remotes { establish(remote, "discovery-remote", remoteExtra) } diff --git a/pkg/relay/pool_goaway_test.go b/pkg/relay/pool_goaway_test.go new file mode 100644 index 00000000..41d6a6ab --- /dev/null +++ b/pkg/relay/pool_goaway_test.go @@ -0,0 +1,102 @@ +package relay_test + +import ( + "context" + "errors" + "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/moqt/track" + "github.com/floatdrop/moq-go/pkg/moqt/wire" + "github.com/floatdrop/moq-go/pkg/relay" + "github.com/floatdrop/moq-go/pkg/relay/discovery" +) + +// staleStore is a Discovery store that never forgets an advertisement, as an +// eventually-consistent backend may still list a relay that is draining. +type staleStore struct{ discovery.DiscoveryStore } + +func (staleStore) UnpublishNamespace(context.Context, wire.TrackNamespace, string) error { return nil } +func (staleStore) UnpublishTrack(context.Context, track.Key, string) error { return nil } +func (staleStore) Withdraw(context.Context, string) error { return nil } + +// TestCrossRelay_DrainingUpstreamRelayAnswersGoingAway: a SUBSCRIBE whose only +// candidate is a relay reached through the upstream pool that has sent GOAWAY +// is refused with GOING_AWAY, as with a draining local publisher: "The +// endpoint has received a GOAWAY and MAY reject new requests" (§10.6.2). The +// relay sends that relay no request (§10.4), so no answer of its own exists. +// +// A local publisher answering DOES_NOT_EXIST alongside does not change that: +// the draining relay's publisher may still have the track. +func TestCrossRelay_DrainingUpstreamRelayAnswersGoingAway(t *testing.T) { + t.Parallel() + t.Run("only candidate", func(t *testing.T) { + t.Parallel() + testDrainingUpstreamRelay(t, false) + }) + t.Run("local publisher refuses", func(t *testing.T) { + t.Parallel() + testDrainingUpstreamRelay(t, true) + }) +} + +// refuseAll answers every request on sess with DOES_NOT_EXIST. +func refuseAll(t *testing.T, sess *session.Session) { + go func() { + for { + req, err := sess.AcceptRequest(t.Context()) + if err != nil { + return + } + _ = req.RejectError(moqt.RequestDoesNotExist, "no cam1") + } + }() +} + +func testDrainingUpstreamRelay(t *testing.T, localRefuser bool) { + store := staleStore{discovery.NewMemoryStore()} + defer store.Close() + relayB := startTestRelay(t.Context(), relay.Config{ + Discovery: store, RelayAddr: "relay-B", GoawayTimeout: 10 * time.Second, + }) + relayA := startTestRelay(t.Context(), relay.Config{ + Discovery: store, RelayAddr: "relay-A", Dialer: dialerTo(nil, relayB), + }) + + pub := dialClient(t, relayB) + publishNS(t, pub, "video") + refuseAll(t, pub) + if localRefuser { + local := dialClient(t, relayA) + publishNS(t, local, "video") + refuseAll(t, local) + } + subSess := dialClient(t, relayA) + subscribe := func() error { + _, err := subSess.Subscribe(t.Context(), &message.Subscribe{Namespace: ns("video"), Name: []byte("cam1")}) + return err + } + // Dials relay B, whose publisher has no cam1. + requireRejectedWithCode(t, subscribe(), moqt.RequestDoesNotExist) + + stopped := make(chan struct{}) + go func() { + defer close(stopped) + _ = relayB.r.Stop(context.Background()) + }() + defer func() { + // B's drain ends once its sessions do: the publisher's and the + // pooled one relayA.stop closes. + _ = pub.Close(moqt.SessionNoError, "done") + relayA.stop(t) + <-stopped + }() + + waitFor(t, 2*time.Second, func() bool { + rej, ok := errors.AsType[*session.RequestRejectedError](subscribe()) + return ok && rej.Code == moqt.RequestGoingAway + }, "a SUBSCRIBE whose only upstream relay is draining was never refused with GOING_AWAY") +} diff --git a/pkg/relay/relay_upstream_pool.go b/pkg/relay/relay_upstream_pool.go index 5ef0c5d9..4093c243 100644 --- a/pkg/relay/relay_upstream_pool.go +++ b/pkg/relay/relay_upstream_pool.go @@ -127,7 +127,9 @@ func newUpstreamPool(cfg upstreamPoolConfig) *upstreamPool { // lease TTL falls through transparently to the next-ranked one. // // Returns nil when Discovery knows no usable remote (none advertised, only this -// relay itself, or every candidate failed to dial). Discovery-lookup and +// relay itself, or every candidate failed to dial). draining reports whether a +// candidate was skipped for having sent GOAWAY (§10.4), which the caller +// answers as it would a draining publisher. Discovery-lookup and // per-peer dial failures are logged and treated as "skip that candidate" — // consistent with the best-effort advertise side: the local registry / a clean // SUBSCRIBE rejection is the fallback, never a torn-down session. @@ -137,9 +139,12 @@ func newUpstreamPool(cfg upstreamPoolConfig) *upstreamPool { // to itself. Duplicate RelayAddrs collapse to one session (the pool keys by // address). Multi-hop cycle detection (A→B→C→A) is out of scope — see the // package limitations. -func (p *upstreamPool) resolveUpstreams(ctx context.Context, ns wire.TrackNamespace) []*session.Session { +func (p *upstreamPool) resolveUpstreams( + ctx context.Context, + ns wire.TrackNamespace, +) (out []*session.Session, draining bool) { if p == nil || p.discovery == nil { - return nil + return nil, false } infos, err := p.discovery.FindNamespace(ctx, ns) if err != nil { @@ -151,7 +156,7 @@ func (p *upstreamPool) resolveUpstreams(ctx context.Context, ns wire.TrackNamesp p.log.LogAttrs(ctx, slog.LevelWarn, "upstream pool: FindNamespace failed", slog.String("namespace", fmt.Sprintf("%v", ns)), slog.String("err", err.Error())) - return nil + return nil, false } p.log.LogAttrs(ctx, slog.LevelInfo, "upstream pool: FindNamespace resolved", slog.String("namespace", fmt.Sprintf("%v", ns)), @@ -161,10 +166,7 @@ func (p *upstreamPool) resolveUpstreams(ctx context.Context, ns wire.TrackNamesp // takes the same top-fanIn upstreams everywhere. rankByAffinity(ns, infos) - var ( - out []*session.Session - seen = make(map[string]bool, len(infos)) - ) + seen := make(map[string]bool, len(infos)) for _, info := range infos { if info.RelayAddr == "" || info.RelayAddr == p.relayAddr || seen[info.RelayAddr] { continue // self / unaddressable / already dialled this address @@ -181,6 +183,7 @@ func (p *upstreamPool) resolveUpstreams(ctx context.Context, ns wire.TrackNamesp if goingAway(sess) { // §10.4: a draining relay takes no new requests, so it must not // hold a fan-in slot. + draining = true continue } out = append(out, sess) @@ -188,7 +191,7 @@ func (p *upstreamPool) resolveUpstreams(ctx context.Context, ns wire.TrackNamesp break // opt-in bound reached; deeper candidates are the fallback pool } } - return out + return out, draining } // rankByAffinity sorts infos in place by descending rendezvous (HRW) weight for diff --git a/pkg/relay/relay_upstream_test.go b/pkg/relay/relay_upstream_test.go index abd57bca..4d6cb282 100644 --- a/pkg/relay/relay_upstream_test.go +++ b/pkg/relay/relay_upstream_test.go @@ -146,7 +146,7 @@ func TestResolveUpstreamsSkipsGoingAwayRelay(t *testing.T) { }) defer p.close() - first := p.resolveUpstreams(ctx, ns) + first, _ := p.resolveUpstreams(ctx, ns) if len(first) != 1 { t.Fatalf("resolveUpstreams = %d sessions, want 1 (fan-in 1)", len(first)) } @@ -170,7 +170,10 @@ func TestResolveUpstreamsSkipsGoingAwayRelay(t *testing.T) { t.Fatal("the pooled session never saw the GOAWAY") } - second := p.resolveUpstreams(ctx, ns) + second, draining := p.resolveUpstreams(ctx, ns) + if !draining { + t.Error("resolveUpstreams did not report the draining relay it skipped") + } if len(second) != 1 || second[0] == first[0] { t.Fatalf("resolveUpstreams after GOAWAY returned the draining %s again; want the next-ranked relay", top) }