mirror of
https://gitcode.com/JianFeeeee/webui4frpc.git
synced 2026-10-02 15:14:01 +00:00
fix(cluster): revoke 在无 topology 条目时静默失效 + 停用状态跨节点不同步
上一提交(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 干净。
This commit is contained in:
@ -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)
|
||||
|
||||
@ -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 {
|
||||
|
||||
@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user