diff --git a/cmd/webui4frpc/main.go b/cmd/webui4frpc/main.go index da733e1..55be602 100644 --- a/cmd/webui4frpc/main.go +++ b/cmd/webui4frpc/main.go @@ -351,20 +351,23 @@ func main() { ring.SetPeerPersist(func(peersJSON string) error { return st.SetClusterPeers(peersJSON) }) // Re-apply local store state onto the ring topology after each state - // adoption, and learn peers' disabled flags back into the local store. + // adoption, and learn the ring's disabled decisions back into the store. // - // Without the first half, group changes made via HTTP handlers + // Without the group half, group changes made via HTTP handlers // (POST /forwards/assign) are overwritten by the next e.state = tk.State - // and never propagate to other nodes. + // and never propagate. // - // The second half closes the loop for stop/start. The links table is - // node-local, so a stop decided on one machine used to leave every other - // node's copy reading "enabled" — including the node that OWNS the forward - // and holds its worker, which would then re-spawn it. Adoption is the point - // at which we have the cluster's view, so fold it into the local store here. - // The flag can only ever be learned from an entry that still exists, so a - // forward already revoked away is (correctly) treated as unknown rather than - // resurrected as stopped. + // The disabled half makes the RING authoritative for stop/start. That is the + // only source all members can agree on: the links table is a per-node copy, so + // treating it as authoritative let two nodes disagree indefinitely (observed + // live: one node reading disabled=1 while the owner still read 0 and kept its + // worker running). This is a write-through, not a listener — it runs on every + // adoption, so any peer's decision lands here within a token cycle. + // + // Ordering matters in one direction only: the engine has already re-applied + // this node's not-yet-confirmed local decisions onto the adopted state (see + // reconcileLocalDisabled), so reading the ring here can never clobber an + // action the user just took on this node. ring.SetTopologySync(func() { links, err := st.ListLinks() if err != nil { @@ -378,11 +381,8 @@ func main() { if !known { continue } - // A local explicit decision must win: the node that served the - // request already wrote its own store, and its topology entry is the - // one the flag is riding on. Anything else is a peer's decision. if ln, found, err := st.LinkByTriple(te.Local.Name, te.Remote.Name, te.Link.RemotePort); err == nil && found && ln.Disabled == disabled { - continue // already agrees + continue // store already agrees with the ring } _ = st.ReconcileLinkDisabled(te.Local.Name, te.Remote.Name, te.Link.RemotePort, disabled) } diff --git a/internal/cluster/ring_disabled_test.go b/internal/cluster/ring_disabled_test.go new file mode 100644 index 0000000..7280d0e --- /dev/null +++ b/internal/cluster/ring_disabled_test.go @@ -0,0 +1,157 @@ +package cluster + +import ( + "context" + "testing" + + "webui4frpc/internal/store" +) + +// staleTokenWith copies the engine's current state but rewrites the given +// forward's disabled flag, standing in for a token a neighbour captured before +// this node's decision. +func staleTokenWith(eng *Engine, local, remote string, port int, disabled bool) State { + s := eng.state + out := State{ + LeaderID: s.LeaderID, + Nodes: s.Nodes, + Cycle: s.Cycle, + Topology: map[string]*TopoEntry{}, + } + for id, e := range s.Topology { + cp := *e + if cp.Local.Name == local && cp.Remote.Name == remote && cp.Link.RemotePort == port { + cp.Link.Disabled = disabled + cp.Active = !disabled + } + out.Topology[id] = &cp + } + return out +} + +// TestLocalDisableSurvivesAdoption is the regression test for the wipe that made +// the ring's authority unachievable. +// +// OnToken/AdoptState replace e.state wholesale. A stop decided between two token +// cycles was therefore erased by the next adoption whenever the incoming token +// had been captured before the decision — so the flag never reached the other +// members and the forward kept running. Written as a probe BEFORE the fix: it +// reported "after adopting a stale token: known=true disabled=false". +func TestLocalDisableSurvivesAdoption(t *testing.T) { + eng := newTestEngine("n1", true) + eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state, SentAt: 1}); err != nil { + t.Fatal(err) + } + + // Local decision: stop it. + eng.UpdateTopologyDisabled("web", "frps1", 18081, true) + if d, _ := eng.TopologyDisabled("web", "frps1", 18081); !d { + t.Fatal("precondition: the local stop should be visible") + } + + // A token captured BEFORE the decision arrives. + stale := staleTokenWith(eng, "web", "frps1", 18081, false) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: stale, SentAt: 2}); err != nil { + t.Fatal(err) + } + // Assert on the STATE, not the TopologyDisabled() accessor: the accessor + // short-circuits to the pending local decision, so it would report the stop + // even when the state about to be forwarded still says "enabled". What + // matters is that the token this node sends on carries the decision. + assertForwardedDisabled(t, eng, "web", "frps1", 18081, true) +} + +// assertForwardedDisabled checks the disabled value a DOWNSTREAM node would see +// in the token this engine forwards — the value that actually propagates. +func assertForwardedDisabled(t *testing.T, eng *Engine, local, remote string, port int, want bool) { + t.Helper() + // Model one hop: a fresh engine adopts exactly what this one would send. + peer := newTestEngine("peer", false) + peer.AdoptState(eng.state) + got, known := peer.TopologyDisabled(local, remote, port) + if !known { + t.Fatalf("peer learned no entry for %s→%s:%d", local, remote, port) + } + if got != want { + t.Fatalf("the forwarded token carries disabled=%v, want %v — the decision did not "+ + "propagate, so the other members would keep the old state", got, want) + } +} + +// TestLocalReEnableSurvivesAdoption: a re-enable needs the same protection. The +// value false is easy to overlook — nothing "looks" stopped, yet losing it just +// as silently strands the forward in the stopped state. +func TestLocalReEnableSurvivesAdoption(t *testing.T) { + eng := newTestEngine("n1", true) + eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state, SentAt: 1}); err != nil { + t.Fatal(err) + } + eng.UpdateTopologyDisabled("web", "frps1", 18081, true) + eng.UpdateTopologyDisabled("web", "frps1", 18081, false) + + // A token that still says disabled (captured mid-stop) arrives. + mid := staleTokenWith(eng, "web", "frps1", 18081, true) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: mid, SentAt: 2}); err != nil { + t.Fatal(err) + } + assertForwardedDisabled(t, eng, "web", "frps1", 18081, false) +} + +// TestConfirmedDecisionStopsBeingAsserted: once the ring reports the same value +// the decision is confirmed everywhere, so this node must stop overriding — +// otherwise it would fight a later decision made elsewhere. +func TestConfirmedDecisionStopsBeingAsserted(t *testing.T) { + eng := newTestEngine("n1", true) + eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state, SentAt: 1}); err != nil { + t.Fatal(err) + } + + eng.UpdateTopologyDisabled("web", "frps1", 18081, true) + if len(eng.localDisabled) != 1 { + t.Fatalf("a pending decision should be tracked, got %v", eng.localDisabled) + } + // The ring comes back agreeing. + agreed := staleTokenWith(eng, "web", "frps1", 18081, true) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: agreed, SentAt: 2}); err != nil { + t.Fatal(err) + } + if len(eng.localDisabled) != 0 { + t.Fatalf("a confirmed decision must stop being asserted, still tracking %v", eng.localDisabled) + } + + // A peer now decides the opposite; this node must follow the ring. + peerSaysEnabled := staleTokenWith(eng, "web", "frps1", 18081, false) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 3, State: peerSaysEnabled, SentAt: 3}); err != nil { + t.Fatal(err) + } + if d, _ := eng.TopologyDisabled("web", "frps1", 18081); d { + t.Fatal("this node kept overriding a peer's later decision; the ring is not authoritative") + } +} + +// TestLocalDecisionsDoNotAccumulate: the tracking map is swept against the live +// topology, so a long-lived cluster cannot grow it without bound. +func TestLocalDecisionsDoNotAccumulate(t *testing.T) { + eng := newTestEngine("n1", true) + eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state, SentAt: 1}); err != nil { + t.Fatal(err) + } + + // Decisions for forwards this node has never seen. + eng.UpdateTopologyDisabled("ghost1", "frps1", 1, true) + eng.UpdateTopologyDisabled("ghost2", "frps1", 2, true) + if len(eng.localDisabled) != 2 { + t.Fatalf("expected 2 tracked decisions, got %d", len(eng.localDisabled)) + } + // A cycle runs: neither ghost is in the topology, so both are swept. + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: eng.state, SentAt: 2}); err != nil { + t.Fatal(err) + } + if len(eng.localDisabled) != 0 { + t.Fatalf("decisions for forwards outside the topology must be swept, got %v", eng.localDisabled) + } +} diff --git a/internal/cluster/ring_engine.go b/internal/cluster/ring_engine.go index 6f87fdd..8e9f11b 100644 --- a/internal/cluster/ring_engine.go +++ b/internal/cluster/ring_engine.go @@ -7,6 +7,7 @@ import ( "encoding/json" "fmt" "log" + "strconv" "sync" "time" @@ -55,6 +56,24 @@ type Engine struct { // group changes made via HTTP handlers between token cycles are overwritten // by the next state adoption and never propagate to other nodes. topologySync func() + // localDisabled records this node's own enabled/disabled decisions that have + // not yet been confirmed by the ring. It is what makes the RING authoritative + // without losing a decision that was made locally and has not yet had a + // chance to reach the other members. + // + // The problem it solves: OnToken/AdoptState do `e.state = tk.State`, a + // wholesale replacement. A stop decided between two token cycles therefore + // disappears on the very next adoption if that token was captured before the + // decision — the flag never reaches the rest of the ring, and the forward + // keeps running. (Verified with a probe before writing this: adopting a + // stale token flipped disabled back to false.) + // + // A key stays in this map until an adoption reports it disabled, at which + // point the whole cluster agrees and the key is dropped. Keys for forwards + // that vanish from the topology are swept on every update, so the map cannot + // grow without bound. The value false is what keeps a RE-ENABLE alive for the + // same reason a stop needs protecting. + localDisabled map[string]bool state State // myAddr maps our Node ID to the address peers dial. @@ -236,6 +255,11 @@ func (e *Engine) OnToken(ctx context.Context, tk *Token) (*Token, error) { // Re-apply local store overrides (group labels, disabled flags) onto // the freshly adopted topology so they survive state adoption and // propagate to all nodes via the next token forward. + // Re-assert decisions this node made locally that the ring has not yet + // confirmed, BEFORE the host's topology sync runs: the sync reads + // TopologyDisabled to learn the cluster's view, so the local decision must + // already be applied or a stop could be reported as "enabled" and dropped. + e.reconcileLocalDisabled() if e.topologySync != nil { e.topologySync() } @@ -802,6 +826,11 @@ func (e *Engine) AdoptState(s State) { for id, t := range kept { e.state.PendingTasks[id] = t } + // Re-assert decisions this node made locally that the ring has not yet + // confirmed, BEFORE the host's topology sync runs: the sync reads + // TopologyDisabled to learn the cluster's view, so the local decision must + // already be applied or a stop could be reported as "enabled" and dropped. + e.reconcileLocalDisabled() if e.topologySync != nil { e.topologySync() } @@ -990,16 +1019,75 @@ func (e *Engine) UpdateTopologyGroup(local, remote string, port int, group strin // UpdateTopologyDisabled marks a forward enabled/disabled in the topology so // the decision rides the next token cycle to every member. Paired with // TopologyDisabled, which the adoption hook uses to learn peers' decisions. +// +// The decision is also remembered locally until the ring confirms it, because +// adoption replaces the whole state: a token captured before this call would +// otherwise wash the decision out on the next cycle and it would never reach +// the other members. See the localDisabled field. func (e *Engine) UpdateTopologyDisabled(local, remote string, port int, disabled bool) bool { + if e.localDisabled == nil { + e.localDisabled = map[string]bool{} + } + e.localDisabled[disableKey(local, remote, port)] = disabled return e.state.UpdateTopologyDisabled(local, remote, port, disabled) } // TopologyDisabled reports the cluster's view of a forward's disabled flag. // The bool is false when the forward has no topology entry (nothing to learn). +// +// A pending local decision takes precedence over the adopted state: it has not +// had a chance to reach the other members yet, and reporting the adopted value +// would make the local store (and the status page) disagree with the user's +// most recent action. func (e *Engine) TopologyDisabled(local, remote string, port int) (bool, bool) { + if d, pending := e.localDisabled[disableKey(local, remote, port)]; pending { + return d, true + } return e.state.TopologyDisabled(local, remote, port) } +// reconcileLocalDisabled re-asserts this node's not-yet-confirmed disabled +// decisions onto the freshly adopted state. Called right after e.state is +// replaced, and paired with reconcileTopologySync (which pushes those decisions +// out to the ring). +// +// A decision is dropped once the adopted state reports the SAME value: at that +// point every member agrees and there is nothing left to protect. Entries whose +// forward no longer exists in the topology are dropped too, so the map tracks +// only live disagreements. +func (e *Engine) reconcileLocalDisabled() { + if len(e.localDisabled) == 0 { + return + } + for id, te := range e.state.Topology { + _ = id + k := disableKey(te.Local.Name, te.Remote.Name, te.Link.RemotePort) + want, pending := e.localDisabled[k] + if !pending { + continue + } + if te.Link.Disabled == want { + // The ring caught up (or agreed independently): stop tracking. + delete(e.localDisabled, k) + continue + } + // Still divergent: keep asserting our decision onto the adopted state. + te.Link.Disabled = want + te.Active = !want + } + // Forget decisions for forwards that left the topology entirely — otherwise + // a long-lived cluster would accumulate dead keys. + live := make(map[string]struct{}, len(e.state.Topology)) + for _, te := range e.state.Topology { + live[disableKey(te.Local.Name, te.Remote.Name, te.Link.RemotePort)] = struct{}{} + } + for k := range e.localDisabled { + if _, ok := live[k]; !ok { + delete(e.localDisabled, k) + } + } +} + // IsLeader reports whether this node is the current ring leader. func (e *Engine) IsLeader() bool { return e.state.LeaderID == e.ID } @@ -1022,6 +1110,14 @@ func (e *Engine) SetPeerPersist(fn func(peersJSON string) error) { e.peerPersist // HTTP handlers are overwritten by the next e.state = tk.State. func (e *Engine) SetTopologySync(fn func()) { e.topologySync = fn } +// disableKey builds the map key for a forward's disabled decision: the natural +// triple (local, remote, remotePort), which is how every other part of the code +// identifies a forward. A NUL separator keeps it unambiguous for names that +// could otherwise collide across the boundaries. +func disableKey(local, remote string, port int) string { + return local + "\x00" + remote + "\x00" + strconv.Itoa(port) +} + // persistPeers extracts all alive peers (addr + nodeKey, excluding self) // from the current ring state and persists them via the peerPersist callback. // Called on every token cycle (OnToken) and on AdoptState so a crashed node