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
18 changes: 6 additions & 12 deletions STATUS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down
9 changes: 8 additions & 1 deletion pkg/relay/handler_subscribe.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down
102 changes: 102 additions & 0 deletions pkg/relay/pool_goaway_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
55 changes: 36 additions & 19 deletions pkg/relay/relay.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -773,29 +785,30 @@ 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) })
}
}

// drainStraggler runs the GOAWAY grace + force-close lifecycle for a single
// 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
Expand All @@ -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")
Expand All @@ -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)
Expand Down
21 changes: 12 additions & 9 deletions pkg/relay/relay_upstream_pool.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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 {
Expand All @@ -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)),
Expand All @@ -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
Expand All @@ -181,14 +183,15 @@ 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)
if p.fanIn > 0 && len(out) >= p.fanIn {
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
Expand Down
7 changes: 5 additions & 2 deletions pkg/relay/relay_upstream_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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))
}
Expand All @@ -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)
}
Expand Down
Loading
Loading