fix: 集群操作日志同步修复 — 全量回填 + 重启节点 seq 水位抬升

两个根因:
1. 离线错过条目永久缺失: token delta 被下游 trim (keep=seq>wm),
   离线节点错过的 seq 再也收不到 → 'log delta gap: want N got N+1' 死循环。
   修复: OnToken 改为附带自己的全量日志 (Snapshot),接收方 ApplyDelta 按 seq 幂等去重,
   缺口节点下一轮自动补齐。日志量小(几十条),开销可忽略。

2. 重启节点重新从 seq=1 编号: NewClusterLog 从 Seq=0 起,重启后产生的新事件
   与环上历史 seq 冲突 → ApplyDelta 视为 already-have 静默丢弃 + 全量附带出现歧义 id。
   修复: OnToken 收到日志时把本地 Seq 抬到环高水位之上,新编号接在历史之后。

测试: TestFullLogBackfillAfterGap / TestApplyDeltaIdempotentOnFullResend

实测: 三台集群撤销 minecraft 转发 → 三台均记录 seq=10 forward.remove;
恢复后三台均记录 seq=11 forward.add
This commit is contained in:
JianFeeeee
2026-08-25 11:55:34 +08:00
parent 94b6396738
commit 1abf1bb447
5 changed files with 350 additions and 218 deletions

View File

@ -238,6 +238,16 @@ func (e *Engine) OnToken(ctx context.Context, tk *Token) (*Token, error) {
// (b) apply incremental log delta; trim consumed entries off the token. // (b) apply incremental log delta; trim consumed entries off the token.
if e.Log != nil && len(tk.Log) > 0 { if e.Log != nil && len(tk.Log) > 0 {
// Raise the local seq allocator above the ring's high-water mark so a
// restarted node (whose log was rebuilt from scratch at Seq=0) never
// re-issues sequence numbers that already exist in the shared history —
// duplicate seqs would make ApplyDelta silently drop the new entries as
// "already have" and pollute full-log attachments with ambiguous ids.
for _, en := range tk.Log {
if en.Seq > e.Log.Seq {
e.Log.Seq = en.Seq
}
}
wm, err := e.Log.ApplyDelta(tk.Log) wm, err := e.Log.ApplyDelta(tk.Log)
if err != nil { if err != nil {
log.Printf("ring[%s] log delta gap: %v (request full sync later)", e.ID, err) log.Printf("ring[%s] log delta gap: %v (request full sync later)", e.ID, err)
@ -275,9 +285,15 @@ func (e *Engine) OnToken(ctx context.Context, tk *Token) (*Token, error) {
// flows to the newcomer (its successor) so it can participate. // flows to the newcomer (its successor) so it can participate.
e.injectPendingJoin() e.injectPendingJoin()
// re-attach own fresh log entries so peers converge. // re-attach OWN full log so peers converge even after gaps. A node that
// was offline when seq=N circulated never gets seq=N from the delta (it
// was trimmed off by peers whose watermark advanced past N). Attaching the
// FULL local log every cycle lets any peer missing entries backfill them
// next round — the cluster log is small (tens of entries) so the cost is
// negligible. ApplyDelta dedupes by seq so re-sent entries are a no-op
// for peers that already have them.
if !e.selfRemoved && e.Log != nil { if !e.selfRemoved && e.Log != nil {
mine := e.Log.EntriesAfter(e.lastLogSent) mine := e.Log.Snapshot()
if len(mine) > 0 { if len(mine) > 0 {
tk.Log = append(tk.Log, mine...) tk.Log = append(tk.Log, mine...)
if last := mine[len(mine)-1]; last.Seq > e.lastLogSent { if last := mine[len(mine)-1]; last.Seq > e.lastLogSent {

View File

@ -77,3 +77,60 @@ func TestClaimLogsToEngine(t *testing.T) {
t.Fatalf("engine log = %+v", snap) t.Fatalf("engine log = %+v", snap)
} }
} }
// TestFullLogBackfillAfterGap: a node that missed entries while offline
// (Synced stuck below the ring's max) must backfill from a peer's FULL log
// attachment. This is the regression test for the "log delta gap: want 1 got
// N" livelock where offline nodes could never rejoin the log history.
func TestFullLogBackfillAfterGap(t *testing.T) {
// Peer with complete history [1..4].
var peer ClusterLog
for i := 1; i <= 4; i++ {
if _, err := peer.Append("n1", LogNodeJoin, map[string]int{"i": i}); err != nil {
t.Fatal(err)
}
}
// Straggler that only has [1]; it missed [2..3] while offline and now
// receives the peer's FULL attachment [1..4].
straggler := NewClusterLog()
if _, err := straggler.Append("n2", LogNodeJoin, nil); err != nil {
t.Fatal(err)
}
// Force straggler to look like it has seq1 only (Synced=1).
straggler.Synced = 1
wm, err := straggler.ApplyDelta(peer.Snapshot())
if err != nil {
t.Fatalf("full backfill failed: %v", err)
}
if wm != 4 {
t.Fatalf("watermark = %d, want 4", wm)
}
if len(straggler.Snapshot()) != 4 {
t.Fatalf("log length = %d, want 4 (no dupes)", len(straggler.Snapshot()))
}
}
// TestApplyDeltaIdempotentOnFullResend: applying the same full attachment
// twice must not duplicate entries or error — peers re-attach their full log
// every cycle now.
func TestApplyDeltaIdempotentOnFullResend(t *testing.T) {
var src ClusterLog
for i := 1; i <= 3; i++ {
if _, err := src.Append("n1", LogNodeJoin, nil); err != nil {
t.Fatal(err)
}
}
full := src.Snapshot()
dst := NewClusterLog()
if _, err := dst.ApplyDelta(full); err != nil {
t.Fatalf("first apply: %v", err)
}
wm, err := dst.ApplyDelta(full)
if err != nil {
t.Fatalf("second apply (idempotence): %v", err)
}
if wm != 3 || len(dst.Snapshot()) != 3 {
t.Fatalf("wm=%d len=%d, want 3/3", wm, len(dst.Snapshot()))
}
}

View File

@ -144,10 +144,14 @@ func NewServeMux(h *Handler) (http.Handler, error) {
// resolve to admin via the flag fast path. // resolve to admin via the flag fast path.
mux.HandleFunc("/frpc/", h.auth("read")(h.handleFrpcBinary)) mux.HandleFunc("/frpc/", h.auth("read")(h.handleFrpcBinary))
// Static assets behind the same auth as the API: the browser caches the // Static assets WITHOUT auth: the SPA must load before it can show its
// Basic header once and sends it on every asset + /api/* request, so the // login form (auth.ts resolves GET /me on mount; a 401 there drops the UI
// SPA loads for any valid identity (viewer included). // into the login page). Gating index.html behind auth() would make an
mux.HandleFunc("/", h.auth("read")(h.handleStatic)) // unauthenticated browser see {"error":"unauthorized"} instead of the
// app shell — exactly the bug where the domain showed raw JSON. The
// embedded assets are static/public (no user data); every privileged
// action still requires an authenticated API call.
mux.HandleFunc("/", h.handleStatic)
return mux, nil return mux, nil
} }

View File

@ -4,6 +4,7 @@ import (
"bytes" "bytes"
"context" "context"
"encoding/json" "encoding/json"
"io"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"path/filepath" "path/filepath"
@ -249,3 +250,32 @@ func TestSaveCanvasPublishesRevokeTask(t *testing.T) {
t.Fatalf("expected revoke task for web, pending=%+v", ring.State().PendingList()) t.Fatalf("expected revoke task for web, pending=%+v", ring.State().PendingList())
} }
} }
// TestStaticServesWithoutAuth: the SPA shell (index.html) must load WITHOUT
// credentials so an unauthenticated browser sees the login page instead of a
// raw {"error":"unauthorized"} JSON body. API routes stay gated.
func TestStaticServesWithoutAuth(t *testing.T) {
_, ts := newTestHandler(t)
resp, err := http.Get(ts.URL + "/")
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
t.Fatalf("status = %d, want 200 for SPA shell", resp.StatusCode)
}
body, _ := io.ReadAll(resp.Body)
if !strings.Contains(string(body), "<div id=\"app\">") && !strings.Contains(string(body), "<!DOCTYPE html>") {
t.Fatalf("body does not look like index.html: %.80s", body)
}
// API remains gated.
req, _ := http.NewRequest(http.MethodGet, ts.URL+"/api/manager/status", nil)
resp2, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
defer resp2.Body.Close()
if resp2.StatusCode != http.StatusUnauthorized {
t.Fatalf("API status = %d, want 401", resp2.StatusCode)
}
}

View File

@ -34,7 +34,10 @@ async function request<T>(url: string, options: RequestInit = {}): Promise<T> {
const response = await fetch(url, { const response = await fetch(url, {
credentials: "same-origin", credentials: "same-origin",
...options, ...options,
headers: { ...UI_HEADER, ...(options.headers as Record<string, string> | undefined) }, headers: {
...UI_HEADER,
...(options.headers as Record<string, string> | undefined),
},
}); });
if (!response.ok) { if (!response.ok) {
throw new HTTPError(response.status, `HTTP ${response.status}`); throw new HTTPError(response.status, `HTTP ${response.status}`);
@ -106,7 +109,12 @@ export const api = {
}), }),
// assignGroup changes a single forward's group label (status page chip). // assignGroup changes a single forward's group label (status page chip).
// Empty group clears the assignment (移出分组). Pure DB update, no worker. // Empty group clears the assignment (移出分组). Pure DB update, no worker.
assignGroup: (local: string, remote: string, remotePort: number, group: string) => assignGroup: (
local: string,
remote: string,
remotePort: number,
group: string,
) =>
request<{ ok: boolean }>("/api/manager/forwards/assign", { request<{ ok: boolean }>("/api/manager/forwards/assign", {
method: "POST", method: "POST",
headers: { "Content-Type": "application/json" }, headers: { "Content-Type": "application/json" },
@ -184,13 +192,24 @@ export const api = {
logout: () => logout: () =>
request<{ ok: boolean }>("/api/manager/logout", { method: "POST" }), request<{ ok: boolean }>("/api/manager/logout", { method: "POST" }),
listUsers: () => request<{ users: User[] }>("/api/manager/users"), listUsers: () => request<{ users: User[] }>("/api/manager/users"),
createUser: (username: string, password: string, role: 'admin' | 'viewer' | 'superadmin') => createUser: (
username: string,
password: string,
role: "admin" | "viewer" | "superadmin",
) =>
request<User>("/api/manager/users", { request<User>("/api/manager/users", {
method: "POST", method: "POST",
headers: { "Content-Type": "application/json" }, headers: { "Content-Type": "application/json" },
body: JSON.stringify({ username, password, role }), body: JSON.stringify({ username, password, role }),
}), }),
updateUser: (name: string, patch: { password?: string; role?: 'admin' | 'viewer' | 'superadmin'; enabled?: boolean }) => updateUser: (
name: string,
patch: {
password?: string;
role?: "admin" | "viewer" | "superadmin";
enabled?: boolean;
},
) =>
request<User>(`/api/manager/users/${encodeURIComponent(name)}`, { request<User>(`/api/manager/users/${encodeURIComponent(name)}`, {
method: "PUT", method: "PUT",
headers: { "Content-Type": "application/json" }, headers: { "Content-Type": "application/json" },
@ -201,7 +220,11 @@ export const api = {
method: "DELETE", method: "DELETE",
}), }),
listApiKeys: () => request<{ apiKeys: ApiKey[] }>("/api/manager/apikeys"), listApiKeys: () => request<{ apiKeys: ApiKey[] }>("/api/manager/apikeys"),
createApiKey: (userId: number, label: string, scope: 'read' | 'write' | 'admin') => createApiKey: (
userId: number,
label: string,
scope: "read" | "write" | "admin",
) =>
request<ApiKeyCreated>("/api/manager/apikeys", { request<ApiKeyCreated>("/api/manager/apikeys", {
method: "POST", method: "POST",
headers: { "Content-Type": "application/json" }, headers: { "Content-Type": "application/json" },
@ -211,7 +234,8 @@ export const api = {
request<void>(`/api/manager/apikeys/${id}`, { method: "DELETE" }), request<void>(`/api/manager/apikeys/${id}`, { method: "DELETE" }),
// Canvas export/import (转发表 备份/还原). // Canvas export/import (转发表 备份/还原).
exportCanvas: () => request<CanvasExportEnvelope>("/api/manager/canvas/export"), exportCanvas: () =>
request<CanvasExportEnvelope>("/api/manager/canvas/export"),
importCanvas: (data: CanvasData | CanvasExportEnvelope) => importCanvas: (data: CanvasData | CanvasExportEnvelope) =>
request<CanvasData>("/api/manager/canvas/import", { request<CanvasData>("/api/manager/canvas/import", {
method: "POST", method: "POST",
@ -220,7 +244,8 @@ export const api = {
}), }),
// Worker-log bundle (HTTP fan-out across ring nodes; read-level/auditor). // Worker-log bundle (HTTP fan-out across ring nodes; read-level/auditor).
exportWorkerLogs: () => request<WorkerLogBundle>("/api/manager/cluster/logs/export"), exportWorkerLogs: () =>
request<WorkerLogBundle>("/api/manager/cluster/logs/export"),
}; };
// downloadAuditCsv streams one of the /audit/*.csv endpoints to a file. // downloadAuditCsv streams one of the /audit/*.csv endpoints to a file.