Files
webui4frpc/internal/httpapi/handlers_forwards.go
JianFeeeee 1c835425de feat(cluster): 停用改为「标记」语义,让 disabled 真正随令牌环跨节点传播
承接用户提问「设计上停用不是本来就会跨节点传输吗」——核实结论:结构上确实
如此(TopoEntry.Link 是完整 store.Link,整个 State 随 token 每轮广播),但
实际路径断了。断点正是「撤销会删掉 topology 条目」:条目是 flag 的载体,
删了就无处传播,于是停用只能靠一次性 revoke 任务投递给 owner,**owner 当时
不在线就收不到**(实测 .60 记 disabled=1 / .106 记 0,就是这么来的)。

## 改为标记而非移除

撤销不再 RemoveTopology,而是 UpdateTopologyDisabled(true),条目保留、
Link.Disabled=true、Active=false。Active 正是为此存在:OfflineReassign()
只处理 Active 条目,所以停用的转发在 owner 掉线时不会被重新排队。

- 新增 UpdateTopologyDisabled / TopologyDisabled(照 UpdateTopologyGroup 的桥)
- 新增 store.ReconcileLinkDisabled 作接收端:adoption 时把环上的 flag 落进
  本地 store;本节点没有该转发时补一条 disabled 占位行(否则日后在本节点被
  claim 会复活),enable 则不建行
- SetTopologySync 由单向(store→环)扩为双向:群组仍上行,disabled 下行
- AddTopology 的 Active 跟随 Link.Disabled(原本硬编码 true,认领一个停用
  转发就会复活它)
- 审计日志细分 forward.stop / forward.start,与 forward.remove 区分

## 语义变更带出的两个新问题(都已修)

1. **「启动」这条路断了**。条目保留 ⇒ SubmitTask 被去重挡下,而认领路径的
   duplicate-claim 防御又会丢弃「已有 owner」的任务 ⇒ 重启任务发不出去,owner
   永远收不到,转发**能停不能起**。
   修:新增 Task.Restart 这一独立任务类型 + SubmitRestart + Handler.RestartFn,
   显式绕过 duplicate-claim 防御并原地复活(不重复建条目、不重跑 claim 簿记)。
   SubmitTask 的守卫同时从 HasTask 收窄为新的 HasActiveTask(跳过 disabled 条目
   与撤销任务);saveCanvas 的判断相应改用 HasActiveTask,避免每次保存都对
   已标记的转发重复发撤销。

2. 原本两处 RemoveTopologyEntry 调用(ClaimFn/RevokeFn 的 disabled 分支)在
   新语义下会把本该保留的条目删掉,改为 UpdateTopologyDisabled。

## 测试(每个都做了「回退修复行→必须变红→还原变绿」双向验证)

- TestStoppedTopologyEntrySurvivesAdoption —— 离线成员也能学到停用,
  一次性 revoke 任务永远做不到这一点
- TestStoppedForwardNotRequeuedOnNodeDeparture / TestAddTopologyRespectsDisabledFlag
  —— 标记而非删除为何安全
- TestSubmitTaskNotBlockedByStoppedEntry / TestSubmitTaskStillDedupesActiveForward
- TestRestartTaskBypassesDuplicateClaimGuard / TestRestartFlagSurvivesTokenSerialization
- TestStopThenStartPublishesRestartTask(HTTP 端到端,断言**任务真的发出**)
- TestReconcileLinkDisabled*(store 侧三条)

★ 两次踩到**假绿**:第一版只断言 store 层(newTestHandler 的 Ring 为 nil,
坏掉的路根本没执行);第二版在 re-enable **之后**才调 SubmitTask,此时新旧
谓词结果相同,测不出差异。都是靠「回退修复行看是否变红」抓出来的 —— 这个
双向验证已经是本项目的固定动作。

go build / go vet / go test ./... 全绿,gofmt 干净。
2026-09-26 10:44:04 +08:00

347 lines
12 KiB
Go

package httpapi
import (
"encoding/json"
"net/http"
"webui4frpc/internal/process"
"webui4frpc/internal/store"
)
// forwardsReq is the {local,remote,remotePort} natural key identifying a single
// forward. The triple is stable across canvas re-saves (unlike Link.ID, which
// ReplaceLinks wholesale-replaces) and matches HasTask's keying.
type forwardsReq struct {
Local string `json:"local"`
Remote string `json:"remote"`
RemotePort int `json:"remotePort"`
}
// groupReq selects a group for one-click start/stop on the forwards page.
type groupReq struct {
Group string `json:"group"`
}
// findLinkByTriple returns the link matching the (local,remote,remotePort)
// natural key, or ok=false.
func findLinkByTriple(s *store.Store, local, remote string, port int) (store.Link, bool) {
links, err := s.ListLinks()
if err != nil {
return store.Link{}, false
}
for _, l := range links {
if l.Local == local && l.Remote == remote && l.RemotePort == port {
return l, true
}
}
return store.Link{}, false
}
// findLinkFromTopology looks up a forward's local/remote/link data from the
// ring topology when it's not in the local SQLite store (forward owned by
// another node). Returns the store.Local, store.Remote, store.Link, and ok.
func (h *Handler) findLinkFromTopology(local, remote string, port int) (store.Local, store.Remote, store.Link, bool) {
if h.Ring == nil {
return store.Local{}, store.Remote{}, store.Link{}, false
}
snap := h.Ring.Snapshot()
for _, t := range snap.Topology {
if t.Local.Name == local && t.Remote.Name == remote && t.Link.RemotePort == port {
return t.Local, t.Remote, t.Link, true
}
}
return store.Local{}, store.Remote{}, store.Link{}, false
}
// startForward flips a forward to enabled and brings it up. A local-only
// forward restarts/starts its local frpc worker so the re-enabled proxy is
// rendered back in; a cluster (non-localOnly) forward is submitted to the ring
// so the lowest-load member claims and spawns it (plan §令牌环协议). SubmitTask
// is idempotent via HasTask. Falls back to ring topology when the forward
// exists only in the cluster (owned by another node, not in local store).
func (h *Handler) startForward(local, remote string, port int) error {
ln, ok := findLinkByTriple(h.Store, local, remote, port)
loc, lok := h.Store.GetLocal(local)
rem, rok := h.Store.GetRemote(remote)
if !ok || !lok || !rok {
// Fall back to ring topology for forwards owned by other nodes.
tLoc, tRem, tLn, tok := h.findLinkFromTopology(local, remote, port)
if !tok {
return store.ErrNotFound
}
ln, loc, rem = tLn, tLoc, tRem
}
_ = h.Store.SetLinkDisabled(local, remote, port, false)
ln.Disabled = false
if loc.LocalOnly {
if h.Process != nil {
key := process.WorkerKey(local, remote, port)
if _, has := h.Process.Status(key); has {
_ = h.Process.Restart(key)
} else {
_ = h.Process.Start(key)
}
}
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)
}
return nil
}
// stopForward flips a forward to disabled and tears it down. A local-only
// forward restarts its worker so renderRemote omits the proxy (sibling forwards
// on the same remote keep running — true per-forward stop); a cluster forward
// is revoked so the owning node cancels its worker and drops it from the
// topology (plan §任务撤销). RevokeTask is idempotent. Falls back to ring
// topology when the forward exists only in the cluster (owned by another node).
func (h *Handler) stopForward(local, remote string, port int) error {
ln, ok := findLinkByTriple(h.Store, local, remote, port)
loc, lok := h.Store.GetLocal(local)
rem, rok := h.Store.GetRemote(remote)
if !ok || !lok || !rok {
tLoc, tRem, tLn, tok := h.findLinkFromTopology(local, remote, port)
if !tok {
return store.ErrNotFound
}
ln, loc, rem = tLn, tLoc, tRem
} else {
_ = h.Store.SetLinkDisabled(local, remote, port, true)
}
// The revoke task travels to whichever node OWNS the forward, and that node
// re-reads the disabled flag from its own store before starting a worker —
// so the flag has to be set on every node that has a copy of this link, not
// just the one handling this request. Propagating Disabled on the task lets
// the owner's RevokeFn stop the worker even if its own store row is stale.
//
// This also fixes a latent inconsistency: `ln` was read BEFORE the
// 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)
if _, has := h.Process.Status(key); has {
_ = h.Process.Restart(key)
}
}
return nil
}
if h.Ring != nil {
h.Ring.RevokeTask(loc, rem, ln)
}
return nil
}
// handleForwardsStart toggles one forward on.
func (h *Handler) handleForwardsStart(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
methodNotAllowed(w)
return
}
var req forwardsReq
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "parse json: "+err.Error(), http.StatusBadRequest)
return
}
if req.Local == "" || req.Remote == "" || req.RemotePort <= 0 {
http.Error(w, "local/remote/remotePort required", http.StatusBadRequest)
return
}
if err := h.startForward(req.Local, req.Remote, req.RemotePort); err != nil {
http.Error(w, err.Error(), http.StatusNotFound)
return
}
writeJSON(w, http.StatusOK, map[string]any{"ok": true})
}
// handleForwardsStop toggles one forward off.
func (h *Handler) handleForwardsStop(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
methodNotAllowed(w)
return
}
var req forwardsReq
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "parse json: "+err.Error(), http.StatusBadRequest)
return
}
if req.Local == "" || req.Remote == "" || req.RemotePort <= 0 {
http.Error(w, "local/remote/remotePort required", http.StatusBadRequest)
return
}
if err := h.stopForward(req.Local, req.Remote, req.RemotePort); err != nil {
http.Error(w, err.Error(), http.StatusNotFound)
return
}
writeJSON(w, http.StatusOK, map[string]any{"ok": true})
}
// handleForwardsGroupStart starts every forward in a group with one click.
func (h *Handler) handleForwardsGroupStart(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
methodNotAllowed(w)
return
}
var req groupReq
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "parse json: "+err.Error(), http.StatusBadRequest)
return
}
links, err := h.Store.ListLinks()
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
for _, ln := range links {
if ln.Group != req.Group {
continue
}
_ = h.startForward(ln.Local, ln.Remote, ln.RemotePort)
}
writeJSON(w, http.StatusOK, map[string]any{"ok": true})
}
// handleForwardsGroupStop stops every forward in a group with one click.
func (h *Handler) handleForwardsGroupStop(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
methodNotAllowed(w)
return
}
var req groupReq
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "parse json: "+err.Error(), http.StatusBadRequest)
return
}
links, err := h.Store.ListLinks()
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
for _, ln := range links {
if ln.Group != req.Group {
continue
}
_ = h.stopForward(ln.Local, ln.Remote, ln.RemotePort)
}
writeJSON(w, http.StatusOK, map[string]any{"ok": true})
}
// assignReq carries the group label to assign to a single forward (status
// page group chip edit). Empty Group clears the assignment (移出分组).
type assignReq struct {
Local string `json:"local"`
Remote string `json:"remote"`
RemotePort int `json:"remotePort"`
Group string `json:"group"`
}
// handleForwardsAssign changes the group label of a single forward. The
// status page group chip is the quick entry; the canvas port editor also
// carries a group field (both persist via the same store column). The local
// DB row is updated if present; the ring topology entry is always updated so
// the group label propagates to all nodes via the next token cycle.
func (h *Handler) handleForwardsAssign(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
methodNotAllowed(w)
return
}
var req assignReq
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "parse json: "+err.Error(), http.StatusBadRequest)
return
}
if req.Local == "" || req.Remote == "" || req.RemotePort <= 0 {
http.Error(w, "local/remote/remotePort required", http.StatusBadRequest)
return
}
found := false
if _, ok := findLinkByTriple(h.Store, req.Local, req.Remote, req.RemotePort); ok {
found = true
if err := h.Store.SetLinkGroup(req.Local, req.Remote, req.RemotePort, req.Group); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
if h.Ring != nil {
if h.Ring.UpdateTopologyGroup(req.Local, req.Remote, req.RemotePort, req.Group) {
found = true
}
}
if !found {
http.Error(w, store.ErrNotFound.Error(), http.StatusNotFound)
return
}
writeJSON(w, http.StatusOK, map[string]any{"ok": true})
}
// handleForwardsGroupDelete dissolves a group: every forward in the named
// group is moved to 未分组 (grp=""). The group itself is not a stored entity
// — it exists only as a label on links — so clearing all members is the
// complete "delete". Clears both local store links and ring topology entries
// so the change syncs to all nodes.
func (h *Handler) handleForwardsGroupDelete(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
methodNotAllowed(w)
return
}
var req groupReq
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "parse json: "+err.Error(), http.StatusBadRequest)
return
}
if req.Group == "" {
http.Error(w, "group required", http.StatusBadRequest)
return
}
links, err := h.Store.ListLinks()
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
for _, ln := range links {
if ln.Group != req.Group {
continue
}
_ = h.Store.SetLinkGroup(ln.Local, ln.Remote, ln.RemotePort, "")
}
// Also clear group on ring topology entries (forwards owned by other nodes).
if h.Ring != nil {
snap := h.Ring.Snapshot()
for _, t := range snap.Topology {
if t.Link.Group == req.Group {
h.Ring.UpdateTopologyGroup(t.Local.Name, t.Remote.Name, t.Link.RemotePort, "")
}
}
}
writeJSON(w, http.StatusOK, map[string]any{"ok": true})
}