From 918ca5d5ba5e7362df0471570f06da4a12cddf4d Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sat, 26 Sep 2026 11:44:28 +0800 Subject: [PATCH] =?UTF-8?q?fix(cluster):=20=E5=81=9C=E7=94=A8=E7=8A=B6?= =?UTF-8?q?=E6=80=81=E4=BB=A5=E3=80=8C=E7=8E=AF=E3=80=8D=E4=B8=BA=E6=9D=83?= =?UTF-8?q?=E5=A8=81=EF=BC=8C=E6=B6=88=E9=99=A4=E9=95=BF=E6=9C=9F=E7=A6=BB?= =?UTF-8?q?=E7=BA=BF=E5=AF=BC=E8=87=B4=E7=9A=84=E6=B0=B8=E4=B9=85=E5=88=86?= =?UTF-8?q?=E6=AD=A7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 承接用户指出的遗留:links 表是每节点本地副本,靠「store→环上行 + 环→store 下行」双向 sync 收敛,但两边都可能被覆盖。其中「本地 store 权威」这条规则有 硬伤,本次改掉。 ## 先复现,再动手 写探针验证「入环 token 会不会冲掉本地未传播的决定」,结果比预期严重: after adopting a stale token: known=true disabled=false LOCAL DISABLE WAS WIPED by an incoming token OnToken/AdoptState 是 `e.state = tk.State` **整体替换**。两次 token 之间做出的 停用决定,只要下一轮到达的 token 是「决定之前」捕获的,就会被整个洗掉 —— 决定永远传不出去,转发照旧运行。Group 之所以看起来没这个问题,是因为它每次 adoption 都被 `SetTopologySync` 从 store 重新推上去,而 disabled 没有对应的 「尚未传播」保护。 ## 改法:环权威 + 本地决定带「未确认」标记 1. **环权威**:adoption 时把环上的 disabled 写穿本地 store(write-through, 不是监听器),任何 peer 的决定都在一个 token 周期内落地。这终结了旧规则 「各人信自己那份」造成的永久分歧。顺序上只有单向要求:引擎先重新断言本地 未确认决定,再跑 host sync,所以读环绝不会覆盖用户刚做的操作。 2. **未确认决定受保护**(RingEngine.localDisabled):`UpdateTopologyDisabled` 记下决定,adoption 后由 `reconcileLocalDisabled` 重新断言到刚采纳的 state 上, 于是它会随下一轮 token 传出去。环报回同值时删除该键(全cluster已一致); 转发从 topology 消失时一并清扫,map 不会无限增长。**false 同样受保护** —— 重新启用也需要传播,丢掉它会把转发永久留在停用态。 3. `TopologyDisabled()` 对未确认决定短路返回本地值,避免状态页与用户刚做的 操作相反。 ## 测试(又抓到一个假绿) 新增 4 条:停用/启用跨 adoption 存活、确认后停止断言(让位给 peer 的后续决定)、 追踪表不累积。 ★ `TestLocalDisableSurvivesAdoption` **第一版是假绿**:它断言 `TopologyDisabled()`,而该访问器会短路到本地决定,于是「即便即将转发的 state 仍是 enabled」它也报 true。改成断言**下游节点会看到什么**(用一个 peer 引擎 AdoptState 本引擎的 state)后,去掉重新断言如期变红: the forwarded token carries disabled=false, want true 双向验证通过,这个教训要记:测「声明」而不是测「实际传播的状态」,等于没测。 go build / go vet / go test ./... 全绿,gofmt 干净。 --- cmd/webui4frpc/main.go | 30 ++--- internal/cluster/ring_disabled_test.go | 157 +++++++++++++++++++++++++ internal/cluster/ring_engine.go | 96 +++++++++++++++ 3 files changed, 268 insertions(+), 15 deletions(-) create mode 100644 internal/cluster/ring_disabled_test.go 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