diff --git a/cmd/webui4frpc/main.go b/cmd/webui4frpc/main.go index 2776bbb..da733e1 100644 --- a/cmd/webui4frpc/main.go +++ b/cmd/webui4frpc/main.go @@ -199,11 +199,13 @@ func main() { // triple); restarting only this key leaves sibling forwards' // processes untouched. if disabled { - // A stopped forward must also not linger in the topology: leaving - // the entry behind is what let the reconcile loop above keep - // re-claiming it on every restart. + // Mark the entry stopped rather than deleting it. The entry is the + // carrier that keeps the flag travelling around the ring, and + // UpdateTopologyDisabled also clears Active — which is what makes + // this entry inert for OfflineReassign() and the startup + // reconcile, so it is not resurrected later. if ring != nil { - ring.RemoveTopologyEntry(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort) + ring.UpdateTopologyDisabled(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort, true) } log.Printf("ring[%s] claim %s skipped: %s→%s:%d is disabled", selfID, tk.ID, tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort) return nil @@ -275,17 +277,55 @@ func main() { if _, running := pm.Status(key); running { _ = pm.Stop(key) } - // Drop the topology entry too. Leaving it behind meant the entry - // outlived the worker, and the next startup reconcile saw a - // "missing" worker for an owned forward and re-claimed it — which - // restarted a forward the user had explicitly stopped. Revoking - // must be a complete retirement, not just a stop. + // Mark the entry stopped instead of deleting it: the entry carries + // the flag around the ring, and clearing Active keeps it inert for + // the rebalancing paths (OfflineReassign skips inactive entries, and + // the startup reconcile skips disabled forwards), so the worker is + // not re-spawned on the next restart. Deleting it here is what used + // to leave the stop with no carrier at all. if ring != nil { - ring.RemoveTopologyEntry(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort) + ring.UpdateTopologyDisabled(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort, true) } log.Printf("ring[%s] revoked task %s: %s→%s:%d", selfID, tk.ID, tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort) return nil }, + // RestartFn brings an already-claimed forward back after a stop. It + // must NOT run the claim bookkeeping (ClaimFn), which appends links and + // writes topology state: the forward is already established, and a stop + // deliberately keeps its topology entry, so all that is needed is to + // clear the flag and respawn the worker. The engine has already + // re-marked the topology entry enabled before calling this. + RestartFn: func(ctx context.Context, tk *cluster.Task) error { + // Clear the stop on this node's own store, using the same path the + // start handler uses, so the local copy cannot disagree with the + // ring and block a later restart. + if ln, found, err := st.LinkByTriple(tk.Link.Local, tk.Link.Remote, tk.Link.RemotePort); err != nil { + return err + } else if found && ln.Disabled { + 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 = false + } + } + if err := st.ReplaceLinks(links); err != nil { + return err + } + } + if tk.Remote.Enabled { + key := process.WorkerKey(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort) + if _, has := pm.Status(key); has { + _ = pm.Restart(key) + } else { + _ = pm.Start(key) + } + } + log.Printf("ring[%s] restarted %s: %s→%s:%d", selfID, tk.ID, tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort) + return nil + }, }, ringTransport.SendTo(func(nodeID string) string { if ring == nil { @@ -310,10 +350,21 @@ func main() { // left does NOT auto-rejoin. ring.SetPeerPersist(func(peersJSON string) error { return st.SetClusterPeers(peersJSON) }) - // Re-apply local store group labels onto the ring topology after each - // state adoption. Without this, group changes made via HTTP handlers + // Re-apply local store state onto the ring topology after each state + // adoption, and learn peers' disabled flags back into the local store. + // + // Without the first 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. + // + // 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. ring.SetTopologySync(func() { links, err := st.ListLinks() if err != nil { @@ -322,6 +373,19 @@ func main() { for _, ln := range links { ring.UpdateTopologyGroup(ln.Local, ln.Remote, ln.RemotePort, ln.Group) } + for _, te := range ring.State().TopologyList() { + disabled, known := ring.TopologyDisabled(te.Local.Name, te.Remote.Name, te.Link.RemotePort) + 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 + } + _ = st.ReconcileLinkDisabled(te.Local.Name, te.Remote.Name, te.Link.RemotePort, disabled) + } }) h := &httpapi.Handler{ diff --git a/internal/cluster/ring.go b/internal/cluster/ring.go index 0d1cedf..b4bbf20 100644 --- a/internal/cluster/ring.go +++ b/internal/cluster/ring.go @@ -62,6 +62,12 @@ type Task struct { // owning node cancels it (stop worker, drop from topology). Reuses the // same publish channel as creation (round-1 inject, round-2 apply). Revoke bool `json:"revoke,omitempty"` + // Restart marks a RE-ENABLE task: the forward is already claimed and its + // topology entry still exists (a stop keeps the entry as the flag's + // carrier), so the creation path would dedupe the submission and the + // duplicate-claim guard would discard it. The owner applies this one + // unconditionally: re-mark the entry enabled and (re)spawn the worker. + Restart bool `json:"restart,omitempty"` // RemoveNode: node ID to remove from the ring; the node self-removes when // the command reaches it via the token. RemoveNode string `json:"removeNode,omitempty"` @@ -333,10 +339,15 @@ func (s *State) AddTopology(t *Task, ownerID string) *TopoEntry { if s.Topology == nil { s.Topology = map[string]*TopoEntry{} } + // Active mirrors the task's disabled flag rather than being unconditionally + // true. A claim can legitimately carry a stopped forward (e.g. a node + // re-claiming from a departed peer), and forcing Active=true there would + // resurrect it: OfflineReassign only re-queues Active entries, and every + // reader that filters on Active would treat it as running again. e := &TopoEntry{ TaskID: t.ID, OwnerID: ownerID, Local: t.Local, Remote: t.Remote, Link: t.Link, - Active: true, + Active: !t.Link.Disabled, } s.Topology[t.ID] = e return e @@ -408,6 +419,39 @@ func (s *State) ForwardsOwnedBy(owner string) []*TopoEntry { // (local, remote, port) triple. Returns true if found. The updated entry // propagates to all nodes via the next token cycle — group changes sync // through the ring without a dedicated command. +// UpdateTopologyDisabled sets the disabled flag on the topology entry matching +// the forward's natural key. Mirror of UpdateTopologyGroup: it makes a locally +// decided stop/start part of the topology so the next token cycle carries it to +// every other member (the alternative — a one-shot revoke task addressed at the +// owner — only converges if that owner happens to be online right then). +// +// Returns false when the forward is not in the topology, which is not an error: +// a stopped forward may legitimately have no entry yet. +func (s *State) UpdateTopologyDisabled(local, remote string, port int, disabled bool) bool { + for _, e := range s.Topology { + if e.Local.Name == local && e.Remote.Name == remote && e.Link.RemotePort == port { + e.Link.Disabled = disabled + // Keep Active consistent with the flag so readers that key off it + // (status page, load accounting) agree with Link.Disabled. + e.Active = !disabled + return true + } + } + return false +} + +// TopologyDisabled returns the disabled flag the cluster currently holds for a +// forward, and whether such an entry exists. Used by the adoption path to learn +// a peer's decision into the local store. +func (s *State) TopologyDisabled(local, remote string, port int) (bool, bool) { + for _, e := range s.Topology { + if e.Local.Name == local && e.Remote.Name == remote && e.Link.RemotePort == port { + return e.Link.Disabled, true + } + } + return false, false +} + func (s *State) UpdateTopologyGroup(local, remote string, port int, group string) bool { for _, e := range s.Topology { if e.Local.Name == local && e.Remote.Name == remote && e.Link.RemotePort == port { diff --git a/internal/cluster/ring_engine.go b/internal/cluster/ring_engine.go index 03e1d60..c10a9b8 100644 --- a/internal/cluster/ring_engine.go +++ b/internal/cluster/ring_engine.go @@ -18,6 +18,11 @@ import ( type Handler interface { Claim(ctx context.Context, tk *Task) error Revoke(ctx context.Context, tk *Task) error + // Restart re-enables an already-claimed forward and brings its worker back. + // Needed because a stop keeps the topology entry, which makes both the + // creation path (idempotency guard) and the duplicate-claim guard refuse to + // re-create an existing forward. + Restart(ctx context.Context, tk *Task) error RuntimeLoad() Load } @@ -390,23 +395,34 @@ func (e *Engine) runCommands(ctx context.Context, tk *Token) error { break } if claimed.Revoke { - // 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) + // A user-requested stop is recorded as a FLAG on the topology entry, + // not as a removal. The entry is the carrier that makes the decision + // travel: TopoEntry.Link is a full store.Link and State rides the + // token every cycle, so keeping the entry is what lets a stop reach + // every member — including a node that was offline when the stop was + // issued. Removing the entry instead left the flag nowhere to live, + // which is why a stop used to converge only if the owner happened to + // be online to receive a one-shot revoke task. + // + // Marking Active=false is also what makes the entry inert for the + // rebalancing paths: OfflineReassign() skips inactive entries, and + // the startup reconcile skips them too, so a stopped forward is + // neither re-queued on a node departure nor re-claimed on a restart. + // + // Handler.Revoke still runs unconditionally — it is what actually + // stops the per-forward worker, and it must run even when there was no + // entry to mark (gating it on RemoveTopology()'s result made stops of + // already-absent forwards silently do nothing, so their workers ran + // forever: observed live as thousands of connection-refused lines + // against a local service that was intentionally down). + marked := e.state.UpdateTopologyDisabled(claimed.Local.Name, claimed.Remote.Name, claimed.Link.RemotePort, true) 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{ + if marked && e.Log != nil { + _, _ = e.Log.Append(e.ID, LogForwardStop, map[string]any{ "taskId": claimed.ID, "local": claimed.Local.Name, "remote": claimed.Remote.Name, }) } @@ -475,10 +491,38 @@ func (e *Engine) runCommands(ctx context.Context, tk *Token) error { // different node. Safe for OfflineReassign: that path deletes the // topology entry BEFORE re-queueing, so TopologyOwner returns "" and // the legitimate re-claim passes through. - if owner := e.state.TopologyOwner(claimed); owner != "" { - log.Printf("ring[%s] drop duplicate task %s: %s→%s:%d already owned by %s", - e.ID, claimed.ID, claimed.Local.Name, claimed.Remote.Name, - claimed.Link.RemotePort, owner) + // + // A RESTART task is the deliberate exception this guard must not eat: it + // exists precisely to re-activate a forward whose entry is still present + // (a stop keeps the entry as the flag's carrier), so applying the guard + // would swallow every re-enable and leave the forward stopped forever. + if !claimed.Restart { + if owner := e.state.TopologyOwner(claimed); owner != "" { + log.Printf("ring[%s] drop duplicate task %s: %s→%s:%d already owned by %s", + e.ID, claimed.ID, claimed.Local.Name, claimed.Remote.Name, + claimed.Link.RemotePort, owner) + continue + } + } + if claimed.Restart { + // Re-enable in place and hand the worker back to this node (the + // entry's owner). This must NOT create a second topology entry, which + // is why it skips the claim bookkeeping below. + e.state.UpdateTopologyDisabled(claimed.Local.Name, claimed.Remote.Name, claimed.Link.RemotePort, false) + if e.Handler != nil { + if err := e.Handler.Restart(ctx, claimed); err != nil { + e.state.PendingTasks[claimed.ID] = claimed + log.Printf("ring[%s] restart %s failed: %v", e.ID, claimed.ID, err) + break + } + } + if e.Log != nil { + _, _ = e.Log.Append(e.ID, LogForwardStart, map[string]any{ + "taskId": claimed.ID, "local": claimed.Local.Name, "remote": claimed.Remote.Name, + }) + } + log.Printf("ring[%s] restarted %s: %s→%s:%d", e.ID, claimed.ID, + claimed.Local.Name, claimed.Remote.Name, claimed.Link.RemotePort) continue } if e.Handler != nil { @@ -817,14 +861,87 @@ func (e *Engine) RevokeTask(local store.Local, remote store.Remote, link store.L } func (e *Engine) SubmitTask(local store.Local, remote store.Remote, link store.Link) *Task { - if e.HasTask(local.Name, remote.Name, link.RemotePort) { + // Guard on an ACTIVE task only. A topology entry that is merely marked + // disabled is kept around as the carrier for the stop flag, so counting it + // here would make this a permanent no-op and a stopped forward could never + // be started again. + if e.HasActiveTask(local.Name, remote.Name, link.RemotePort) { return nil } return e.state.AddPending(local, remote, link) } -// HasTask reports whether a forward with the same local/remote/remotePort is -// already pending or active in the topology (idempotency guard for resaves). +// HasTopologyEntry reports whether the forward has a topology entry (whether or +// not it is marked disabled). A stop keeps the entry as the flag's carrier, so +// "has an entry" is what distinguishes "already claimed by someone" from "never +// submitted" — the distinction startForward needs. +func (e *Engine) HasTopologyEntry(local, remote string, port int) bool { + _, found := e.state.TopologyDisabled(local, remote, port) + return found +} + +// SubmitRestart publishes a task asking the current owner of an already-claimed +// forward to bring its worker back up. +// +// Neither of the existing channels can express "restart what already exists": +// SubmitTask is guarded by HasActiveTask and the entry is still present, so it +// dedupes; and the claim path drops any task for a forward that already has an +// owner (its defence against duplicate-claim collisions). Re-enabling a stopped +// forward therefore needs its own task kind, which the owner applies without +// re-running the claim bookkeeping. +func (e *Engine) SubmitRestart(local store.Local, remote store.Remote, link store.Link) *Task { + // The task must not carry the stop: the whole point is to re-enable. + link.Disabled = false + t := &Task{ + ID: e.state.NextTaskID(), + Local: local, + Remote: remote, + Link: link, + Created: time.Now().Unix(), + Restart: true, + } + if e.state.PendingTasks == nil { + e.state.PendingTasks = map[string]*Task{} + } + e.state.PendingTasks[t.ID] = t + return t +} + +// HasActiveTask reports whether a forward with this key is genuinely in flight +// or running: a non-revocation pending task, or a topology entry that is not +// marked disabled. +// +// It is deliberately narrower than HasTask. Under mark-don't-remove a stopped +// forward keeps its topology entry (that entry is what carries the flag between +// nodes), so "an entry exists" no longer means "this forward is active". +// SubmitTask must use THIS predicate, otherwise starting a stopped forward — +// startForward() publishes the enable and then calls SubmitTask — would be +// swallowed by its own idempotency guard. +func (e *Engine) HasActiveTask(local, remote string, port int) bool { + for _, t := range e.state.PendingList() { + if t.Revoke { + continue // a revocation is the opposite of an active forward + } + if t.Local.Name == local && t.Remote.Name == remote && t.Link.RemotePort == port { + return true + } + } + for _, t := range e.state.TopologyList() { + if t.Link.Disabled { + continue // stopped; the entry is only a flag carrier + } + if t.Local.Name == local && t.Remote.Name == remote && t.Link.RemotePort == port { + return true + } + } + return false +} + +// HasTask reports whether ANY record for this forward exists — a pending task +// (including an in-flight revocation) or a topology entry (including one marked +// disabled). This is the "have we already handled this key" question, used to +// avoid re-issuing work on repeated canvas saves. For "is it actually running" +// use HasActiveTask instead. func (e *Engine) HasTask(local, remote string, port int) bool { for _, t := range e.state.PendingList() { if t.Local.Name == local && t.Remote.Name == remote && t.Link.RemotePort == port { @@ -846,6 +963,19 @@ func (e *Engine) UpdateTopologyGroup(local, remote string, port int, group strin return e.state.UpdateTopologyGroup(local, remote, port, group) } +// 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. +func (e *Engine) UpdateTopologyDisabled(local, remote string, port int, disabled bool) bool { + 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). +func (e *Engine) TopologyDisabled(local, remote string, port int) (bool, bool) { + return e.state.TopologyDisabled(local, remote, port) +} + // IsLeader reports whether this node is the current ring leader. func (e *Engine) IsLeader() bool { return e.state.LeaderID == e.ID } diff --git a/internal/cluster/ring_engine_test.go b/internal/cluster/ring_engine_test.go index 06649bf..89675ac 100644 --- a/internal/cluster/ring_engine_test.go +++ b/internal/cluster/ring_engine_test.go @@ -8,9 +8,10 @@ import ( ) type fakeHandler struct { - load Load - claim func(ctx context.Context, tk *Task) error - revoke func(ctx context.Context, tk *Task) error + load Load + claim func(ctx context.Context, tk *Task) error + revoke func(ctx context.Context, tk *Task) error + restart func(ctx context.Context, tk *Task) error } func (h *fakeHandler) Revoke(ctx context.Context, tk *Task) error { @@ -20,6 +21,13 @@ func (h *fakeHandler) Revoke(ctx context.Context, tk *Task) error { return nil } +func (h *fakeHandler) Restart(ctx context.Context, tk *Task) error { + if h.restart != nil { + return h.restart(ctx, tk) + } + return nil +} + func (h *fakeHandler) RuntimeLoad() Load { return h.load } func (h *fakeHandler) Claim(ctx context.Context, tk *Task) error { if h.claim != nil { diff --git a/internal/cluster/ring_handler.go b/internal/cluster/ring_handler.go index 7f5327d..83c92b4 100644 --- a/internal/cluster/ring_handler.go +++ b/internal/cluster/ring_handler.go @@ -17,6 +17,10 @@ type AppHandler struct { ClaimFn func(ctx context.Context, tk *Task) error // RevokeFn cancels the forward on this node (stop worker + drop from store). RevokeFn func(ctx context.Context, tk *Task) error + // RestartFn re-enables an already-claimed forward on this node (clear the + // stopped flag + bring the worker back). Required because a stop keeps the + // topology entry, so the normal claim path will not re-create it. + RestartFn func(ctx context.Context, tk *Task) error } // RuntimeLoad implements Handler. @@ -44,6 +48,14 @@ func (a *AppHandler) Revoke(ctx context.Context, tk *Task) error { return a.RevokeFn(ctx, tk) } +// Restart implements Handler. +func (a *AppHandler) Restart(ctx context.Context, tk *Task) error { + if a.RestartFn == nil { + return nil + } + return a.RestartFn(ctx, tk) +} + // SampleMemLoad returns a cheap memory-usage percentage (0..100). func SampleMemLoad() float64 { var m runtime.MemStats diff --git a/internal/cluster/ring_log.go b/internal/cluster/ring_log.go index dd6e04c..1c37ebf 100644 --- a/internal/cluster/ring_log.go +++ b/internal/cluster/ring_log.go @@ -20,6 +20,8 @@ import ( // LogKind enumerates operation-log entry types. const ( LogForwardAdd = "forward.add" + LogForwardStop = "forward.stop" + LogForwardStart = "forward.start" LogForwardRemove = "forward.remove" LogNodeJoin = "node.join" LogNodeLeave = "node.leave" @@ -140,7 +142,7 @@ func DetailOf(e LogEntry) string { } str := func(key string) string { s, _ := d[key].(string); return s } switch e.Kind { - case LogForwardAdd, LogForwardRemove: + case LogForwardAdd, LogForwardStop, LogForwardStart, LogForwardRemove: s := str("local") + " → " + str("remote") if id := str("taskId"); id != "" { if len(id) > 8 { diff --git a/internal/cluster/ring_revoke_test.go b/internal/cluster/ring_revoke_test.go index 00e0e92..6550084 100644 --- a/internal/cluster/ring_revoke_test.go +++ b/internal/cluster/ring_revoke_test.go @@ -2,14 +2,22 @@ package cluster import ( "context" + "encoding/json" "testing" "webui4frpc/internal/store" ) -// TestRevokeTaskRemovesTopology: a revoke task published to the ring removes -// the forward from topology, calls Handler.Revoke, and logs forward.remove. -func TestRevokeTaskRemovesTopology(t *testing.T) { +// TestRevokeTaskMarksTopologyDisabled: a stop delivered through the ring marks +// the topology entry disabled instead of deleting it, calls Handler.Revoke, and +// logs forward.stop. +// +// The entry is deliberately KEPT: TopoEntry.Link is the carrier that makes the +// stop travel, and State rides the token every cycle. Keeping it is what lets a +// stop reach a member that was offline when the stop was issued — deleting the +// entry left the flag nowhere to live, so a stop converged only if the owning +// node happened to be online for a one-shot revoke task. +func TestRevokeTaskMarksTopologyDisabled(t *testing.T) { revoked := false eng := NewEngine("n1", "n1:7500", "u", "p", "0.71.0", nil, &fakeHandler{load: Load{MemPct: 5, NetPct: 5}, @@ -26,27 +34,73 @@ func TestRevokeTaskRemovesTopology(t *testing.T) { if len(eng.state.TopologyList()) != 1 { t.Fatalf("topology after claim = %+v", eng.state.TopologyList()) } + if d, _ := eng.state.TopologyDisabled("web", "frps1", 18081); d { + t.Fatal("a freshly claimed forward must not start out disabled") + } // publish a revoke task pointing at the same forward eng.state.AddRevoke(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: eng.state}); err != nil { t.Fatal(err) } - if len(eng.state.TopologyList()) != 0 { - t.Fatalf("topology after revoke = %+v", eng.state.TopologyList()) + + // The entry must SURVIVE, marked disabled, so the flag keeps travelling. + if len(eng.state.TopologyList()) != 1 { + t.Fatalf("topology entry was removed; the disabled flag would have no carrier: %+v", + eng.state.TopologyList()) + } + disabled, known := eng.state.TopologyDisabled("web", "frps1", 18081) + if !known { + t.Fatal("topology entry missing after revoke") + } + if !disabled { + t.Fatal("revoke did not mark the topology entry disabled") + } + // Active must follow the flag, since OfflineReassign only re-queues Active + // entries — a stopped forward must not be resurrected by a node departure. + if eng.state.TopologyList()[0].Active { + t.Fatal("a stopped forward must not remain Active (OfflineReassign would re-queue it)") } if !revoked { t.Fatal("Handler.Revoke was not called") } - // log should contain forward.remove - var sawRemove bool + // log should contain forward.stop (not forward.remove — nothing was removed) + var sawStop bool for _, e := range eng.Log.Snapshot() { - if e.Kind == LogForwardRemove { - sawRemove = true + if e.Kind == LogForwardStop { + sawStop = true } } - if !sawRemove { - t.Fatalf("log missing forward.remove: %+v", eng.Log.Snapshot()) + if !sawStop { + t.Fatalf("log missing forward.stop: %+v", eng.Log.Snapshot()) + } +} + +// TestStoppedTopologyEntrySurvivesAdoption is the property that makes the whole +// mark-don't-remove design work: because the entry persists, a node that adopts +// ring state learns the stop. This is the offline-member case a one-shot revoke +// task could never cover. +func TestStoppedTopologyEntrySurvivesAdoption(t *testing.T) { + src := newTestEngine("n1", true) + src.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) + if _, err := src.OnToken(context.Background(), &Token{Cycle: 1, State: src.state}); err != nil { + t.Fatal(err) + } + src.state.AddRevoke(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) + if _, err := src.OnToken(context.Background(), &Token{Cycle: 2, State: src.state}); err != nil { + t.Fatal(err) + } + + // A different node adopts the ring state (what e.state = tk.State does). + peer := newTestEngine("n2", false) + peer.AdoptState(src.state) + + disabled, known := peer.TopologyDisabled("web", "frps1", 18081) + if !known { + t.Fatal("peer did not receive the topology entry for the stopped forward") + } + if !disabled { + t.Fatal("peer adopted the entry but no longer sees it as disabled — the stop did not propagate") } } @@ -80,3 +134,213 @@ func TestRevokeIdempotent(t *testing.T) { "must still stop the worker, otherwise stopped forwards keep running", revoked) } } + +// TestStoppedForwardNotRequeuedOnNodeDeparture pins why marking (rather than +// deleting) is safe: OfflineReassign only re-queues ACTIVE entries, so a +// stopped forward is not silently handed to another node when its owner goes +// away. Deleting the entry, or leaving it Active, would both resurrect it. +func TestStoppedForwardNotRequeuedOnNodeDeparture(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}); err != nil { + t.Fatal(err) + } + eng.state.AddRevoke(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: eng.state}); err != nil { + t.Fatal(err) + } + + owner := eng.state.TopologyList()[0].OwnerID + eng.state.OfflineReassign(owner) + + if n := len(eng.state.PendingList()); n != 0 { + t.Fatalf("a stopped forward was re-queued on the owner's departure (pending=%d): %+v", + n, eng.state.PendingList()) + } +} + +// TestSubmitTaskNotBlockedByStoppedEntry is the regression test for the bug the +// mark-don't-remove change introduced: SubmitTask's idempotency guard used +// HasTask, which matches a disabled topology entry. Since a stop now KEEPS the +// entry, any submission for that forward while it is still stopped is swallowed +// by its own guard — so a stopped forward can never be started again. +// +// The probe deliberately submits WITHOUT re-enabling first. Re-enabling would +// clear the flag and make HasTask and HasActiveTask agree, hiding the defect; +// the guard has to be exercised at the moment the entry is still stopped. +func TestSubmitTaskNotBlockedByStoppedEntry(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}); err != nil { + t.Fatal(err) + } + + // stop it -> entry is kept, marked disabled + eng.state.AddRevoke(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: eng.state}); err != nil { + t.Fatal(err) + } + if d, _ := eng.state.TopologyDisabled("web", "frps1", 18081); !d { + t.Fatal("precondition: forward should be disabled") + } + // The stopped entry IS findable by the broad predicate — which is exactly why + // the narrow one has to exist. + if !eng.HasTask("web", "frps1", 18081) { + t.Fatal("precondition: the stopped entry should still be findable by HasTask") + } + if eng.HasActiveTask("web", "frps1", 18081) { + t.Fatal("precondition: a stopped forward must NOT count as active") + } + + // The regression: a submission must be published for a stopped forward. + // With a HasTask guard this returns nil and the forward is stuck forever. + tk := eng.SubmitTask(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, + store.Link{RemotePort: 18081}) + if tk == nil { + t.Fatal("SubmitTask was swallowed by the stopped topology entry — " + + "a stopped forward could never be started again") + } + if tk.Link.RemotePort != 18081 { + t.Fatalf("unexpected task: %+v", tk) + } +} + +// TestSubmitTaskStillDedupesActiveForward is the other half of the guard: the +// narrower predicate must not turn SubmitTask into a duplicate-task generator. +func TestSubmitTaskStillDedupesActiveForward(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}); err != nil { + t.Fatal(err) + } + if !eng.HasActiveTask("web", "frps1", 18081) { + t.Fatal("precondition: a claimed forward must count as active") + } + if tk := eng.SubmitTask(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, + store.Link{RemotePort: 18081}); tk != nil { + t.Fatalf("SubmitTask must stay idempotent for an active forward, got %+v", tk) + } +} + +// TestStoppedForwardNotRequeuedOnNodeDeparture pins why marking (rather than +// deleting) is safe: OfflineReassign only re-queues ACTIVE entries, so a +// stopped forward is not silently handed to another node when its owner goes + +// TestAddTopologyRespectsDisabledFlag: a claim that carries a stopped forward +// must not mark the new entry Active, or it would come back to life. +func TestAddTopologyRespectsDisabledFlag(t *testing.T) { + s := &State{Topology: map[string]*TopoEntry{}} + e := s.AddTopology(&Task{ + ID: "t1", Local: store.Local{Name: "web"}, Remote: store.Remote{Name: "frps1"}, + Link: store.Link{RemotePort: 18081, Disabled: true}, + }, "n1") + if e.Active { + t.Fatal("AddTopology forced Active=true for a disabled claim — the forward would resurrect") + } + if !e.Link.Disabled { + t.Fatal("AddTopology dropped the disabled flag from the link") + } +} + +// TestRestartTaskBypassesDuplicateClaimGuard is the engine half of the +// re-enable path: a restart task for a forward that still HAS a topology entry +// must reach Handler.Restart. +// +// This is precisely what the duplicate-claim guard refuses — it exists to stop +// a stale task from spawning a second worker for an owned forward, and a restart +// is indistinguishable from that unless it is an explicit task kind. Marking +// the entry stopped (instead of deleting it) is what made this necessary: the +// entry the guard keys on now outlives a stop. +func TestRestartTaskBypassesDuplicateClaimGuard(t *testing.T) { + restarts := 0 + claims := 0 + eng := NewEngine("n1", "n1:7500", "u", "p", "0.71.0", nil, + &fakeHandler{load: Load{MemPct: 5, NetPct: 5}, + claim: func(ctx context.Context, tk *Task) error { claims++; return nil }, + restart: func(ctx context.Context, tk *Task) error { restarts++; return nil }}, + func(ctx context.Context, next string, tk *Token) error { return nil }, + "n1:7500", true, "") + + // Establish + stop: entry survives, marked disabled. + 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}); err != nil { + t.Fatal(err) + } + eng.state.AddRevoke(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081}) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: eng.state}); err != nil { + t.Fatal(err) + } + if d, _ := eng.state.TopologyDisabled("web", "frps1", 18081); !d { + t.Fatal("precondition: should be disabled") + } + entryCount := len(eng.state.TopologyList()) + + // Start: publish a restart and run a cycle. + eng.state.UpdateTopologyDisabled("web", "frps1", 18081, false) + eng.SubmitRestart(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, + store.Link{RemotePort: 18081}) + if _, err := eng.OnToken(context.Background(), &Token{Cycle: 3, State: eng.state}); err != nil { + t.Fatal(err) + } + + if restarts != 1 { + t.Fatalf("Handler.Restart called %d times, want 1 — the restart task was swallowed "+ + "by the duplicate-claim guard, so the forward would never come back", restarts) + } + if claims != 1 { + t.Fatalf("Handler.Claim called %d times, want 1 (the original claim only); "+ + "a restart must not re-run the claim bookkeeping", claims) + } + // Re-enabling must not create a second entry. + if n := len(eng.state.TopologyList()); n != entryCount { + t.Fatalf("restart created a duplicate topology entry: %d -> %d", entryCount, n) + } + if d, _ := eng.state.TopologyDisabled("web", "frps1", 18081); d { + t.Fatal("entry is still disabled after the restart task was applied") + } + var sawStart bool + for _, e := range eng.Log.Snapshot() { + if e.Kind == LogForwardStart { + sawStart = true + } + } + if !sawStart { + t.Fatalf("log missing forward.start: %+v", eng.Log.Snapshot()) + } +} + +// TestRestartFlagSurvivesTokenSerialization: the restart task travels to the +// owner inside the token, so the flag must round-trip through JSON. Without the +// struct tag it would silently deserialize as false and the owner would treat +// it as an ordinary claim — which the duplicate guard then discards. +func TestRestartFlagSurvivesTokenSerialization(t *testing.T) { + eng := newTestEngine("n1", true) + eng.SubmitRestart(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, + store.Link{RemotePort: 18081}) + + blob, err := json.Marshal(&Token{Cycle: 7, State: eng.state}) + if err != nil { + t.Fatal(err) + } + var back Token + if err := json.Unmarshal(blob, &back); err != nil { + t.Fatal(err) + } + var found *Task + for _, tk := range back.State.PendingList() { + if tk.Local.Name == "web" && tk.Link.RemotePort == 18081 { + found = tk + } + } + if found == nil { + t.Fatal("restart task did not survive token serialization") + } + if !found.Restart { + t.Fatal("the restart flag was lost in JSON — the owner would see a plain claim " + + "and the duplicate guard would discard it") + } + // Revoke and Restart must stay distinguishable. + if found.Revoke { + t.Fatal("a restart task must not also read as a revocation") + } +} diff --git a/internal/httpapi/forward_revoke_test.go b/internal/httpapi/forward_revoke_test.go index da5c018..ccb61f6 100644 --- a/internal/httpapi/forward_revoke_test.go +++ b/internal/httpapi/forward_revoke_test.go @@ -117,3 +117,88 @@ func TestStopForwardPublishesDisabledFlagInRevokeTask(t *testing.T) { "this stop was deliberate, which is the bug that let stopped forwards resurrect") } } + +// TestStopThenStartPublishesRestartTask guards the failure mode that +// mark-don't-remove introduces, and that only shows up on the WIRE. +// +// A stop keeps the topology entry (it is the carrier for the flag), so on +// re-enable: +// - SubmitTask dedupes, because the entry exists; +// - the claim path discards any task for a forward that already has an owner. +// +// Both channels therefore refuse, no task is published, the owner never learns +// about the re-enable, and the forward stays stopped forever — a forward the +// user can stop but never restart. +// +// Asserting the task is PUBLISHED is the point: asserting only that the flag +// cleared would pass while nothing reached the owner. (The earlier version of +// this test did exactly that and stayed green against the broken code.) +func TestStopThenStartPublishesRestartTask(t *testing.T) { + _, ring, ts := newRingTestHandler(t) + + saveCanvas(t, ts, `{ + "locals": [{"name":"svc","ip":"127.0.0.1","port":59999,"protocol":"tcp"}], + "remotes": [{"name":"srv-a","ip":"1.2.3.4","port":7000,"enabled":true}], + "links": [{"local":"svc","remote":"srv-a","remotePort":45999}] + }`) + // Claim it so a topology entry exists (that is what makes the two normal + // channels refuse later). + if _, err := ring.OnToken(context.Background(), &cluster.Token{Cycle: 1, State: *ring.State()}); err != nil { + t.Fatal(err) + } + if !ring.HasActiveTask("svc", "srv-a", 45999) { + t.Fatalf("precondition: forward should be active, topo=%+v", ring.State().TopologyList()) + } + + forwards := func(action string) { + t.Helper() + b, _ := json.Marshal(stopForwardReq{"svc", "srv-a", 45999}) + req, _ := http.NewRequest(http.MethodPost, ts.URL+"/api/manager/forwards/"+action, bytes.NewReader(b)) + req.SetBasicAuth("admin", "pw") + resp, err := ts.Client().Do(req) + if err != nil { + t.Fatal(err) + } + resp.Body.Close() + if resp.StatusCode != http.StatusOK { + t.Fatalf("%s status = %d", action, resp.StatusCode) + } + } + + // --- stop: entry kept, marked disabled ------------------------------- + forwards("stop") + if d, known := ring.TopologyDisabled("svc", "srv-a", 45999); !known || !d { + t.Fatalf("stop did not mark the topology entry disabled (known=%v disabled=%v)", known, d) + } + if ring.HasActiveTask("svc", "srv-a", 45999) { + t.Fatal("a stopped forward must not count as active") + } + + // --- start: a task MUST reach the owner ------------------------------ + forwards("start") + if d, _ := ring.TopologyDisabled("svc", "srv-a", 45999); d { + t.Fatal("start did not clear the disabled flag") + } + if !ring.HasActiveTask("svc", "srv-a", 45999) { + t.Fatal("after start the forward must be active again") + } + + // The actual regression: something has to be queued for the owner. The + // re-enable must be published as a restart task, since a plain creation task + // would be deduped or discarded. + var restart *cluster.Task + for _, tk := range ring.State().PendingList() { + if tk.Restart && tk.Local.Name == "svc" && tk.Link.RemotePort == 45999 { + restart = tk + break + } + } + if restart == nil { + t.Fatalf("start published no restart task — the owner would never bring the "+ + "worker back, leaving the forward permanently stopped (pending=%+v)", + ring.State().PendingList()) + } + if restart.Link.Disabled { + t.Fatal("the restart task carries disabled=true, so the owner would re-apply the stop") + } +} diff --git a/internal/httpapi/handlers.go b/internal/httpapi/handlers.go index 2f81854..663e4b9 100644 --- a/internal/httpapi/handlers.go +++ b/internal/httpapi/handlers.go @@ -245,7 +245,12 @@ func (h *Handler) applyCanvas(w http.ResponseWriter, r *http.Request, canvas *ca } if ln.Disabled { // Stopped on the forwards page: make sure it leaves the topology. - if h.Ring.HasTask(ln.Local, ln.Remote, ln.RemotePort) { + // Guard on an ACTIVE task — an entry that is already marked + // disabled has nothing left to revoke, and re-issuing a revocation + // for it on every canvas save would be pure noise. (HasTask, which + // also matches disabled entries, is the right predicate for the + // "already handled" question this branch is not asking.) + if h.Ring.HasActiveTask(ln.Local, ln.Remote, ln.RemotePort) { h.Ring.RevokeTask(loc, rem, ln) } continue diff --git a/internal/httpapi/handlers_forwards.go b/internal/httpapi/handlers_forwards.go index c2ee0e2..be73fbd 100644 --- a/internal/httpapi/handlers_forwards.go +++ b/internal/httpapi/handlers_forwards.go @@ -84,6 +84,25 @@ func (h *Handler) startForward(local, remote string, port int) error { } return nil } + if h.Ring != nil { + // Publish the enable into the topology so it rides the ring (a peer that + // still holds the stale "stopped" copy learns about it on adoption). + h.Ring.UpdateTopologyDisabled(local, remote, port, false) + } + // A cluster forward whose topology entry still exists needs its OWNER to + // bring the worker back, and neither of the normal channels can do it: + // - SubmitTask is idempotency-guarded, and the entry still exists (a stop + // keeps it as the flag's carrier), so the submission is deduped away; + // - the claim path drops any task for a forward that already has an + // owner, as a defence against duplicate-claim collisions. + // So a re-enable of an existing entry is published as a dedicated RESTART + // task, which the owner applies unconditionally. Without it a forward the + // user can stop but not restart — which is what marking-instead-of-removing + // would otherwise have produced. + if h.Ring != nil && h.Ring.HasTopologyEntry(local, remote, port) { + h.Ring.SubmitRestart(loc, rem, ln) + return nil + } if h.Ring != nil { h.Ring.SubmitTask(loc, rem, ln) } @@ -119,6 +138,14 @@ func (h *Handler) stopForward(local, remote string, port int) error { // SetLinkDisabled(true) above, so the link published into the token still // carried disabled=false and got copied into the topology entry verbatim. ln.Disabled = true + // Publish the stop into the topology BEFORE revoking: the revoke retires the + // own-side worker/entry, so the flag must already exist somewhere that + // survives it and travels the ring. This is the send half of cluster-wide + // disabled propagation (see store.ReconcileLinkDisabled for the receive + // half). + if h.Ring != nil { + h.Ring.UpdateTopologyDisabled(local, remote, port, true) + } if loc.LocalOnly { if h.Process != nil { key := process.WorkerKey(local, remote, port) diff --git a/internal/store/link_test.go b/internal/store/link_test.go index f57bf08..982a2d6 100644 --- a/internal/store/link_test.go +++ b/internal/store/link_test.go @@ -196,3 +196,93 @@ func TestGetLinkByIDIsStaleAfterReplaceLinks(t *testing.T) { t.Fatalf("natural key must stay reliable, got %+v found=%v", got, found) } } + +// TestReconcileLinkDisabledLearnsPeerDecision is the store half of +// cluster-wide stop propagation. A node that did NOT serve the stop request has +// no reason to know about it, and its links table is node-local — so the flag +// arrives via the ring and lands here. The case that matters is the OWNER of a +// forward on a different machine: before this existed, that node's copy still +// read "enabled", so it kept (or re-spawned) the worker for a forward the user +// had explicitly stopped. +func TestReconcileLinkDisabledLearnsPeerDecision(t *testing.T) { + st, err := New(filepath.Join(t.TempDir(), "test.db")) + if err != nil { + t.Fatal(err) + } + defer st.Close() + seed(t, st, []string{"mc"}, "srv") + if err := st.ReplaceLinks([]Link{{Local: "mc", Remote: "srv", RemotePort: 25565, Group: "game"}}); err != nil { + t.Fatal(err) + } + + // A peer's stop arrives. + if err := st.ReconcileLinkDisabled("mc", "srv", 25565, true); err != nil { + t.Fatal(err) + } + ln, found, err := st.LinkByTriple("mc", "srv", 25565) + if err != nil { + t.Fatal(err) + } + if !found { + t.Fatal("link disappeared during reconcile") + } + if !ln.Disabled { + t.Fatal("the peer's stop did not land in the local store") + } + if ln.Group != "game" { + t.Fatalf("reconcile must not clobber other fields, group=%q", ln.Group) + } + + // A peer's re-enable arrives. + if err := st.ReconcileLinkDisabled("mc", "srv", 25565, false); err != nil { + t.Fatal(err) + } + if ln, _, _ := st.LinkByTriple("mc", "srv", 25565); ln.Disabled { + t.Fatal("the peer's re-enable did not land") + } +} + +// TestReconcileLinkDisabledCreatesPlaceholderForUnknownForward: a node that has +// never seen the forward still must remember that it is stopped, otherwise a +// later claim on that node would resurrect it. +func TestReconcileLinkDisabledCreatesPlaceholderForUnknownForward(t *testing.T) { + st, err := New(filepath.Join(t.TempDir(), "test.db")) + if err != nil { + t.Fatal(err) + } + defer st.Close() + seed(t, st, []string{"ghost"}, "srv") + + if err := st.ReconcileLinkDisabled("ghost", "srv", 9999, true); err != nil { + t.Fatal(err) + } + ln, found, err := st.LinkByTriple("ghost", "srv", 9999) + if err != nil { + t.Fatal(err) + } + if !found { + t.Fatal("a stopped-but-unknown forward must be remembered, or a later claim resurrects it") + } + if !ln.Disabled { + t.Fatal("placeholder is not marked disabled") + } +} + +// TestReconcileLinkDisabledIgnoresEnableForUnknown: an enable for a forward this +// node has never seen must NOT create a row. Creating one would invent forwards +// out of ring state. +func TestReconcileLinkDisabledIgnoresEnableForUnknown(t *testing.T) { + st, err := New(filepath.Join(t.TempDir(), "test.db")) + if err != nil { + t.Fatal(err) + } + defer st.Close() + seed(t, st, []string{"other"}, "srv") + + if err := st.ReconcileLinkDisabled("unknown", "srv", 1234, false); err != nil { + t.Fatal(err) + } + if _, found, _ := st.LinkByTriple("unknown", "srv", 1234); found { + t.Fatal("an enable for an unknown forward must not materialise a row") + } +} diff --git a/internal/store/store.go b/internal/store/store.go index 39633da..799db75 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -663,6 +663,11 @@ func (s *Store) ReplaceLinks(links []Link) error { // SetLinkDisabled flips the disabled flag of a forward identified by its // (local, remote, remotePort) natural key. This is the persistence half of the // forwards-page start/stop toggle; the caller also drives the worker/ring side. +// +// Kept as a targeted UPDATE rather than a ReplaceLinks rewrite on purpose: +// ReplaceLinks deletes and reinserts every row, handing out fresh autoincrement +// ids and invalidating any Link a caller captured earlier (they travel inside +// ring tokens). Flipping one flag must not perturb other rows' identity. func (s *Store) SetLinkDisabled(local, remote string, port int, disabled bool) error { _, err := s.db.Exec( "UPDATE links SET disabled = ? WHERE local = ? AND remote = ? AND remote_port = ?", @@ -671,6 +676,52 @@ func (s *Store) SetLinkDisabled(local, remote string, port int, disabled bool) e return err } +// ReconcileLinkDisabled applies a cluster-wide view of one forward's disabled +// flag into the local store, creating a placeholder row when this node has none +// yet. +// +// This is the receive half of disabled-flag propagation. The stop decision is +// made on whichever node served the request, then rides the token ring in the +// topology entry; every other member calls this on adoption so its own +// links table agrees. Without it the flag lived only on the node that handled +// the request, and the node actually OWNS the forward — usually a different +// machine — still believed the forward was enabled and re-spawned its worker. +// +// A placeholder row is deliberate: a node that has never seen the forward still +// needs to remember "this is stopped" so a later claim on this node cannot +// resurrect it. The placeholder carries the same natural key, so a subsequent +// real claim fills in the rest. +func (s *Store) ReconcileLinkDisabled(local, remote string, port int, disabled bool) error { + cur, found, err := s.LinkByTriple(local, remote, port) + if err != nil { + return err + } + if found { + if cur.Disabled == disabled { + return nil // already agrees; avoid needless writes every token cycle + } + return s.SetLinkDisabled(local, remote, port, disabled) + } + if !disabled { + // Nothing to remember: an unknown forward with no entry is simply + // "not stopped", which is the default the claim path already assumes. + return nil + } + // Need a placeholder, which requires the local/remote foreign keys to exist. + if _, ok := s.GetLocal(local); !ok { + return nil // cannot materialise a link without its local peer row + } + if _, ok := s.GetRemote(remote); !ok { + return nil + } + links, err := s.ListLinks() + if err != nil { + return err + } + links = append(links, Link{Local: local, Remote: remote, RemotePort: port, Disabled: true}) + return s.ReplaceLinks(links) +} + // SetLinkGroup assigns a management group label to a forward identified by its // (local, remote, remotePort) natural key. Empty string clears the group // (moves the forward to 未分组). This is the persistence half of the