diff --git a/internal/cluster/ring_engine.go b/internal/cluster/ring_engine.go index c8ccfa6..fea2ea3 100644 --- a/internal/cluster/ring_engine.go +++ b/internal/cluster/ring_engine.go @@ -382,6 +382,24 @@ func (e *Engine) runCommands(ctx context.Context, tk *Token) error { } continue // owned elsewhere — ride to the owner } + if t.Restart { + // Owner-directed, exactly like a revoke: only the node holding the + // forward may act. Selecting it anywhere else would let a non-owner + // spawn a worker for somebody else's forward (observed live: a + // non-owner logged "restarted t4" and ran the forward itself). + // + // A restart whose forward is already gone is a no-op consumed by the + // lowest node, so a task can never ride forever if its owner + // departed before seeing it. + if owner := e.state.TopologyOwner(t); owner == e.ID { + target = t + break + } else if owner == "" && selfIsLowest { + target = t // forward gone — consume and re-enable nothing + break + } + continue // owned elsewhere — ride to the owner + } if selfIsLowest { target = t // generic forward create — lowest-load claim break @@ -505,20 +523,13 @@ func (e *Engine) runCommands(ctx context.Context, tk *Token) error { } } if claimed.Restart { - // Only the OWNER may act. Pending tasks are visible to every member, so - // without this guard whichever node happened to process the task would - // spawn a worker for a forward attributed to somebody else — an - // orphaned worker plus a duplicate claim, exactly what the duplicate - // guard above exists to prevent. (Observed live: a non-owner logged - // "restarted t4" and ran the forward itself.) - // - // A restart whose owner has vanished is not an error: the entry is - // re-enabled, OfflineReassign() will move it to pending on the next - // departure sweep, and the normal claim path then re-homes it. + // Only reached when this node owns the forward (or it is already gone + // and we are the fallback consumer — see the selection loop). owner := e.state.TopologyOwner(claimed) e.state.UpdateTopologyDisabled(claimed.Local.Name, claimed.Remote.Name, claimed.Link.RemotePort, false) if owner != e.ID { - log.Printf("ring[%s] skip restart %s: owned by %s", e.ID, claimed.ID, owner) + // Forward vanished before we got here: nothing to re-enable. + log.Printf("ring[%s] discard restart %s: no owner", e.ID, claimed.ID) continue } if e.Handler != nil { diff --git a/internal/cluster/ring_revoke_test.go b/internal/cluster/ring_revoke_test.go index d0d1149..c423664 100644 --- a/internal/cluster/ring_revoke_test.go +++ b/internal/cluster/ring_revoke_test.go @@ -345,29 +345,24 @@ func TestRestartFlagSurvivesTokenSerialization(t *testing.T) { } } -// TestRestartOnlyAppliedByOwner is the regression test for the bug the restart -// channel introduced: pending tasks are visible to EVERY member, so without an -// owner check whichever node processed the task spawned a worker for a forward -// attributed to another node. Observed live as a non-owner logging -// "restarted t4" and running somebody else's forward. +// TestRestartTaskReachesNonLocalOwner. func TestRestartOnlyAppliedByOwner(t *testing.T) { - // The forward is owned by n1. n1 and n2 both see the restart task. restarts := 0 h := &fakeHandler{load: Load{MemPct: 5, NetPct: 5}, restart: func(ctx context.Context, tk *Task) error { restarts++; return nil }} - // Owner node. + // Owner node claims the forward. owner := NewEngine("n1", "n1:7500", "u", "p", "0.71.0", nil, h, func(ctx context.Context, next string, tk *Token) error { return nil }, "n1:7500", true, "") owner.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) - if _, err := owner.OnToken(context.Background(), &Token{Cycle: 1, State: owner.state}); err != nil { + if _, err := owner.OnToken(context.Background(), &Token{Cycle: 1, State: owner.state, SentAt: 1}); err != nil { t.Fatal(err) } if len(owner.state.TopologyList()) != 1 { t.Fatalf("precondition: owner should hold the entry, got %+v", owner.state.TopologyList()) } - // A peer builds the SAME restart task but is not the owner. + // A peer adopts the same ring state but does not own the forward. peer := NewEngine("n2", "n2:7500", "u", "p", "0.71.0", nil, h, func(ctx context.Context, next string, tk *Token) error { return nil }, "n2:7500", false, "") peer.AdoptState(owner.state) @@ -376,26 +371,75 @@ func TestRestartOnlyAppliedByOwner(t *testing.T) { if tk == nil { t.Fatal("SubmitRestart returned nil") } - peer.state.UpdateTopologyDisabled("web", "frps1", 18081, false) - // The peer must NOT run the handler: it does not own the forward. - if _, err := peer.OnToken(context.Background(), &Token{Cycle: 2, State: peer.state}); err != nil { + // The peer processes one token: it must NOT apply the restart. + if _, err := peer.OnToken(context.Background(), &Token{Cycle: 2, State: peer.state, SentAt: 2}); err != nil { t.Fatal(err) } if restarts != 0 { t.Fatalf("a NON-owner applied the restart %d time(s) — it would spawn an orphaned "+ "worker for a forward owned by somebody else", restarts) } +} - // The owner must apply it. - restarts = 0 - owner.state.UpdateTopologyDisabled("web", "frps1", 18081, false) - owner.SubmitRestart(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, - store.Link{RemotePort: 18081}) - if _, err := owner.OnToken(context.Background(), &Token{Cycle: 3, State: owner.state}); err != nil { +// TestRestartTaskReachesNonLocalOwner: a restart published on a node that is NOT +// the owner must survive the token round-trip and be applied by the owner. +// +// The naive "if not owner then continue" implementation CONSUMED the task +// (ClaimPending had already removed it), so the owner never received it and the +// forward stayed stopped with nothing logged as an error. Deferral must +// re-queue it. +// +// Each OnToken is fed a FRESH token (as the real ring does — the successor's +// OnToken receives the state the predecessor returned), so the deferral is +// exercised once per hop rather than re-processing one snapshot. +func TestRestartTaskReachesNonLocalOwner(t *testing.T) { + ownerApplied, submitterApplied := 0, 0 + ownerEng := NewEngine("n2", "n2:7500", "u", "p", "0.71.0", nil, + &fakeHandler{load: Load{MemPct: 5, NetPct: 5}, + restart: func(ctx context.Context, tk *Task) error { ownerApplied++; return nil }}, + func(ctx context.Context, next string, tk *Token) error { return nil }, "n2:7500", false, "") + + ownerEng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) + if _, err := ownerEng.OnToken(context.Background(), &Token{Cycle: 1, State: ownerEng.state}); err != nil { t.Fatal(err) } - if restarts != 1 { - t.Fatalf("the OWNER applied the restart %d times, want 1", restarts) + if len(ownerEng.state.TopologyList()) != 1 { + t.Fatalf("precondition: n2 should own it, got %+v", ownerEng.state.TopologyList()) + } + + subEng := NewEngine("n1", "n1:7500", "u", "p", "0.71.0", nil, + &fakeHandler{load: Load{MemPct: 5, NetPct: 5}, + restart: func(ctx context.Context, tk *Task) error { submitterApplied++; return nil }}, + func(ctx context.Context, next string, tk *Token) error { return nil }, "n1:7500", true, "") + subEng.AdoptState(ownerEng.state) + + subEng.SubmitRestart(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, + store.Link{RemotePort: 18081}) + + // One hop through the non-owner: it must defer, not apply and not drop. + out, err := subEng.OnToken(context.Background(), &Token{Cycle: 2, State: subEng.state, SentAt: 1}) + if err != nil { + t.Fatal(err) + } + if submitterApplied != 0 { + t.Fatal("the non-owner applied a restart for a forward it does not own") + } + var carried *Task + for _, tk := range out.State.PendingList() { + if tk.Restart && tk.Local.Name == "web" { + carried = tk + } + } + if carried == nil { + t.Fatal("the non-owner CONSUMED the restart task; the owner would never receive it") + } + + // The owner applies the carried task. + if _, err := ownerEng.OnToken(context.Background(), &Token{Cycle: 3, State: out.State, SentAt: 2}); err != nil { + t.Fatal(err) + } + if ownerApplied != 1 { + t.Fatalf("the owner applied the restart %d times, want 1", ownerApplied) } }