mirror of
https://gitcode.com/JianFeeeee/webui4frpc.git
synced 2026-10-02 23:24:00 +00:00
承接用户提问「设计上停用不是本来就会跨节点传输吗」——核实结论:结构上确实 如此(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 干净。
173 lines
5.0 KiB
Go
173 lines
5.0 KiB
Go
// Cluster operation log + incremental sync via the token.
|
|
//
|
|
// Every node keeps an append-only operation log (increasing Seq). Local
|
|
// mutations (forward created/removed, node joined/left, leader change, ...)
|
|
// are appended to the node's own log. During the token's sync round (phase 2)
|
|
// the token carries the log DELTA — entries with Seq > lastSyncedSeq — so
|
|
// every other node can append-and-replay them in order, converging on the same
|
|
// full cluster picture. Because each node has the full log, a newcomer can
|
|
// pull + replay the whole log to reconstruct identical state.
|
|
package cluster
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"sort"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// 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"
|
|
LogLeaderChange = "leader.change"
|
|
LogTaskClaimed = "task.claimed"
|
|
)
|
|
|
|
// LogEntry is one immutable, append-only cluster operation.
|
|
type LogEntry struct {
|
|
Seq int64 `json:"seq"`
|
|
Node string `json:"node"`
|
|
Kind string `json:"kind"`
|
|
Data json.RawMessage `json:"data,omitempty"`
|
|
At int64 `json:"at"`
|
|
}
|
|
|
|
// ClusterLog is a node's local operation log with a sync watermark.
|
|
type ClusterLog struct {
|
|
mu sync.Mutex
|
|
Seq int64 // highest local seq issued
|
|
Log []LogEntry // append-only
|
|
Synced int64 // watermark: entries <= this are known to remote nodes
|
|
}
|
|
|
|
// NewClusterLog creates an empty log with initial sequence.
|
|
func NewClusterLog() *ClusterLog {
|
|
return &ClusterLog{}
|
|
}
|
|
|
|
// Append adds an entry with the next sequence number (caller supplies kind/data).
|
|
// The entry is treated as already-synced (Synced advanced with Seq): a locally
|
|
// appended entry exists on this node and will be broadcast; the next full-log
|
|
// attachment from a peer must NOT re-apply it as a duplicate. Without this,
|
|
// a node that appends seq=N while its Synced lags behind would later receive
|
|
// its own seq=N in a peer's full attachment and ApplyDelta would re-add it
|
|
// (duplicate rows in the audit log).
|
|
func (l *ClusterLog) Append(node, kind string, data any) (LogEntry, error) {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
l.Seq++
|
|
raw, err := json.Marshal(data)
|
|
if err != nil {
|
|
l.Seq--
|
|
return LogEntry{}, err
|
|
}
|
|
e := LogEntry{Seq: l.Seq, Node: node, Kind: kind, Data: raw, At: time.Now().Unix()}
|
|
l.Log = append(l.Log, e)
|
|
if e.Seq > l.Synced {
|
|
l.Synced = e.Seq
|
|
}
|
|
return e, nil
|
|
}
|
|
|
|
// EntriesAfter returns entries with seq > after (delta for sync round).
|
|
func (l *ClusterLog) EntriesAfter(after int64) []LogEntry {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
var out []LogEntry
|
|
for _, e := range l.Log {
|
|
if e.Seq > after {
|
|
out = append(out, e)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// ApplyDelta appends-and-replays remote delta entries in order; returns the new
|
|
// watermark. Entries with seq <= existing watermark are skipped (idempotent).
|
|
func (l *ClusterLog) ApplyDelta(delta []LogEntry) (newWatermark int64, err error) {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
// ensure ordered by seq
|
|
sort.Slice(delta, func(i, j int) bool { return delta[i].Seq < delta[j].Seq })
|
|
last := l.Synced
|
|
for _, e := range delta {
|
|
if e.Seq <= last {
|
|
continue // already have
|
|
}
|
|
if e.Seq != last+1 {
|
|
return last, fmt.Errorf("gap in log delta: want %d got %d", last+1, e.Seq)
|
|
}
|
|
l.Log = append(l.Log, e)
|
|
last = e.Seq
|
|
}
|
|
l.Synced = last
|
|
if last > l.Seq {
|
|
l.Seq = last
|
|
}
|
|
return last, nil
|
|
}
|
|
|
|
// Replay applies a set of log entries locally (used at join to reconstruct
|
|
// state). Same idempotent-by-seq semantics as ApplyDelta.
|
|
func (l *ClusterLog) Replay(entries []LogEntry) error {
|
|
_, err := l.ApplyDelta(entries)
|
|
return err
|
|
}
|
|
|
|
// Snapshot returns a copy of the log (for debugging / join bootstrap).
|
|
func (l *ClusterLog) Snapshot() []LogEntry {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
out := make([]LogEntry, len(l.Log))
|
|
copy(out, l.Log)
|
|
return out
|
|
}
|
|
|
|
// DetailOf produces a human-readable one-line summary of a log entry's
|
|
// payload, matching the frontend's detailOf formatting. Used by the audit CSV
|
|
// export and anywhere a flat text rendering of an entry is needed.
|
|
func DetailOf(e LogEntry) string {
|
|
if len(e.Data) == 0 {
|
|
return ""
|
|
}
|
|
var d map[string]any
|
|
if err := json.Unmarshal(e.Data, &d); err != nil {
|
|
return string(e.Data)
|
|
}
|
|
str := func(key string) string { s, _ := d[key].(string); return s }
|
|
switch e.Kind {
|
|
case LogForwardAdd, LogForwardStop, LogForwardStart, LogForwardRemove:
|
|
s := str("local") + " → " + str("remote")
|
|
if id := str("taskId"); id != "" {
|
|
if len(id) > 8 {
|
|
id = id[len(id)-8:]
|
|
}
|
|
s += " · " + id
|
|
}
|
|
return s
|
|
case LogNodeJoin:
|
|
if addr := str("addr"); addr != "" {
|
|
return str("node") + " @ " + addr
|
|
}
|
|
return str("node")
|
|
case LogNodeLeave:
|
|
return str("node")
|
|
case LogLeaderChange:
|
|
return "→ " + str("leader")
|
|
case LogTaskClaimed:
|
|
return fmt.Sprintf("%s→%s:%v", str("local"), str("remote"), d["port"])
|
|
default:
|
|
if len(d) > 0 {
|
|
b, _ := json.Marshal(d)
|
|
return string(b)
|
|
}
|
|
return ""
|
|
}
|
|
}
|