From 041cc04dd6f0df9a775dc303405b605b404f8158 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sat, 26 Sep 2026 10:05:55 +0800 Subject: [PATCH] =?UTF-8?q?fix(cluster):=20revoke=20=E5=9C=A8=E6=97=A0=20t?= =?UTF-8?q?opology=20=E6=9D=A1=E7=9B=AE=E6=97=B6=E9=9D=99=E9=BB=98?= =?UTF-8?q?=E5=A4=B1=E6=95=88=20+=20=E5=81=9C=E7=94=A8=E7=8A=B6=E6=80=81?= =?UTF-8?q?=E8=B7=A8=E8=8A=82=E7=82=B9=E4=B8=8D=E5=90=8C=E6=AD=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 上一提交(46e8bc3)上线后实测:`.106` 重启后,被用户停用的 `minecraft` worker **仍然被拉起**。继续挖出两处,均为静默失效型。 ## 1. Handler.Revoke 被 RemoveTopology 的返回值挡住了 if claimed.Revoke { if e.state.RemoveTopology(...) { // ← false 时整个块跳过 if e.Handler != nil { e.Handler.Revoke(...) } } } 停 worker 的唯一动作被关在「topology 里确实有条目」的条件里。而 RemoveTopology 对**不在 topology 的转发返回 false** —— 也就是说,越是 需要清理的僵尸转发(条目已丢、worker 还在),revoke 越是完全不做。 这正是 minecraft 的处境:停用请求到达时它已不在 topology ⇒ 返回 false ⇒ Revoke 不执行 ⇒ worker 永远不停 ⇒ 每 33 秒刷一次 connection refused。 修法:停 worker 与摘条目是**两件独立的事**,条目不在也照样停。 (审计日志 forward.remove 仍只在真的摘掉条目时写,这是对的。) ## 2. 停用状态是节点本地的,从不跨节点同步 `links` 表每节点各一份。实测同一时刻 `.60` 记 disabled=1、`.106` 记 disabled=0 —— 因为只有收到停用请求的那个节点写了 store,而**持有 worker 的 owner 往往是另一台机器**,它的副本仍是"启用"。Claim 只读本地 store ⇒ owner 认为该转发是启用的 ⇒ 又被拉起。 修法:task 已携带 Disabled,以它为准,RevokeFn 在本节点把它落库 (已有行改写;没有行则补一条 disabled 记录,防止日后在本节点被 claim 时复活)。 ## 测试 - TestRevokeIdempotent 原先只断言"不 panic",正好漏掉这个 bug —— 它 允许「Handler.Revoke 从不调用」通过。现补上断言:拓扑里没有条目时 **仍必须调用** Handler.Revoke。 - 验证该测试有效:把 `&& removed` gate 加回去 → 如期变红;还原后变绿。 go build / go vet / go test ./... 全绿,gofmt 干净。 --- cmd/webui4frpc/main.go | 44 ++++++++++++++++++++++++++++ internal/cluster/ring_engine.go | 28 +++++++++++------- internal/cluster/ring_revoke_test.go | 20 ++++++++++++- 3 files changed, 81 insertions(+), 11 deletions(-) diff --git a/cmd/webui4frpc/main.go b/cmd/webui4frpc/main.go index 05e13f8..2776bbb 100644 --- a/cmd/webui4frpc/main.go +++ b/cmd/webui4frpc/main.go @@ -227,6 +227,50 @@ func main() { // stop/start: it wiped sibling forwards sharing the local or remote. RevokeFn: func(ctx context.Context, tk *cluster.Task) error { _ = st.SetLinkDisabled(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort, true) + // The links table is node-local: only the node that received the + // stop request has the flag persisted, while the node that OWNS the + // forward is the one holding the worker (and is often a different + // machine). Trusting only the local row meant the owner's copy could + // still read disabled=false. The task now carries the flag, so use it + // as the authority and write it here. + ln, found, err := st.LinkByTriple(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort) + switch { + case err != nil: + return err + case found: + ln.Disabled = true + // Persist via the same rewrite path the canvas uses, so the flag + // survives a restart on THIS node too. + links, err := st.ListLinks() + if err != nil { + return err + } + for i := range links { + if links[i].Local == ln.Local && links[i].Remote == ln.Remote && links[i].RemotePort == ln.RemotePort { + links[i].Disabled = true + } + } + if err := st.ReplaceLinks(links); err != nil { + return err + } + case tk.Link.Disabled: + // No local row yet (this node never claimed the forward) but the + // task asserts the stop — record it so a later claim on this node + // cannot resurrect the forward. + loc, rem := tk.Local, tk.Remote + links, err := st.ListLinks() + if err != nil { + return err + } + links = append(links, store.Link{ + Local: loc.Name, Remote: rem.Name, RemotePort: tk.Link.RemotePort, + OffsetX: tk.Link.OffsetX, OffsetY: tk.Link.OffsetY, + Group: tk.Link.Group, Disabled: true, + }) + if err := st.ReplaceLinks(links); err != nil { + return err + } + } key := process.WorkerKey(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort) if _, running := pm.Status(key); running { _ = pm.Stop(key) diff --git a/internal/cluster/ring_engine.go b/internal/cluster/ring_engine.go index 8e9a5c2..03e1d60 100644 --- a/internal/cluster/ring_engine.go +++ b/internal/cluster/ring_engine.go @@ -390,18 +390,26 @@ func (e *Engine) runCommands(ctx context.Context, tk *Token) error { break } if claimed.Revoke { - if e.state.RemoveTopology(claimed.Local.Name, claimed.Remote.Name, claimed.Link.RemotePort) { - if e.Handler != nil { - if err := e.Handler.Revoke(ctx, claimed); err != nil { - log.Printf("ring[%s] revoke %s: %v", e.ID, claimed.ID, err) - } - } - if e.Log != nil { - _, _ = e.Log.Append(e.ID, LogForwardRemove, map[string]any{ - "taskId": claimed.ID, "local": claimed.Local.Name, "remote": claimed.Remote.Name, - }) + // The revoke handler must run REGARDLESS of whether a topology entry + // existed. RemoveTopology() returns false when the forward is not in + // the topology, and gating the handler on it made the revoke a silent + // no-op in exactly the case that matters: a forward the user stopped + // while it was already absent from the topology kept its worker + // running forever (observed live — a disabled forward kept dialing a + // local service that was intentionally down, thousands of + // connection-refused lines). Stopping the worker and retiring the + // entry are independent duties; the entry may simply not be there. + removed := e.state.RemoveTopology(claimed.Local.Name, claimed.Remote.Name, claimed.Link.RemotePort) + if e.Handler != nil { + if err := e.Handler.Revoke(ctx, claimed); err != nil { + log.Printf("ring[%s] revoke %s: %v", e.ID, claimed.ID, err) } } + if removed && e.Log != nil { + _, _ = e.Log.Append(e.ID, LogForwardRemove, map[string]any{ + "taskId": claimed.ID, "local": claimed.Local.Name, "remote": claimed.Remote.Name, + }) + } continue } if claimed.RemoveNode != "" && claimed.RemoveNode == e.ID { diff --git a/internal/cluster/ring_revoke_test.go b/internal/cluster/ring_revoke_test.go index 1467d7b..00e0e92 100644 --- a/internal/cluster/ring_revoke_test.go +++ b/internal/cluster/ring_revoke_test.go @@ -51,8 +51,21 @@ func TestRevokeTaskRemovesTopology(t *testing.T) { } // TestRevokeIdempotent: revoking an already-missing forward does not error. +// +// It must ALSO still call Handler.Revoke. The handler is what actually stops +// the per-forward frpc worker; the topology entry is only bookkeeping. Gating +// the handler on RemoveTopology()'s return value (as this test used to permit) +// turned every revoke of an already-absent forward into a silent no-op — the +// worker kept running, which is how a forward the user had stopped kept +// dialling a local service that was intentionally down, forever. func TestRevokeIdempotent(t *testing.T) { - eng := newTestEngine("n1", true) + revoked := 0 + eng := NewEngine("n1", "n1:7500", "u", "p", "0.71.0", nil, + &fakeHandler{load: Load{MemPct: 5, NetPct: 5}, + revoke: func(ctx context.Context, tk *Task) error { revoked++; return nil }}, + func(ctx context.Context, next string, tk *Token) error { return nil }, + "n1:7500", true, "") + eng.state.AddRevoke(store.Local{Name: "ghost"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 1}) if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state}); err != nil { t.Fatalf("revoke missing: %v", err) @@ -61,4 +74,9 @@ func TestRevokeIdempotent(t *testing.T) { if len(eng.state.TopologyList()) != 0 { t.Fatal("should be empty") } + // The stop side-effect must have happened even though there was no entry. + if revoked != 1 { + t.Fatalf("Handler.Revoke called %d times, want 1 — a revoke with no topology entry "+ + "must still stop the worker, otherwise stopped forwards keep running", revoked) + } }