feat: per-forward worker 模型 — 每条转发独立 frpc 配置+独立进程

核心改动(ring 协议不变,Task 本来就是 {Local,Remote,Link} 三元组):
- process.WorkerKey: worker 键改为 'local~remote~port' 三元组
- renderForward: 每条 forward 渲染只含 1 个 proxy 的独立 frpc 配置
  (替代 renderRemote 把该 remote 全部 forwards 塞进一个进程)
- ClaimFn/RevokeFn: 认领/撤销只操作这一条 forward 自己的进程
- SyncWorkers/auto-start: 只拉起 localOnly 转发的独立进程,
  集群转发由 ring 认领节点拉起 (消除三台争抢 proxy already exists)
- RemoteStatus/StartRemote/StopRemote/RestartRemote: remote 级聚合辅助,
  保持 /profiles API 形状不变, 前端零改动
- handleStatus/logs: 按 worker key 解析真实 remote 名

效果: 三台节点的 frpc 各自只注册自己拥有的 proxy, 单条转发故障不再波及兄弟
This commit is contained in:
JianFeeeee
2026-08-24 01:42:58 +08:00
parent 2292ee7f3a
commit e3a978888a
9 changed files with 340 additions and 101 deletions

View File

@ -8,6 +8,7 @@ import (
"context"
"encoding/json"
"flag"
"fmt"
"log"
"net"
"net/http"
@ -90,7 +91,7 @@ func main() {
}
return fallbackBin
},
Render: renderRemote(st),
Render: renderForward(st),
AutoRestart: func(name string) bool {
s, _ := st.Settings()
if !s.RestartOnExit {
@ -172,44 +173,32 @@ func main() {
// disabled (by RevokeFn); clear it so renderRemote renders the
// proxy back in. No-op for a fresh claim.
_ = st.SetLinkDisabled(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort, false)
// Start (or restart) the worker for the remote we just claimed.
// Restart re-renders the frpc config with all of this remote's
// links (including the just-upserted one) and respawns the process.
// We deliberately do NOT iterate every stored remote here: this
// node's store may hold remotes for cluster forwards it does NOT
// own (the submitter persists the whole canvas), and starting
// those here would duplicate workers across the cluster.
// Start (or restart) the per-forward worker for exactly this link.
// Each forward has its own frpc process (keyed by the forward
// triple); restarting only this key leaves sibling forwards'
// processes untouched.
if tk.Remote.Enabled {
if _, has := pm.Status(tk.Remote.Name); has {
_ = pm.Restart(tk.Remote.Name)
key := process.WorkerKey(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort)
if _, has := pm.Status(key); has {
_ = pm.Restart(key)
} else {
_ = pm.Start(tk.Remote.Name)
_ = pm.Start(key)
}
}
log.Printf("ring[%s] claimed task %s: %s→%s:%d", selfID, tk.ID, tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort)
return nil
},
// Revoke cancels the forward NON-destructively: flag the exact link
// disabled (so renderRemote drops just its proxy and a later forward-
// centric start can re-enable it) and right-size the worker — stop it
// when the remote has no remaining enabled forwards, else restart so
// sibling proxies survive. The old delete-local/remote/drop-all-
// matching-links form served the topology-derived-canvas model but
// broke per-forward stop/start: it wiped sibling forwards sharing the
// local or remote (group stop) and made re-start fail (link gone).
// Revoke cancels the forward NON-destructively: stop THIS forward's
// own worker process (per-forward model — no siblings share it) and
// flag the exact link disabled so a later forward-centric start can
// re-enable it. The old delete-local/remote/drop-all-matching-links
// form served the topology-derived-canvas model but broke per-forward
// 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)
fwd, _ := st.LinksForRemote(tk.Remote.Name)
enabled := 0
for _, f := range fwd {
if !f.Disabled {
enabled++
}
}
if enabled == 0 {
_ = pm.Stop(tk.Remote.Name)
} else if _, running := pm.Status(tk.Remote.Name); running {
_ = pm.Restart(tk.Remote.Name)
key := process.WorkerKey(tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort)
if _, running := pm.Status(key); running {
_ = pm.Stop(key)
}
log.Printf("ring[%s] revoked task %s: %s→%s:%d", selfID, tk.ID, tk.Local.Name, tk.Remote.Name, tk.Link.RemotePort)
return nil
@ -297,53 +286,72 @@ func main() {
return path, ver, nil
},
SyncWorkers: func() {
// Only forwards that stay on this node (localOnly) get their worker
// started here. Cluster-distributed forwards (non-localOnly) are claimed
// and spawned by the owning cluster node — starting them locally too
// would duplicate the worker.
remotes, err := st.ListRemotes()
// PER-FORWARD worker model: every enabled, non-disabled link on THIS
// node gets its own frpc process keyed by the forward triple.
// - localOnly forwards: always owned here (loopback stays local).
// - cluster forwards: only started by the ring claim (ClaimFn) —
// SyncWorkers deliberately skips them so a node that merely holds
// a canvas copy does not spawn workers for forwards it doesn't own
// (the old "every node registers every proxy" bug).
links, err := st.ListLinks()
if err != nil {
return
}
for _, r := range remotes {
if !r.Enabled {
for _, ln := range links {
loc, ok := st.GetLocal(ln.Local)
if !ok || !loc.LocalOnly {
continue // cluster-owned or unknown local
}
rem, ok := st.GetRemote(ln.Remote)
if !ok || !rem.Enabled {
continue
}
// worker starts here only if every link of this remote is localOnly
fwd, err := st.LinksForRemote(r.Name)
if err != nil {
key := process.WorkerKey(ln.Local, ln.Remote, ln.RemotePort)
if ln.Disabled {
_ = pm.Stop(key) // per-forward stop kills exactly its process
continue
}
allLocal := len(fwd) > 0
for _, f := range fwd {
loc, ok := st.GetLocal(f.Service)
if !ok || !loc.LocalOnly {
allLocal = false
break
}
}
if !allLocal {
// has at least one cluster-distributed forward; cluster owns it
continue
}
if _, has := pm.Status(r.Name); has {
_ = pm.Restart(r.Name)
if _, has := pm.Status(key); has {
_ = pm.Restart(key)
} else {
_ = pm.Start(r.Name)
_ = pm.Start(key)
}
}
},
}
// Auto-start enabled remotes on boot.
// Auto-start localOnly forwards on boot. Cluster forwards are NOT started
// here — they are re-submitted by the leader's startup auto-submit below
// and claimed (spawning their own worker) wherever the ring places them.
s, _ := st.Settings()
if s.AutoStartProfiles {
remotes, _ := st.ListRemotes()
for _, r := range remotes {
if r.Enabled {
_ = pm.Start(r.Name)
links, _ := st.ListLinks()
remotesByName := map[string]store.Remote{}
if remotes, err := st.ListRemotes(); err == nil {
for _, r := range remotes {
remotesByName[r.Name] = r
}
}
localsByName := map[string]store.Local{}
if locals, err := st.ListLocals(); err == nil {
for _, l := range locals {
localsByName[l.Name] = l
}
}
for _, ln := range links {
if ln.Disabled {
continue
}
loc, ok := localsByName[ln.Local]
if !ok || !loc.LocalOnly {
continue
}
rem, ok := remotesByName[ln.Remote]
if !ok || !rem.Enabled {
continue
}
_ = pm.Start(process.WorkerKey(ln.Local, ln.Remote, ln.RemotePort))
}
}
handler, err := httpapi.NewServeMux(h)
@ -633,35 +641,44 @@ func splitCSV(s string) []string {
return out
}
// renderRemote builds the process.Render closure from store data.
func renderRemote(st *store.Store) func(string) ([]byte, error) {
return func(remoteName string) ([]byte, error) {
// renderForward builds the process.Render closure for the PER-FORWARD worker
// model: each (local, remote, remotePort) link gets its OWN frpc config
// containing exactly ONE proxy. The worker key is the forward triple
// (process.WorkerKey), so every node only ever registers its own proxies on
// frps — this eliminates the multi-node "proxy already exists" fight of the
// old per-remote model, where every node rendered ALL forwards of a remote.
// A disabled link renders an empty proxy list: frpc exits cleanly with no
// proxies (the supervisor sees a clean exit and stops restarting it), which
// is exactly the per-forward stop semantics.
func renderForward(st *store.Store) func(string) ([]byte, error) {
return func(key string) ([]byte, error) {
localName, remoteName, port, ok := process.ParseWorkerKey(key)
if !ok {
return nil, fmt.Errorf("invalid worker key %q", key)
}
rem, ok := st.GetRemote(remoteName)
if !ok {
return nil, store.ErrNotFound
}
forwards, err := st.LinksForRemote(remoteName)
if err != nil {
return nil, err
loc, ok := st.GetLocal(localName)
if !ok {
return nil, store.ErrNotFound
}
proxies := make([]render.Proxy, 0, len(forwards))
for _, f := range forwards {
// A disabled forward is omitted from the generated frpc config so
// a per-forward stop drops just this proxy on worker restart,
// leaving sibling forwards on the same remote untouched.
if f.Disabled {
continue
}
loc, ok := st.GetLocal(f.Service)
if !ok {
continue
disabled := false
for _, ln := range mustListLinks(st) {
if ln.Local == localName && ln.Remote == remoteName && ln.RemotePort == port {
disabled = ln.Disabled
break
}
}
proxies := []render.Proxy(nil)
if !disabled {
proxies = append(proxies, render.Proxy{
Name: loc.Name,
Type: loc.Protocol,
LocalIP: loc.IP,
LocalPort: loc.Port,
RemotePort: f.RemotePort,
RemotePort: port,
UseEncryption: loc.UseEncryption,
UseCompression: loc.UseCompression,
BandwidthLimit: loc.BandwidthLimit,
@ -690,6 +707,12 @@ func renderRemote(st *store.Store) func(string) ([]byte, error) {
}
}
// mustListLinks is ListLinks with errors swallowed (best-effort reads).
func mustListLinks(st *store.Store) []store.Link {
links, _ := st.ListLinks()
return links
}
// multiFlag collects repeated string flags (-peer a -peer b ...).
type multiFlag []string