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) + } }