Files
webui4frpc/internal/cluster/ring_revoke_test.go
JianFeeeee 5cdc052d01 fix(cluster): restart 任务必须只由 owner 执行
上线实测抓到的:从 .60(非 owner)启动 owner 在 .30 的 portal,日志出现
**两条** `.60 restarted t4` —— 关键帧是 `[192.168.2.60:7500] restarted t4`。
pending 任务对所有成员可见,而 restart 分支没有 owner 判定,于是谁先轮到
谁就执行:非 owner 给自己的机器起了一个属于别人的转发 worker,而真正的
owner(.30)什么也没做,worker 一直没起来。

这正是上面 duplicate-claim 防御本来要防的「孤儿 worker + 重复认领」,只是
restart 分支为了绕开那层防御,把 owner 校验也一并跳过了 —— 绕开的是
「重复建条目」的必要性,不是「只有 owner 能动手」的必要性。

修法与 RemoveNode 分支一致:`owner != e.ID` 时只清标志、不执行 handler。
owner 已消失的情况不是错误:条目已被重新标为启用,OfflineReassign() 会在
下一次离线清理时把它转入 pending,再由正常认领流程重新安置。

测试 TestRestartOnlyAppliedByOwner:同一份 restart 任务分别交给 owner 与非
owner,断言非 owner 调用 0 次、owner 恰好 1 次。双向验证 —— 去掉守卫后
如期变红("a NON-owner applied the restart 1 time(s)"),还原后变绿。

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

402 lines
17 KiB
Go

package cluster
import (
"context"
"encoding/json"
"testing"
"webui4frpc/internal/store"
)
// TestRevokeTaskMarksTopologyDisabled: a stop delivered through the ring marks
// the topology entry disabled instead of deleting it, calls Handler.Revoke, and
// logs forward.stop.
//
// The entry is deliberately KEPT: TopoEntry.Link is the carrier that makes the
// stop travel, and State rides the token every cycle. Keeping it is what lets a
// stop reach a member that was offline when the stop was issued — deleting the
// entry left the flag nowhere to live, so a stop converged only if the owning
// node happened to be online for a one-shot revoke task.
func TestRevokeTaskMarksTopologyDisabled(t *testing.T) {
revoked := false
eng := NewEngine("n1", "n1:7500", "u", "p", "0.71.0", nil,
&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, "")
// establish a forward
eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
// lowest-load is n1 (only node): it claims + builds topology
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state}); err != nil {
t.Fatal(err)
}
if len(eng.state.TopologyList()) != 1 {
t.Fatalf("topology after claim = %+v", eng.state.TopologyList())
}
if d, _ := eng.state.TopologyDisabled("web", "frps1", 18081); d {
t.Fatal("a freshly claimed forward must not start out disabled")
}
// publish a revoke task pointing at the same forward
eng.state.AddRevoke(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: eng.state}); err != nil {
t.Fatal(err)
}
// The entry must SURVIVE, marked disabled, so the flag keeps travelling.
if len(eng.state.TopologyList()) != 1 {
t.Fatalf("topology entry was removed; the disabled flag would have no carrier: %+v",
eng.state.TopologyList())
}
disabled, known := eng.state.TopologyDisabled("web", "frps1", 18081)
if !known {
t.Fatal("topology entry missing after revoke")
}
if !disabled {
t.Fatal("revoke did not mark the topology entry disabled")
}
// Active must follow the flag, since OfflineReassign only re-queues Active
// entries — a stopped forward must not be resurrected by a node departure.
if eng.state.TopologyList()[0].Active {
t.Fatal("a stopped forward must not remain Active (OfflineReassign would re-queue it)")
}
if !revoked {
t.Fatal("Handler.Revoke was not called")
}
// log should contain forward.stop (not forward.remove — nothing was removed)
var sawStop bool
for _, e := range eng.Log.Snapshot() {
if e.Kind == LogForwardStop {
sawStop = true
}
}
if !sawStop {
t.Fatalf("log missing forward.stop: %+v", eng.Log.Snapshot())
}
}
// TestStoppedTopologyEntrySurvivesAdoption is the property that makes the whole
// mark-don't-remove design work: because the entry persists, a node that adopts
// ring state learns the stop. This is the offline-member case a one-shot revoke
// task could never cover.
func TestStoppedTopologyEntrySurvivesAdoption(t *testing.T) {
src := newTestEngine("n1", true)
src.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := src.OnToken(context.Background(), &Token{Cycle: 1, State: src.state}); err != nil {
t.Fatal(err)
}
src.state.AddRevoke(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := src.OnToken(context.Background(), &Token{Cycle: 2, State: src.state}); err != nil {
t.Fatal(err)
}
// A different node adopts the ring state (what e.state = tk.State does).
peer := newTestEngine("n2", false)
peer.AdoptState(src.state)
disabled, known := peer.TopologyDisabled("web", "frps1", 18081)
if !known {
t.Fatal("peer did not receive the topology entry for the stopped forward")
}
if !disabled {
t.Fatal("peer adopted the entry but no longer sees it as disabled — the stop did not propagate")
}
}
// TestRevokeIdempotent: revoking an already-missing forward does not error.
//
// It must ALSO still call Handler.Revoke. The handler is what actually stops
// the per-forward frpc worker; the topology entry is only bookkeeping. Gating
// the handler on RemoveTopology()'s return value (as this test used to permit)
// turned every revoke of an already-absent forward into a silent no-op — the
// worker kept running, which is how a forward the user had stopped kept
// dialling a local service that was intentionally down, forever.
func TestRevokeIdempotent(t *testing.T) {
revoked := 0
eng := NewEngine("n1", "n1:7500", "u", "p", "0.71.0", nil,
&fakeHandler{load: Load{MemPct: 5, NetPct: 5},
revoke: func(ctx context.Context, tk *Task) error { revoked++; return nil }},
func(ctx context.Context, next string, tk *Token) error { return nil },
"n1:7500", true, "")
eng.state.AddRevoke(store.Local{Name: "ghost"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 1})
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state}); err != nil {
t.Fatalf("revoke missing: %v", err)
}
// no topology entry, no panic
if len(eng.state.TopologyList()) != 0 {
t.Fatal("should be empty")
}
// The stop side-effect must have happened even though there was no entry.
if revoked != 1 {
t.Fatalf("Handler.Revoke called %d times, want 1 — a revoke with no topology entry "+
"must still stop the worker, otherwise stopped forwards keep running", revoked)
}
}
// TestStoppedForwardNotRequeuedOnNodeDeparture pins why marking (rather than
// deleting) is safe: OfflineReassign only re-queues ACTIVE entries, so a
// stopped forward is not silently handed to another node when its owner goes
// away. Deleting the entry, or leaving it Active, would both resurrect it.
func TestStoppedForwardNotRequeuedOnNodeDeparture(t *testing.T) {
eng := newTestEngine("n1", true)
eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state}); err != nil {
t.Fatal(err)
}
eng.state.AddRevoke(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: eng.state}); err != nil {
t.Fatal(err)
}
owner := eng.state.TopologyList()[0].OwnerID
eng.state.OfflineReassign(owner)
if n := len(eng.state.PendingList()); n != 0 {
t.Fatalf("a stopped forward was re-queued on the owner's departure (pending=%d): %+v",
n, eng.state.PendingList())
}
}
// TestSubmitTaskNotBlockedByStoppedEntry is the regression test for the bug the
// mark-don't-remove change introduced: SubmitTask's idempotency guard used
// HasTask, which matches a disabled topology entry. Since a stop now KEEPS the
// entry, any submission for that forward while it is still stopped is swallowed
// by its own guard — so a stopped forward can never be started again.
//
// The probe deliberately submits WITHOUT re-enabling first. Re-enabling would
// clear the flag and make HasTask and HasActiveTask agree, hiding the defect;
// the guard has to be exercised at the moment the entry is still stopped.
func TestSubmitTaskNotBlockedByStoppedEntry(t *testing.T) {
eng := newTestEngine("n1", true)
eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state}); err != nil {
t.Fatal(err)
}
// stop it -> entry is kept, marked disabled
eng.state.AddRevoke(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: eng.state}); err != nil {
t.Fatal(err)
}
if d, _ := eng.state.TopologyDisabled("web", "frps1", 18081); !d {
t.Fatal("precondition: forward should be disabled")
}
// The stopped entry IS findable by the broad predicate — which is exactly why
// the narrow one has to exist.
if !eng.HasTask("web", "frps1", 18081) {
t.Fatal("precondition: the stopped entry should still be findable by HasTask")
}
if eng.HasActiveTask("web", "frps1", 18081) {
t.Fatal("precondition: a stopped forward must NOT count as active")
}
// The regression: a submission must be published for a stopped forward.
// With a HasTask guard this returns nil and the forward is stuck forever.
tk := eng.SubmitTask(store.Local{Name: "web"}, store.Remote{Name: "frps1"},
store.Link{RemotePort: 18081})
if tk == nil {
t.Fatal("SubmitTask was swallowed by the stopped topology entry — " +
"a stopped forward could never be started again")
}
if tk.Link.RemotePort != 18081 {
t.Fatalf("unexpected task: %+v", tk)
}
}
// TestSubmitTaskStillDedupesActiveForward is the other half of the guard: the
// narrower predicate must not turn SubmitTask into a duplicate-task generator.
func TestSubmitTaskStillDedupesActiveForward(t *testing.T) {
eng := newTestEngine("n1", true)
eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state}); err != nil {
t.Fatal(err)
}
if !eng.HasActiveTask("web", "frps1", 18081) {
t.Fatal("precondition: a claimed forward must count as active")
}
if tk := eng.SubmitTask(store.Local{Name: "web"}, store.Remote{Name: "frps1"},
store.Link{RemotePort: 18081}); tk != nil {
t.Fatalf("SubmitTask must stay idempotent for an active forward, got %+v", tk)
}
}
// TestStoppedForwardNotRequeuedOnNodeDeparture pins why marking (rather than
// deleting) is safe: OfflineReassign only re-queues ACTIVE entries, so a
// stopped forward is not silently handed to another node when its owner goes
// TestAddTopologyRespectsDisabledFlag: a claim that carries a stopped forward
// must not mark the new entry Active, or it would come back to life.
func TestAddTopologyRespectsDisabledFlag(t *testing.T) {
s := &State{Topology: map[string]*TopoEntry{}}
e := s.AddTopology(&Task{
ID: "t1", Local: store.Local{Name: "web"}, Remote: store.Remote{Name: "frps1"},
Link: store.Link{RemotePort: 18081, Disabled: true},
}, "n1")
if e.Active {
t.Fatal("AddTopology forced Active=true for a disabled claim — the forward would resurrect")
}
if !e.Link.Disabled {
t.Fatal("AddTopology dropped the disabled flag from the link")
}
}
// TestRestartTaskBypassesDuplicateClaimGuard is the engine half of the
// re-enable path: a restart task for a forward that still HAS a topology entry
// must reach Handler.Restart.
//
// This is precisely what the duplicate-claim guard refuses — it exists to stop
// a stale task from spawning a second worker for an owned forward, and a restart
// is indistinguishable from that unless it is an explicit task kind. Marking
// the entry stopped (instead of deleting it) is what made this necessary: the
// entry the guard keys on now outlives a stop.
func TestRestartTaskBypassesDuplicateClaimGuard(t *testing.T) {
restarts := 0
claims := 0
eng := NewEngine("n1", "n1:7500", "u", "p", "0.71.0", nil,
&fakeHandler{load: Load{MemPct: 5, NetPct: 5},
claim: func(ctx context.Context, tk *Task) error { claims++; return nil },
restart: func(ctx context.Context, tk *Task) error { restarts++; return nil }},
func(ctx context.Context, next string, tk *Token) error { return nil },
"n1:7500", true, "")
// Establish + stop: entry survives, marked disabled.
eng.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 1, State: eng.state}); err != nil {
t.Fatal(err)
}
eng.state.AddRevoke(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 2, State: eng.state}); err != nil {
t.Fatal(err)
}
if d, _ := eng.state.TopologyDisabled("web", "frps1", 18081); !d {
t.Fatal("precondition: should be disabled")
}
entryCount := len(eng.state.TopologyList())
// Start: publish a restart and run a cycle.
eng.state.UpdateTopologyDisabled("web", "frps1", 18081, false)
eng.SubmitRestart(store.Local{Name: "web"}, store.Remote{Name: "frps1"},
store.Link{RemotePort: 18081})
if _, err := eng.OnToken(context.Background(), &Token{Cycle: 3, State: eng.state}); err != nil {
t.Fatal(err)
}
if restarts != 1 {
t.Fatalf("Handler.Restart called %d times, want 1 — the restart task was swallowed "+
"by the duplicate-claim guard, so the forward would never come back", restarts)
}
if claims != 1 {
t.Fatalf("Handler.Claim called %d times, want 1 (the original claim only); "+
"a restart must not re-run the claim bookkeeping", claims)
}
// Re-enabling must not create a second entry.
if n := len(eng.state.TopologyList()); n != entryCount {
t.Fatalf("restart created a duplicate topology entry: %d -> %d", entryCount, n)
}
if d, _ := eng.state.TopologyDisabled("web", "frps1", 18081); d {
t.Fatal("entry is still disabled after the restart task was applied")
}
var sawStart bool
for _, e := range eng.Log.Snapshot() {
if e.Kind == LogForwardStart {
sawStart = true
}
}
if !sawStart {
t.Fatalf("log missing forward.start: %+v", eng.Log.Snapshot())
}
}
// TestRestartFlagSurvivesTokenSerialization: the restart task travels to the
// owner inside the token, so the flag must round-trip through JSON. Without the
// struct tag it would silently deserialize as false and the owner would treat
// it as an ordinary claim — which the duplicate guard then discards.
func TestRestartFlagSurvivesTokenSerialization(t *testing.T) {
eng := newTestEngine("n1", true)
eng.SubmitRestart(store.Local{Name: "web"}, store.Remote{Name: "frps1"},
store.Link{RemotePort: 18081})
blob, err := json.Marshal(&Token{Cycle: 7, State: eng.state})
if err != nil {
t.Fatal(err)
}
var back Token
if err := json.Unmarshal(blob, &back); err != nil {
t.Fatal(err)
}
var found *Task
for _, tk := range back.State.PendingList() {
if tk.Local.Name == "web" && tk.Link.RemotePort == 18081 {
found = tk
}
}
if found == nil {
t.Fatal("restart task did not survive token serialization")
}
if !found.Restart {
t.Fatal("the restart flag was lost in JSON — the owner would see a plain claim " +
"and the duplicate guard would discard it")
}
// Revoke and Restart must stay distinguishable.
if found.Revoke {
t.Fatal("a restart task must not also read as a revocation")
}
}
// TestRestartOnlyAppliedByOwner is the regression test for the bug the restart
// channel introduced: pending tasks are visible to EVERY member, so without an
// owner check whichever node processed the task spawned a worker for a forward
// attributed to another node. Observed live as a non-owner logging
// "restarted t4" and running somebody else's forward.
func TestRestartOnlyAppliedByOwner(t *testing.T) {
// The forward is owned by n1. n1 and n2 both see the restart task.
restarts := 0
h := &fakeHandler{load: Load{MemPct: 5, NetPct: 5},
restart: func(ctx context.Context, tk *Task) error { restarts++; return nil }}
// Owner node.
owner := NewEngine("n1", "n1:7500", "u", "p", "0.71.0", nil, h,
func(ctx context.Context, next string, tk *Token) error { return nil }, "n1:7500", true, "")
owner.state.AddPending(store.Local{Name: "web"}, store.Remote{Name: "frps1"}, store.Link{RemotePort: 18081})
if _, err := owner.OnToken(context.Background(), &Token{Cycle: 1, State: owner.state}); err != nil {
t.Fatal(err)
}
if len(owner.state.TopologyList()) != 1 {
t.Fatalf("precondition: owner should hold the entry, got %+v", owner.state.TopologyList())
}
// A peer builds the SAME restart task but is not the owner.
peer := NewEngine("n2", "n2:7500", "u", "p", "0.71.0", nil, h,
func(ctx context.Context, next string, tk *Token) error { return nil }, "n2:7500", false, "")
peer.AdoptState(owner.state)
tk := peer.SubmitRestart(store.Local{Name: "web"}, store.Remote{Name: "frps1"},
store.Link{RemotePort: 18081})
if tk == nil {
t.Fatal("SubmitRestart returned nil")
}
peer.state.UpdateTopologyDisabled("web", "frps1", 18081, false)
// The peer must NOT run the handler: it does not own the forward.
if _, err := peer.OnToken(context.Background(), &Token{Cycle: 2, State: peer.state}); err != nil {
t.Fatal(err)
}
if restarts != 0 {
t.Fatalf("a NON-owner applied the restart %d time(s) — it would spawn an orphaned "+
"worker for a forward owned by somebody else", restarts)
}
// The owner must apply it.
restarts = 0
owner.state.UpdateTopologyDisabled("web", "frps1", 18081, false)
owner.SubmitRestart(store.Local{Name: "web"}, store.Remote{Name: "frps1"},
store.Link{RemotePort: 18081})
if _, err := owner.OnToken(context.Background(), &Token{Cycle: 3, State: owner.state}); err != nil {
t.Fatal(err)
}
if restarts != 1 {
t.Fatalf("the OWNER applied the restart %d times, want 1", restarts)
}
}