diff --git a/STATUS.md b/STATUS.md index 9a117cb0..364ce8b9 100644 --- a/STATUS.md +++ b/STATUS.md @@ -546,19 +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. -- 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). -- 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. +- 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.go b/pkg/relay/relay.go index 04d40963..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 @@ -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 @@ -695,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 @@ -773,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) }) } } @@ -794,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 @@ -806,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") @@ -832,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/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) } 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: 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") + } +}