feat: cluster reliability (leader failover, crash rejoin, key exchange) + auth/users + canvas/forwards enhancements + comprehensive README + API docs

- Cluster: forwardToNext offline detection (leader+non-leader), WatchLeader 1s heartbeat fallback, 409 for standalone nodes, Node.NodeKey key exchange via token ring, ClusterPeers persistence + auto-rejoin, Forward delegates to forwardToNext (bugfix)
- Auth: Basic Auth (flag-creds fast path) + bcrypt users (admin/viewer) + Bearer API keys (read/write/admin scope)
- Frontend: UsersView (accounts+API keys), ClusterView (ring/nodeKey/tasks/topology/log), StatusView (group management, per-proxy status), CanvasView (edge toggle/group), PortEdge (disabled/group labels)
- API: handlers split (canvas/forwards/users/logs), canvas export/import, forwards group start/stop/assign/delete, cluster endpoints
- Docs: comprehensive README rewrite (all flags/APIs/auth/cluster), docs/cluster-api.md (cluster management API reference)
- Deploy: run-cluster.sh now 4-node ring + 1 isolated standalone, test-forward.sh updated for 4 nodes
- Removed plan.md (design notes consolidated into README + API docs)
This commit is contained in:
2026-08-19 21:09:24 +08:00
parent b518a13446
commit eda9bb9597
55 changed files with 5462 additions and 793 deletions

View File

@ -30,15 +30,24 @@ type Node struct {
Version string `json:"version,omitempty"`
Cache []string `json:"cache,omitempty"`
LastSeen int64 `json:"lastSeen,omitempty"`
// NodeKey is this member's cluster admission key. It travels in the
// token so every node knows every peer's key — a crashed node can
// rejoin via ANY cached peer by presenting that peer's key. Without
// this, a rejoiner would only know its original sponsor's key (from
// the -join-key flag) and couldn't rejoin through a different peer.
NodeKey string `json:"nodeKey,omitempty"`
}
// Load is the combined load metric used to pick the task claimer.
type Load struct {
MemPct float64 `json:"memPct"`
NetPct float64 `json:"netPct"`
MemPct float64 `json:"memPct"`
NetPct float64 `json:"netPct"`
Forwards int `json:"forwards,omitempty"` // active forwards owned by this node (primary signal)
}
func (l Load) Score() float64 { return l.MemPct + l.NetPct }
// Score weights owned forwards heavily so the node with the fewest active
// forwards is picked first; mem/net only break ties at equal forward count.
func (l Load) Score() float64 { return float64(l.Forwards)*100 + l.MemPct + l.NetPct }
// Task is a PENDING forward request circulated in the token. It carries the
// intermediate forwarding intent (local/remote/link) — NOT a rendered frpc
@ -84,7 +93,6 @@ type State struct {
PendingTasks map[string]*Task `json:"pendingTasks,omitempty"`
// Topology: active forwards owned by members (full cluster view).
Topology map[string]*TopoEntry `json:"topology,omitempty"`
Seq int64 `json:"seq"`
}
// Token is the circulating message: one physical token per cycle (single
@ -207,8 +215,42 @@ func (s *State) LowestAlive() *Node {
}
func (s *State) NextTaskID() string {
s.Seq++
return fmt.Sprintf("t%d", s.Seq)
// Collision-free id allocation by scanning the ids actually in flight
// (pending + topology) and taking max+1. The ring is mutually exclusive
// (one token holder at a time), so when a node mints an id it has the
// authoritative full view — scanning existing ids guarantees a fresh id.
// n is tiny (handful of forwards). No Seq counter is needed: a cross-node
// Seq was previously adopted wholesale on every OnToken (e.state = tk.State),
// which dropped local increments and could regress below an id still in use,
// recycling it and overwriting an active topology entry.
max := int64(0)
for id := range s.PendingTasks {
if n := taskIDNum(id); n > max {
max = n
}
}
for _, e := range s.Topology {
if n := taskIDNum(e.TaskID); n > max {
max = n
}
}
return fmt.Sprintf("t%d", max+1)
}
// taskIDNum extracts the numeric suffix of a task id "t12" -> 12 (0 if it does
// not parse). Used only to keep NextTaskID collision-free.
func taskIDNum(id string) int64 {
if len(id) < 2 || id[0] != 't' {
return 0
}
var n int64
for _, c := range id[1:] {
if c < '0' || c > '9' {
return 0
}
n = n*10 + int64(c-'0')
}
return n
}
// AddRemoveNode publishes a node-removal command via the token; the target

View File

@ -4,6 +4,7 @@ package cluster
import (
"context"
"encoding/json"
"fmt"
"log"
"sync"
@ -30,6 +31,20 @@ type Engine struct {
Cache []string
Handler Handler
// nodeKey is this node's cluster admission key. A newcomer must present
// the sponsor's nodeKey (as JoinInfo.JoinKey) to join via it. Persisted
// in store.Settings; stable across restarts so -join-key stays valid.
// Generated lazily: on CreateCluster (creator) or AdoptState (joiner),
// NOT at startup — a fresh node that hasn't created/joined has no key.
// Cleared on detachAsStandalone (leaving the cluster invalidates the key).
nodeKey string
keyPersist func(key string) error // persists nodeKey to store (nil in tests)
// peerPersist saves the cached peer list (JSON of [{addr,key},...]) so a
// crashed node can auto-rejoin on restart via any cached peer. Called on
// every token cycle (OnToken) and on AdoptState. Cleared (pass "") on
// detachAsStandalone — an explicit leave must NOT auto-rejoin.
peerPersist func(peersJSON string) error
state State
// myAddr maps our Node ID to the address peers dial.
myAddr string
@ -81,10 +96,11 @@ type Engine struct {
// NewEngine builds the engine; state holds this node as initial leader unless
// a peer list says otherwise (creation node starts the ring).
func NewEngine(id, addr, user, pass, version string, cache []string, h Handler, send func(ctx context.Context, next string, tk *Token) error, selfAddr string, isLeader bool) *Engine {
func NewEngine(id, addr, user, pass, version string, cache []string, h Handler, send func(ctx context.Context, next string, tk *Token) error, selfAddr string, isLeader bool, nodeKey string) *Engine {
e := &Engine{
ID: id, Addr: addr, User: user, Pass: pass,
Version: version, Cache: cache, Handler: h,
nodeKey: nodeKey,
state: State{
LeaderID: "",
Cycle: 0,
@ -99,7 +115,7 @@ func NewEngine(id, addr, user, pass, version string, cache []string, h Handler,
published: map[string]struct{}{},
}
n := Node{ID: id, Addr: selfAddr, Alive: true, IsLeader: isLeader,
Load: Load{MemPct: 10, NetPct: 10}, Version: version, Cache: cache}
Load: Load{MemPct: 10, NetPct: 10}, Version: version, Cache: cache, NodeKey: nodeKey}
e.state.UpsertNode(n)
return e
}
@ -121,12 +137,21 @@ func (e *Engine) myNode() Node {
return e.state.Nodes[i]
}
// loadSnapshot reads our runtime load (mem+net) from the handler.
// loadSnapshot reads our runtime load. The primary signal is the count of
// forwards this node currently owns (Forwards) — this is what makes the
// lowest-load claim actually distribute tasks across nodes instead of the
// leader hogging every claimable task in a single token pass (its stored
// load was never refreshed between claims, so it stayed "lowest"). mem/net
// from the handler only break ties at equal forward count.
func (e *Engine) loadSnapshot() Load {
l := Load{Forwards: len(e.state.ForwardsOwnedBy(e.ID))}
if e.Handler != nil {
return e.Handler.RuntimeLoad()
b := e.Handler.RuntimeLoad()
l.MemPct, l.NetPct = b.MemPct, b.NetPct
} else {
l.MemPct, l.NetPct = 20, 20
}
return Load{MemPct: 20, NetPct: 20}
return l
}
// OnToken is the SINGLE-ROUND token handler. Per the authoritative design
@ -153,7 +178,15 @@ func (e *Engine) OnToken(ctx context.Context, tk *Token) (*Token, error) {
log.Printf("ring[%s] OnToken cycle=%d", e.ID, tk.Cycle)
// Parallel rhythm timer: operations run while the pace clock ticks.
rhythm := time.NewTimer(ringHopDelay)
// Delay scales with alive node count (more nodes → lower per-hop delay,
// keeping the round time ~constant for real-time sync).
alive := 0
for i := range tk.State.Nodes {
if tk.State.Nodes[i].Alive {
alive++
}
}
rhythm := time.NewTimer(hopDelayFor(alive))
defer rhythm.Stop()
// (a) ADOPT the cluster picture. Incoming state is authoritative: joins
@ -219,7 +252,7 @@ func (e *Engine) OnToken(ctx context.Context, tk *Token) (*Token, error) {
ID: e.ID, Addr: e.myAddr, Alive: true,
IsLeader: e.state.LeaderID == e.ID,
Load: e.loadSnapshot(),
Version: e.Version, Cache: e.Cache,
Version: e.Version, Cache: e.Cache, NodeKey: e.nodeKey,
})
}
@ -243,6 +276,8 @@ func (e *Engine) OnToken(ctx context.Context, tk *Token) (*Token, error) {
// Publish our consolidated state back into the token.
tk.State = e.state
e.lastSyncAt = time.Now().Unix()
// Persist the cached peer list so a crash/restart can auto-rejoin.
e.persistPeers()
// Record the tasks leaving on this token so their absence from the next
// incoming token is recognized as "consumed downstream" rather than
// "never sent" — otherwise the localPending re-merge above would
@ -352,16 +387,33 @@ func (e *Engine) runCommands(ctx context.Context, tk *Token) error {
}
}
}
wasLeader := e.state.LeaderID == e.ID
if succ, ok := e.state.AliveSuccessor(e.ID); ok && succ != e.ID {
e.removedNext = succ
}
e.state.SelfRemove(e.ID)
e.selfRemoved = true
// Leader failover: if we were the leader, the ring would run
// leaderless after our departure — LeaderID would be "" (SelfRemove
// clears it), no node would call Send (cycle never advances), and
// WatchLeader can't find AlivePredecessor("") to promote a
// successor. Designate the captured successor as the new leader
// so the token carries a valid LeaderID downstream; the successor
// then calls Send on its turn and the ring keeps cycling.
if wasLeader && e.removedNext != "" {
e.state.LeaderID = e.removedNext
for i := range e.state.Nodes {
e.state.Nodes[i].IsLeader = e.state.Nodes[i].ID == e.removedNext
}
if e.Log != nil {
_, _ = e.Log.Append(e.ID, LogLeaderChange, map[string]string{"leader": e.removedNext})
}
}
if e.Log != nil {
_, _ = e.Log.Append(e.ID, LogNodeLeave, map[string]string{"node": e.ID})
}
log.Printf("ring[%s] self-removed from cluster (token command %s); next=%s",
e.ID, claimed.ID, e.removedNext)
log.Printf("ring[%s] self-removed from cluster (token command %s); next=%s leader=%s",
e.ID, claimed.ID, e.removedNext, e.state.LeaderID)
continue
}
if claimed.RemoveNode != "" {
@ -376,6 +428,20 @@ func (e *Engine) runCommands(ctx context.Context, tk *Token) error {
e.ID, claimed.ID, claimed.RemoveNode)
continue
}
// Defense-in-depth against duplicate claims: if an active topology
// entry for this forward already exists (owned by us or another
// node), this task is a stale resurrected copy or a multi-token
// collision — drop it WITHOUT spawning, so we never end up with an
// orphaned worker running a forward the topology attributes to a
// 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)
continue
}
if e.Handler != nil {
if err := e.Handler.Claim(ctx, claimed); err != nil {
e.state.PendingTasks[claimed.ID] = claimed
@ -389,6 +455,13 @@ func (e *Engine) runCommands(ctx context.Context, tk *Token) error {
})
}
e.state.AddTopology(claimed, e.ID)
// Refresh our own stored load so the next selfIsLowest check in this
// same pass sees the incremented Forwards count — otherwise we'd keep
// claiming (stored load stays stale until we forward the token) and
// hog every claimable task, defeating lowest-load distribution.
if i := e.state.Find(e.ID); i >= 0 {
e.state.Nodes[i].Load = e.loadSnapshot()
}
}
return nil
}
@ -397,29 +470,72 @@ func (e *Engine) runCommands(ctx context.Context, tk *Token) error {
// It is the transport hook used by the HTTP handler after OnToken. If this
// node just self-removed, its ID is no longer in state.Nodes so
// AliveSuccessor would fail — use the original successor captured before
// removal (plan §移除节点 step 3: "令牌传递给自身原本的下一家").
// removal (plan §移除节点 step 3: "令牌传递给自身原本的下一家"). Otherwise
// delegates to forwardToNext which handles send-failure → mark offline →
// reassign → try next hop (plan §故障自幽).
func (e *Engine) Forward(ctx context.Context, tk *Token) error {
if e.removedNext != "" {
next := e.removedNext
e.removedNext = ""
var err error
if e.send != nil {
log.Printf("ring[%s] forward (self-removed) cycle=%d to %s", e.ID, tk.Cycle, next)
return e.send(ctx, next, tk)
err = e.send(ctx, next, tk)
}
return nil
// After handing off the token to the old successor, detach to a
// fresh standalone state so Snapshot() no longer serves the old
// cluster picture (members, topology, log) after self-leave. The
// token already carries the published state (with the new leader if
// we designated one); the local reset does not affect the sent token.
if e.selfRemoved {
e.detachAsStandalone()
}
return err
}
next, ok := e.state.AliveSuccessor(e.ID)
if !ok {
return nil // single-node ring
return e.forwardToNext(ctx, tk)
}
// detachAsStandalone resets the engine to a fresh standalone state after a
// self-leave has completed (the token was handed to the old successor). This
// prevents Snapshot() from serving the old cluster picture — other members,
// the full topology, pending tasks, and the cluster log — after the node has
// permanently left the ring. Equivalent to CreateCluster minus the IsMember
// guard (we are already detached) plus a fresh log (old cluster events stale).
func (e *Engine) detachAsStandalone() {
// Leaving the cluster invalidates this node's admission key — a
// standalone node has no key until it creates/joins again. Clear both
// the in-memory key and the persisted copy (so a restart doesn't
// resurrect a stale key for a node that's no longer in any cluster).
e.nodeKey = ""
if e.keyPersist != nil {
_ = e.keyPersist("")
}
if next == e.ID {
return nil // never forward to ourselves
// Clear the cached peer list so this node does NOT auto-rejoin on
// restart — it explicitly left the cluster.
e.clearPeers()
e.state = State{
LeaderID: e.ID,
PendingTasks: map[string]*Task{},
Topology: map[string]*TopoEntry{},
RoundDelay: 2 * time.Second,
}
if e.send != nil {
log.Printf("ring[%s] forward cycle=%d to %s", e.ID, tk.Cycle, next)
return e.send(ctx, next, tk)
e.state.UpsertNode(Node{
ID: e.ID, Addr: e.myAddr, Alive: true, IsLeader: true,
Load: e.loadSnapshot(), Version: e.Version, Cache: e.Cache, NodeKey: e.nodeKey,
})
e.selfRemoved = false
e.removedNext = ""
e.lastRingStart = time.Time{}
e.lastTokenAt = 0
e.lastSyncAt = 0
e.inflight.clear()
e.failCount = map[string]int{}
e.published = map[string]struct{}{}
if e.Log != nil {
e.Log = NewClusterLog()
_, _ = e.Log.Append(e.ID, LogLeaderChange, map[string]string{"leader": e.ID})
}
return nil
log.Printf("ring[%s] detached to standalone after self-leave", e.ID)
}
// Snapshot returns a serializable view of the ring for the frontend.
@ -433,6 +549,9 @@ type RingSnapshot struct {
Pending []*Task `json:"pending"`
Topology []*TopoEntry `json:"topology"`
Log []LogEntry `json:"log,omitempty"`
// NodeKey: this node's cluster admission key. The cluster page displays it
// so the operator can copy it for newcomers joining via this node.
NodeKey string `json:"nodeKey,omitempty"`
}
func (e *Engine) Snapshot() *RingSnapshot {
@ -445,6 +564,7 @@ func (e *Engine) Snapshot() *RingSnapshot {
Nodes: e.state.Nodes,
Pending: e.state.PendingList(),
Topology: e.state.TopologyList(),
NodeKey: e.nodeKey,
}
if e.Log != nil {
snap.Log = e.Log.Snapshot()
@ -509,6 +629,10 @@ type JoinInfo struct {
Addr string `json:"addr"`
Version string `json:"version,omitempty"`
Cache []string `json:"cache,omitempty"`
// JoinKey is the sponsor's nodeKey — the newcomer must present it to
// prove it is authorized to join via the sponsor. The sponsor verifies
// ji.JoinKey == e.nodeKey; mismatch → 403.
JoinKey string `json:"joinKey,omitempty"`
}
// injectPendingJoin writes every queued newcomer into state right after
@ -572,10 +696,17 @@ func (e *Engine) AdoptState(s State) {
e.state.PendingTasks[id] = t
}
e.state.UpsertNode(Node{ID: e.ID, Addr: e.myAddr, Alive: true,
Load: e.loadSnapshot(), Version: e.Version, Cache: e.Cache})
Load: e.loadSnapshot(), Version: e.Version, Cache: e.Cache, NodeKey: e.nodeKey})
if e.Log != nil {
e.Log.Append(e.ID, LogNodeJoin, map[string]string{"node": e.ID, "addr": e.myAddr})
}
// Newcomer generates its own admission key after joining, so future
// nodes can join via it. Per the user's design: "每个节点加入集群后
// 生成自身密钥". A node that already has a persisted key (restart
// re-join) keeps it.
e.ensureNodeKey()
// Persist the cached peer list from the adopted ring state.
e.persistPeers()
}
// CreateCluster reseeds this node as a fresh standalone leader (single-node
@ -588,17 +719,17 @@ func (e *Engine) CreateCluster() error {
if e.IsMember() {
return fmt.Errorf("node is a multi-node cluster member; leave first")
}
e.ensureNodeKey()
e.state = State{
LeaderID: e.ID,
Cycle: 0,
PendingTasks: map[string]*Task{},
Topology: map[string]*TopoEntry{},
RoundDelay: 2 * time.Second,
Seq: e.state.Seq,
}
e.state.UpsertNode(Node{
ID: e.ID, Addr: e.myAddr, Alive: true, IsLeader: true,
Load: e.loadSnapshot(), Version: e.Version, Cache: e.Cache,
Load: e.loadSnapshot(), Version: e.Version, Cache: e.Cache, NodeKey: e.nodeKey,
})
e.selfRemoved = false
e.removedNext = ""
@ -652,6 +783,80 @@ func (e *Engine) HasTask(local, remote string, port int) bool {
// IsLeader reports whether this node is the current ring leader.
func (e *Engine) IsLeader() bool { return e.state.LeaderID == e.ID }
// NodeKey returns this node's cluster admission key (for the frontend to
// display so the operator can copy it for newcomers).
func (e *Engine) NodeKey() string { return e.nodeKey }
// SetKeyPersist installs the callback used to persist the nodeKey to durable
// storage (store.SetNodeKey). Called once from main.go after NewEngine. Tests
// leave it nil — ensureNodeKey still generates the key in-memory.
func (e *Engine) SetKeyPersist(fn func(key string) error) { e.keyPersist = fn }
// SetPeerPersist installs the callback used to persist the cached peer list
// to durable storage (store.SetClusterPeers). Called once from main.go.
func (e *Engine) SetPeerPersist(fn func(peersJSON string) error) { e.peerPersist = fn }
// persistPeers extracts all alive peers (addr + nodeKey, excluding self)
// from the current ring state and persists them via the peerPersist callback.
// Called on every token cycle (OnToken) and on AdoptState so a crashed node
// always has the latest peer list to rejoin through. Skipped for standalone
// (single-node) rings — a standalone node has no peers to cache.
func (e *Engine) persistPeers() {
if e.peerPersist == nil {
return
}
type peerEntry struct {
Addr string `json:"addr"`
Key string `json:"key"`
}
var peers []peerEntry
for _, n := range e.state.Nodes {
if n.ID == e.ID || !n.Alive {
continue
}
if n.Addr == "" || n.NodeKey == "" {
continue
}
peers = append(peers, peerEntry{Addr: n.Addr, Key: n.NodeKey})
}
if len(peers) == 0 {
return // standalone or all-offline: don't overwrite a good cache
}
data, err := json.Marshal(peers)
if err != nil {
return
}
if err := e.peerPersist(string(data)); err != nil {
log.Printf("ring[%s] persist cluster peers failed: %v", e.ID, err)
}
}
// clearPeers wipes the cached peer list (called from detachAsStandalone so
// an explicit leave does NOT auto-rejoin on restart).
func (e *Engine) clearPeers() {
if e.peerPersist != nil {
_ = e.peerPersist("")
}
}
// ensureNodeKey generates a random admission key if this node doesn't have one
// yet, and persists it via the keyPersist callback (so it survives restarts).
// Called from CreateCluster (the creator generates a key so others can join
// via it) and AdoptState (a newcomer generates its own key after joining, so
// future nodes can join via it). Per the user's design: "每个节点加入集群后
// 生成自身密钥" — the key is born with cluster membership, not at startup.
func (e *Engine) ensureNodeKey() {
if e.nodeKey != "" {
return
}
e.nodeKey = GenerateNodeKey()
if e.keyPersist != nil {
if err := e.keyPersist(e.nodeKey); err != nil {
log.Printf("ring[%s] persist nodeKey failed: %v", e.ID, err)
}
}
}
// IsMember reports whether this node is currently an active multi-node member
// (self is in the ring alongside others). Used by the create/join gates to
// refuse actions that would split an active ring. A detached node (self not

View File

@ -3,6 +3,8 @@ package cluster
import (
"context"
"testing"
"webui4frpc/internal/store"
)
type fakeHandler struct {
@ -33,7 +35,7 @@ func sendNull(ctx context.Context, next string, tk *Token) error { return nil }
// task gets claimed (lowest load = only node).
func TestSingleNodeCycle(t *testing.T) {
eng := NewEngine("n1", "n1:7500", "u", "p", "0.71.0", []string{"0.71.0"},
&fakeHandler{load: Load{MemPct: 30, NetPct: 30}}, sendNull, "n1:7500", true)
&fakeHandler{load: Load{MemPct: 30, NetPct: 30}}, sendNull, "n1:7500", true, "")
eng.StartRing(context.Background())
// single node: token stays local; no successor so nothing travels.
@ -46,7 +48,7 @@ func TestSingleNodeCycle(t *testing.T) {
// stamps; then phase flips and second round syncs.
func TestTwoNodesSingleRound(t *testing.T) {
eng2 := NewEngine("n2", "n2:7500", "u", "p", "0.71.0", nil,
&fakeHandler{load: Load{MemPct: 10, NetPct: 10}}, sendNull, "n2:7500", false)
&fakeHandler{load: Load{MemPct: 10, NetPct: 10}}, sendNull, "n2:7500", false, "")
eng2.state.UpsertNode(Node{ID: "n1", Addr: "n1:7500", Alive: true, IsLeader: true, Load: Load{MemPct: 50, NetPct: 50}})
// n1 sends a token carrying the cluster picture; single-round processing:
@ -69,3 +71,119 @@ func TestTwoNodesSingleRound(t *testing.T) {
t.Fatal("nil token after processing")
}
}
// TestLeaderRoundTripNoResurrection: when the leader submits a task and the
// token completes a full round (leader→n2 claims→n3→leader), the consumed task
// MUST NOT be resurrected on the leader's next OnToken, and the topology must
// still attribute the forward to the single claimer (n2).
//
// This holds NOT via a published-set mark on StartRing, but via Go map
// reference semantics: StartRing sets tk.State = e.state, so the token's
// PendingTasks map IS the leader's own map. When n2 calls ClaimPending it
// deletes from that shared map — the deletion is visible to the leader too.
// By the time the token returns, the claimed task is already gone from the
// leader's e.state.PendingTasks, so the localPending re-merge has nothing to
// re-inject. Single token + remove-on-claim = single claim (plan §M6).
func TestLeaderRoundTripNoResurrection(t *testing.T) {
// 3-node ring: n1 (leader+submitter), n2 (lowest, claims), n3 (idle).
var sent *Token
sendCap := func(ctx context.Context, next string, tk *Token) error {
sent = tk
return nil
}
n1 := NewEngine("n1", "n1:7500", "u", "p", "v", nil,
&fakeHandler{load: Load{MemPct: 50, NetPct: 50}}, sendCap, "n1:7500", true, "")
// n2/n3 carry a lower stored load than n1's NewEngine default ({10,10})
// so LowestAlive picks n2 (first among the tied low nodes) as the claimer.
n1.state.UpsertNode(Node{ID: "n2", Addr: "n2:7500", Alive: true, Load: Load{MemPct: 1, NetPct: 1}})
n1.state.UpsertNode(Node{ID: "n3", Addr: "n3:7500", Alive: true, Load: Load{MemPct: 1, NetPct: 1}})
task := n1.SubmitTask(store.Local{Name: "l1"}, store.Remote{Name: "r1"}, store.Link{RemotePort: 100})
if task == nil {
t.Fatal("submit returned nil")
}
n1.StartRing(context.Background())
if sent == nil {
t.Fatal("StartRing did not send a token")
}
// n2 receives, claims t1 (lowest load), establishes topology.
n2 := NewEngine("n2", "n2:7500", "u", "p", "v", nil,
&fakeHandler{load: Load{MemPct: 10, NetPct: 10}}, sendCap, "n2:7500", false, "")
out2, err := n2.OnToken(context.Background(), sent)
if err != nil {
t.Fatal(err)
}
if got := len(n2.state.TopologyList()); got != 1 {
t.Fatalf("n2 should own 1 forward, topo=%+v", n2.state.TopologyList())
}
if len(n2.state.PendingList()) != 0 {
t.Fatalf("pending should be empty after n2 claim, got %+v", n2.state.PendingList())
}
// n3 receives, nothing to claim, forwards.
n3 := NewEngine("n3", "n3:7500", "u", "p", "v", nil,
&fakeHandler{load: Load{MemPct: 10, NetPct: 10}}, sendCap, "n3:7500", false, "")
out3, err := n3.OnToken(context.Background(), out2)
if err != nil {
t.Fatal(err)
}
// Token returns to n1. The consumed task t1 MUST NOT be resurrected.
if _, err := n1.OnToken(context.Background(), out3); err != nil {
t.Fatal(err)
}
if got := len(n1.state.PendingList()); got != 0 {
t.Fatalf("n1 resurrected a consumed task: pending=%+v (single-token + shared-map must prevent this)",
n1.state.PendingList())
}
// Topology still attributes t1 to n2 (not overwritten by a re-claim).
topo := n1.state.TopologyList()
if len(topo) != 1 || topo[0].OwnerID != "n2" {
t.Fatalf("topology should be t1@n2, got %+v", topo)
}
}
// TestClaimGuardDropsDuplicateForward: when a node (lowest load) receives a
// pending task whose forward is ALREADY in the topology owned by another node
// (a resurrected/stale copy), the claim path must drop it WITHOUT spawning a
// worker or overwriting the topology. Without the guard a second worker for
// the same forward would be spawned and orphaned (the topology entry is keyed
// by task id, so AddTopology would silently overwrite the owner).
func TestClaimGuardDropsDuplicateForward(t *testing.T) {
// n1 is lowest load and would otherwise claim; the fake Claim hook fails
// the test if ever called.
n1 := NewEngine("n1", "n1:7500", "u", "p", "v", nil,
&fakeHandler{
load: Load{MemPct: 10, NetPct: 10},
claim: func(ctx context.Context, tk *Task) error {
t.Fatalf("Claim must not be called for an already-owned forward: %s", tk.ID)
return nil
},
}, sendNull, "n1:7500", true, "")
dup := &Task{ID: "t1", Local: store.Local{Name: "l1"}, Remote: store.Remote{Name: "r1"}, Link: store.Link{RemotePort: 100}}
tk := &Token{Cycle: 1, State: State{
Nodes: []Node{
{ID: "n1", Addr: "n1:7500", Alive: true, Load: Load{MemPct: 10, NetPct: 10}},
{ID: "n2", Addr: "n2:7500", Alive: true, Load: Load{MemPct: 50, NetPct: 50}},
},
PendingTasks: map[string]*Task{"t1": dup},
Topology: map[string]*TopoEntry{"t1": {
TaskID: "t1", OwnerID: "n2", Local: store.Local{Name: "l1"},
Remote: store.Remote{Name: "r1"}, Link: store.Link{RemotePort: 100}, Active: true,
}},
}}
if _, err := n1.OnToken(context.Background(), tk); err != nil {
t.Fatal(err)
}
// Pending drained (the stale task was consumed/dropped, not left to ride).
if got := len(n1.state.PendingList()); got != 0 {
t.Fatalf("stale task should be dropped, pending=%+v", n1.state.PendingList())
}
// Topology untouched: still owned by n2.
topo := n1.state.TopologyList()
if len(topo) != 1 || topo[0].OwnerID != "n2" {
t.Fatalf("topology should remain t1@n2, got %+v", topo)
}
}

View File

@ -12,11 +12,29 @@ import (
"time"
)
// ringHopDelay paces each token hop (real-world cadence, avoids busy-loop).
// 500ms keeps a 3-node ring well under the LossTimeout floor (2500ms) so a
// healthy round is never misjudged lost, while not idly burning CPU/HTTP at
// 20Hz like the old 50ms did.
const ringHopDelay = 500 * time.Millisecond
// Per-hop pacing bounds. The hop delay scales DOWN as the ring grows so the
// round time stays ~ringHopDelayMax regardless of node count — a static
// 500ms made large clusters slow (3 nodes=1.5s, 5 nodes=2.5s, 10 nodes=5s
// per round); now 3/5/10 nodes all round at ~500ms (until the floor bites),
// keeping sync real-time without a token storm (round freq ≈2Hz).
const (
ringHopDelayMax = 500 * time.Millisecond
ringHopDelayMin = 50 * time.Millisecond
)
// hopDelayFor returns the per-hop pace for a ring of aliveNodes members.
// Nodes越多延迟越低: delay = ringHopDelayMax / aliveNodes, floored at min.
// n=2→250ms, n=3→167ms, n=5→100ms, n=10→50ms(floor) — round time ≈500ms.
func hopDelayFor(aliveNodes int) time.Duration {
if aliveNodes < 2 {
aliveNodes = 2
}
d := ringHopDelayMax / time.Duration(aliveNodes)
if d < ringHopDelayMin {
d = ringHopDelayMin
}
return d
}
// LossTimeout is the token-loss threshold per design: roundDelay/2 + 20ms,
// floored so a healthy fast ring is never misjudged.
@ -92,8 +110,14 @@ func (e *Engine) Send(ctx context.Context, tk *Token) error {
return e.forwardToNext(ctx, tk)
}
// forwardToNext sends the token to the next alive successor; on failure it
// marks that node offline, reattaches its tasks, and tries the next hop.
// forwardToNext sends the token to the next alive successor. On send
// failure (no receipt within the HTTP timeout = neighbor unreachable),
// the normal node-death procedure fires: mark offline, reassign the dead
// node's tasks, and try the next hop. If the dead node was the leader,
// the predecessor reuses the same procedure and additionally promotes
// itself to leader + starts a fresh cycle (the token destined for the
// dead leader is lost; a new cycle must begin). This is the PRIMARY
// leader-death detection path per plan §故障自幽 + §leader 补充/监控.
func (e *Engine) forwardToNext(ctx context.Context, tk *Token) error {
for hops := 0; hops < len(e.state.Nodes); hops++ {
next, ok := e.nextRecipient()
@ -104,26 +128,32 @@ func (e *Engine) forwardToNext(ctx context.Context, tk *Token) error {
if e.send == nil {
return nil
}
log.Printf("ring[%s] forward cycle=%d to %s", e.ID, tk.Cycle, next)
err := e.send(ctx, next, tk)
if err == nil {
if e.state.LeaderID == e.ID {
e.inflight.mark(e.state.RoundDelay)
}
// Pace the ring so tokens circulate at a realistic cadence.
// Pacing handled by the parallel rhythm timer in OnToken.
return nil
}
log.Printf("ring[%s] send to %s failed: %v", e.ID, next, err)
// Do not declare a neighbor offline on a single timeout: transient
// send failures (network jitter, busy handler) must not break the
// ring. Only after consecutive failures do we evict the node.
e.failMu.Lock()
e.failCount[next]++
if e.failCount[next] >= 2 {
e.state.MarkOffline(next)
e.state.OfflineReassign(next)
delete(e.failCount, next)
// Send failed = no receipt within timeout = neighbor offline.
// Normal node-death: mark offline, reassign tasks to pending.
log.Printf("ring[%s] send to %s failed (no receipt): %v", e.ID, next, err)
e.state.MarkOffline(next)
e.state.OfflineReassign(next)
// If the dead node was the leader, promote self and start a new
// cycle. The token was going to the leader; with the leader dead
// the token is lost — start fresh as the new leader (plan §leader
// 补充/监控: "上家邻居探测到 leader 崩溃 → 自身成为新 leader").
if next == e.state.LeaderID {
e.becomeLeader()
if e.Log != nil {
_, _ = e.Log.Append(e.ID, LogLeaderChange, map[string]string{"leader": e.ID})
}
e.StartRing(ctx)
return nil
}
// Non-leader neighbor death: continue to the next recipient.
}
e.inflight.clear()
return nil
@ -138,10 +168,26 @@ func (e *Engine) becomeLeader() {
log.Printf("ring[%s] promoted to leader", e.ID)
}
// WatchLeader runs the leader liveness monitor: the predecessor of the leader
// pings it; on failure it marks the leader offline and promotes itself.
// WatchLeader runs the FALLBACK leader liveness monitor. The PRIMARY path is
// forwardToNext: when the predecessor sends a token to the leader and the send
// fails (leader's HTTP server down), forwardToNext marks the leader offline and
// promotes self. WatchLeader covers the case forwardToNext CANNOT detect:
// the leader received the token (POST returned 200) but then crashed/restarted/
// detached before forwarding it — the send succeeded, so forwardToNext sees no
// error. In this case the predecessor pings the leader; if the leader is down
// (connection refused) or restarted/detached (standalone → 409), the heartbeat
// fails and the predecessor takes over.
//
// The predecessor role is NOT permanent — it shifts as the ring topology
// changes (nodes join/leave). Each tick re-evaluates AlivePredecessor(LeaderID)
// so the correct node monitors the leader at all times. Per design:
// "上邻居也不是永久的,也要有普通节点按照令牌传递的拓扑变换转换为上邻居的逻辑".
//
// Interval = 1s so worst-case detection (tick + 1.5s ping timeout ≈ 2.5s)
// aligns with LossTimeout (roundDelay/2 + 20ms, floored at 2500ms), per design:
// "与leader超时重发时间一致".
func (e *Engine) WatchLeader(ctx context.Context) {
tick := time.NewTicker(2 * time.Second)
tick := time.NewTicker(1 * time.Second)
defer tick.Stop()
for {
select {
@ -164,6 +210,11 @@ func (e *Engine) WatchLeader(ctx context.Context) {
e.state.MarkOffline(e.state.LeaderID)
e.state.OfflineReassign(e.state.LeaderID)
e.becomeLeader()
// Kick off a fresh cycle: the ring died with the old
// leader (no token inflight → WatchTokenLoss won't
// fire). Without this the newly promoted leader would
// sit idle and the ring would stay dead.
e.StartRing(ctx)
}
}
}

View File

@ -12,7 +12,7 @@ func newTestEngine(id string, isLeader bool) *Engine {
return NewEngine(id, id+":7500", "u", "p", "0.1.0", nil,
&fakeHandler{load: Load{MemPct: 20, NetPct: 20}},
func(ctx context.Context, next string, tk *Token) error { return nil },
id+":7500", isLeader)
id+":7500", isLeader, "")
}
// TestLogAppendDeltaReplay: append entries, extract delta after a watermark,

View File

@ -104,7 +104,7 @@ func TestRemoveCommandFulfilledNoPhantom(t *testing.T) {
eng := NewEngine("n1", "n1:7500", "u", "p", "0.1.0", nil,
&fakeHandler{load: Load{MemPct: 5, NetPct: 5},
claim: func(ctx context.Context, tk *Task) error { claimed = append(claimed, tk); return nil }},
sendNull, "n1:7500", true)
sendNull, "n1:7500", true, "")
rm := eng.state.AddRemoveNode("node-x:7500") // target not in the ring
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state}); err != nil {
t.Fatal(err)
@ -119,3 +119,102 @@ func TestRemoveCommandFulfilledNoPhantom(t *testing.T) {
t.Fatalf("phantom topology entry created: %+v", eng.state.TopologyList())
}
}
// TestLeaderSelfRemoveDesignatesSuccessor: when the LEADER self-removes, it
// must designate its successor as the new leader before publishing the token.
// Without this, LeaderID would be "" (SelfRemove clears it), no node would
// call Send (cycle never advances), and WatchLeader can't find
// AlivePredecessor("") to promote anyone — the ring runs leaderless and dies.
func TestLeaderSelfRemoveDesignatesSuccessor(t *testing.T) {
eng := newTestEngine("n1", true)
eng.state.InsertAfter("n1", Node{
ID: "n2", Addr: "n2:7500", Alive: true, Load: Load{MemPct: 90, NetPct: 90},
})
// n1 (leader) publishes a remove-node command for itself.
rm := eng.state.AddRemoveNode("n1")
eng.state.PendingTasks = map[string]*Task{rm.ID: rm}
out, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state})
if err != nil {
t.Fatal(err)
}
if out == nil {
t.Fatal("nil token after OnToken")
}
// n1 removed itself from the ring.
if eng.state.Find("n1") >= 0 {
t.Fatalf("n1 still in ring: %+v", eng.state.Nodes)
}
// n2 designated as the new leader (not "" — the old bug).
if eng.state.LeaderID != "n2" {
t.Fatalf("LeaderID = %q want n2 (successor should be designated as leader)", eng.state.LeaderID)
}
// n2 marked IsLeader in Nodes.
if i := eng.state.Find("n2"); i >= 0 && !eng.state.Nodes[i].IsLeader {
t.Fatalf("n2 not marked IsLeader: %+v", eng.state.Nodes[i])
}
// removedNext captured so Forward can hand the token to the old successor.
if eng.removedNext != "n2" {
t.Fatalf("removedNext = %q want n2", eng.removedNext)
}
}
// TestDetachAfterForward: after self-leave + Forward (token handed to the old
// successor), the engine resets to a fresh standalone state so Snapshot() no
// longer serves the old cluster picture (members, topology, pending, log).
func TestDetachAfterForward(t *testing.T) {
eng := newTestEngine("n1", true)
eng.state.InsertAfter("n1", Node{
ID: "n2", Addr: "n2:7500", Alive: true, Load: Load{MemPct: 90, NetPct: 90},
})
// Give n1 an owned forward so the old state has a non-empty topology.
p := eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
eng.state.ClaimPending(p.ID)
eng.state.AddTopology(p, "n1")
rm := eng.state.AddRemoveNode("n1")
eng.state.PendingTasks = map[string]*Task{rm.ID: rm}
out, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state})
if err != nil {
t.Fatal(err)
}
if out == nil {
t.Fatal("nil token after OnToken")
}
// Before Forward: self-removed, removedNext captured, old state retained
// (n2 still in Nodes, n1's forward re-queued as pending by SelfRemove).
if !eng.selfRemoved {
t.Fatal("selfRemoved not set before Forward")
}
if eng.removedNext != "n2" {
t.Fatalf("removedNext = %q want n2", eng.removedNext)
}
if eng.state.Find("n2") < 0 {
t.Fatal("n2 (old member) missing before Forward — test setup wrong")
}
// Forward hands off the token (sendNull no-op) then detaches to standalone.
if err := eng.Forward(context.Background(), out); err != nil {
t.Fatal(err)
}
// After Forward: fresh standalone state.
if eng.state.LeaderID != "n1" {
t.Fatalf("LeaderID = %q want n1 (standalone leader after detach)", eng.state.LeaderID)
}
if len(eng.state.Nodes) != 1 || eng.state.Nodes[0].ID != "n1" {
t.Fatalf("state not standalone (want just n1): %+v", eng.state.Nodes)
}
if len(eng.state.Topology) != 0 {
t.Fatalf("topology not cleared after detach: %+v", eng.state.Topology)
}
if len(eng.state.PendingTasks) != 0 {
t.Fatalf("pending not cleared after detach: %+v", eng.state.PendingList())
}
if eng.selfRemoved {
t.Fatal("selfRemoved should be cleared after detach")
}
if eng.removedNext != "" {
t.Fatalf("removedNext should be empty after detach, got %q", eng.removedNext)
}
}

View File

@ -15,7 +15,7 @@ func TestRevokeTaskRemovesTopology(t *testing.T) {
&fakeHandler{load: Load{MemPct: 5, NetPct: 5},
revoke: func(ctx context.Context, tk *Task) error { revoked = true; return nil }},
func(ctx context.Context, next string, tk *Token) error { return nil },
"n1:7500", true)
"n1:7500", true, "")
// establish a forward
eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})

View File

@ -6,12 +6,23 @@ package cluster
import (
"bytes"
"context"
"crypto/rand"
"encoding/json"
"fmt"
"net/http"
"time"
)
// GenerateNodeKey returns a random 16-byte hex string for use as a cluster
// admission key. Called on first startup when no key is persisted yet.
func GenerateNodeKey() string {
b := make([]byte, 16)
if _, err := rand.Read(b); err != nil {
return fmt.Sprintf("%x", time.Now().UnixNano())
}
return fmt.Sprintf("%x", b)
}
// JoinRing asks target to add us to its ring and returns the adopted state.
func (e *Engine) JoinRing(ctx context.Context, targetAddr string, ji JoinInfo) error {
url := fmt.Sprintf("http://%s/api/manager/cluster/join", targetAddr)
@ -49,9 +60,10 @@ func (e *Engine) JoinRing(ctx context.Context, targetAddr string, ji JoinInfo) e
// JoinRingAddr wraps JoinRing for HTTP handlers: it builds the newcomer's
// JoinInfo from this node's own id/addr/version/cache (so the caller doesn't
// touch the engine's unexported fields) and asks targetAddr to sponsor us
// into its ring. This is the runtime "加入集群" path (vs. the startup
// into its ring. joinKey is the sponsor's nodeKey — the sponsor verifies it
// before admitting. This is the runtime "加入集群" path (vs. the startup
// bootstrap call in cmd/webui4frpc/main.go).
func (e *Engine) JoinRingAddr(ctx context.Context, targetAddr string) error {
ji := JoinInfo{ID: e.ID, Addr: e.myAddr, Version: e.Version, Cache: e.Cache}
func (e *Engine) JoinRingAddr(ctx context.Context, targetAddr, joinKey string) error {
ji := JoinInfo{ID: e.ID, Addr: e.myAddr, Version: e.Version, Cache: e.Cache, JoinKey: joinKey}
return e.JoinRing(ctx, targetAddr, ji)
}