diff --git a/internal/cluster/ring_engine.go b/internal/cluster/ring_engine.go index c10a9b8..c8ccfa6 100644 --- a/internal/cluster/ring_engine.go +++ b/internal/cluster/ring_engine.go @@ -505,10 +505,22 @@ func (e *Engine) runCommands(ctx context.Context, tk *Token) error { } } if claimed.Restart { - // Re-enable in place and hand the worker back to this node (the - // entry's owner). This must NOT create a second topology entry, which - // is why it skips the claim bookkeeping below. + // 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. + 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) + continue + } if e.Handler != nil { if err := e.Handler.Restart(ctx, claimed); err != nil { e.state.PendingTasks[claimed.ID] = claimed diff --git a/internal/cluster/ring_revoke_test.go b/internal/cluster/ring_revoke_test.go index 6550084..d0d1149 100644 --- a/internal/cluster/ring_revoke_test.go +++ b/internal/cluster/ring_revoke_test.go @@ -344,3 +344,58 @@ func TestRestartFlagSurvivesTokenSerialization(t *testing.T) { t.Fatal("a restart task must not also read as a revocation") } } + +// 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. +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 := 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 { + 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. + 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) + tk := peer.SubmitRestart(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, + store.Link{RemotePort: 18081}) + 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 { + 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 { + t.Fatal(err) + } + if restarts != 1 { + t.Fatalf("the OWNER applied the restart %d times, want 1", restarts) + } +}