From 6fea93e821e937121846af6c01028a2ca2520433 Mon Sep 17 00:00:00 2001 From: pingqiu Date: Sat, 4 Apr 2026 00:59:10 -0700 Subject: [PATCH] =?UTF-8?q?feat:=20Task=20H=20=E2=80=94=20PendingCoordinat?= =?UTF-8?q?or=20extracted=20to=20sw-block/engine/replication/runtime?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit New reusable pending-execution coordinator with fail-closed command matching: - Store/TakeCatchUp/TakeRebuild/Cancel/Has/Peek - TakeCatchUp: fail-closed on target LSN mismatch (cancel + return nil) - TakeRebuild: same fail-closed semantics - Cancel callback invoked on mismatch or explicit cancellation 9 tests prove boundary behavior: - match succeeds, mismatch cancels, explicit cancel, noop on empty, peek non-destructive, store replaces, take from empty No weed/ imports. Pure coordination logic reusable by any adapter shell. weed/server/block_recovery.go rebinding deferred to Task I. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../engine/replication/runtime/pending.go | 139 ++++++++++++++++ .../replication/runtime/pending_test.go | 155 ++++++++++++++++++ 2 files changed, 294 insertions(+) create mode 100644 sw-block/engine/replication/runtime/pending.go create mode 100644 sw-block/engine/replication/runtime/pending_test.go diff --git a/sw-block/engine/replication/runtime/pending.go b/sw-block/engine/replication/runtime/pending.go new file mode 100644 index 000000000..000545043 --- /dev/null +++ b/sw-block/engine/replication/runtime/pending.go @@ -0,0 +1,139 @@ +// Package runtime provides reusable recovery-coordination helpers for the +// V2 engine. These helpers are independent of weed/ host specifics and can +// be used by any adapter shell. +package runtime + +import "sync" + +// PendingExecution holds the state needed to execute a planned recovery +// action. The coordinator stores one pending execution per volume and +// matches it against incoming commands. +type PendingExecution struct { + VolumeID string + ReplicaID string + + // CatchUpTarget is the target LSN for catch-up execution. + // Used by ExecutePendingCatchUp for fail-closed target matching. + CatchUpTarget uint64 + + // RebuildTargetLSN is the target LSN for rebuild execution. + // Used by ExecutePendingRebuild for fail-closed target matching. + RebuildTargetLSN uint64 + + // Opaque handles — the coordinator stores these but doesn't interpret them. + // The host adapter supplies concrete types (driver, plan, IO bindings). + Driver interface{} + Plan interface{} + CatchUpIO interface{} + RebuildIO interface{} +} + +// CancelFunc is called when a pending execution is cancelled due to +// mismatch or supersession. The reason string explains why. +type CancelFunc func(pending *PendingExecution, reason string) + +// PendingCoordinator manages pending recovery executions with fail-closed +// command matching. It is safe for concurrent use. +// +// Flow: +// 1. Store() caches a planned execution after recovery planning completes +// 2. TakeCatchUp/TakeRebuild matches an incoming command against the cached plan +// 3. If the target doesn't match, the pending plan is cancelled (fail-closed) +// 4. Cancel() explicitly cancels a pending execution +type PendingCoordinator struct { + mu sync.Mutex + pending map[string]*PendingExecution + cancelFn CancelFunc +} + +// NewPendingCoordinator creates a coordinator with the given cancel callback. +// The cancel function is called when a pending execution is cancelled due to +// target mismatch or explicit cancellation. +func NewPendingCoordinator(cancelFn CancelFunc) *PendingCoordinator { + return &PendingCoordinator{ + pending: make(map[string]*PendingExecution), + cancelFn: cancelFn, + } +} + +// Store caches a pending execution for a volume, replacing any previous one. +func (pc *PendingCoordinator) Store(volumeID string, pe *PendingExecution) { + pc.mu.Lock() + defer pc.mu.Unlock() + pc.pending[volumeID] = pe +} + +// TakeCatchUp takes the pending execution for the volume if the catch-up +// target matches. If there's a mismatch, the pending execution is cancelled +// (fail-closed) and nil is returned. If no pending execution exists, nil +// is returned. +func (pc *PendingCoordinator) TakeCatchUp(volumeID string, targetLSN uint64) *PendingExecution { + pc.mu.Lock() + pe, ok := pc.pending[volumeID] + if ok { + delete(pc.pending, volumeID) + } + pc.mu.Unlock() + + if !ok || pe == nil { + return nil + } + if pe.CatchUpTarget != targetLSN { + if pc.cancelFn != nil { + pc.cancelFn(pe, "start_catchup_target_mismatch") + } + return nil + } + return pe +} + +// TakeRebuild takes the pending execution for the volume if the rebuild +// target matches. Same fail-closed semantics as TakeCatchUp. +func (pc *PendingCoordinator) TakeRebuild(volumeID string, targetLSN uint64) *PendingExecution { + pc.mu.Lock() + pe, ok := pc.pending[volumeID] + if ok { + delete(pc.pending, volumeID) + } + pc.mu.Unlock() + + if !ok || pe == nil { + return nil + } + if pe.RebuildTargetLSN != targetLSN { + if pc.cancelFn != nil { + pc.cancelFn(pe, "start_rebuild_target_mismatch") + } + return nil + } + return pe +} + +// Has returns true if a pending execution exists for the volume. +func (pc *PendingCoordinator) Has(volumeID string) bool { + pc.mu.Lock() + defer pc.mu.Unlock() + _, ok := pc.pending[volumeID] + return ok +} + +// Cancel explicitly cancels and removes the pending execution for a volume. +func (pc *PendingCoordinator) Cancel(volumeID, reason string) { + pc.mu.Lock() + pe, ok := pc.pending[volumeID] + if ok { + delete(pc.pending, volumeID) + } + pc.mu.Unlock() + + if ok && pe != nil && pc.cancelFn != nil { + pc.cancelFn(pe, reason) + } +} + +// Peek returns the pending execution without removing it. Returns nil if none. +func (pc *PendingCoordinator) Peek(volumeID string) *PendingExecution { + pc.mu.Lock() + defer pc.mu.Unlock() + return pc.pending[volumeID] +} diff --git a/sw-block/engine/replication/runtime/pending_test.go b/sw-block/engine/replication/runtime/pending_test.go new file mode 100644 index 000000000..6f480ea47 --- /dev/null +++ b/sw-block/engine/replication/runtime/pending_test.go @@ -0,0 +1,155 @@ +package runtime + +import "testing" + +func TestPendingCoordinator_TakeCatchUp_MatchSucceeds(t *testing.T) { + pc := NewPendingCoordinator(nil) + pc.Store("vol1", &PendingExecution{ + VolumeID: "vol1", + CatchUpTarget: 100, + }) + + pe := pc.TakeCatchUp("vol1", 100) + if pe == nil { + t.Fatal("matching take should succeed") + } + if pe.CatchUpTarget != 100 { + t.Fatalf("target=%d", pe.CatchUpTarget) + } + + // Second take should return nil (already consumed). + if pc.TakeCatchUp("vol1", 100) != nil { + t.Fatal("second take should return nil") + } +} + +func TestPendingCoordinator_TakeCatchUp_MismatchCancels(t *testing.T) { + var cancelledReason string + pc := NewPendingCoordinator(func(pe *PendingExecution, reason string) { + cancelledReason = reason + }) + pc.Store("vol1", &PendingExecution{ + VolumeID: "vol1", + CatchUpTarget: 100, + }) + + pe := pc.TakeCatchUp("vol1", 999) // wrong target + if pe != nil { + t.Fatal("mismatched take should return nil") + } + if cancelledReason != "start_catchup_target_mismatch" { + t.Fatalf("cancel reason=%q", cancelledReason) + } + + // Pending should be consumed (removed even on mismatch). + if pc.Has("vol1") { + t.Fatal("pending should be removed after mismatch") + } +} + +func TestPendingCoordinator_TakeRebuild_MatchSucceeds(t *testing.T) { + pc := NewPendingCoordinator(nil) + pc.Store("vol1", &PendingExecution{ + VolumeID: "vol1", + RebuildTargetLSN: 50, + }) + + pe := pc.TakeRebuild("vol1", 50) + if pe == nil { + t.Fatal("matching rebuild take should succeed") + } +} + +func TestPendingCoordinator_TakeRebuild_MismatchCancels(t *testing.T) { + var cancelledReason string + pc := NewPendingCoordinator(func(pe *PendingExecution, reason string) { + cancelledReason = reason + }) + pc.Store("vol1", &PendingExecution{ + VolumeID: "vol1", + RebuildTargetLSN: 50, + }) + + pe := pc.TakeRebuild("vol1", 999) + if pe != nil { + t.Fatal("mismatched rebuild take should return nil") + } + if cancelledReason != "start_rebuild_target_mismatch" { + t.Fatalf("cancel reason=%q", cancelledReason) + } +} + +func TestPendingCoordinator_Cancel_ExplicitRemoval(t *testing.T) { + var cancelledReason string + pc := NewPendingCoordinator(func(pe *PendingExecution, reason string) { + cancelledReason = reason + }) + pc.Store("vol1", &PendingExecution{VolumeID: "vol1"}) + + pc.Cancel("vol1", "superseded") + if cancelledReason != "superseded" { + t.Fatalf("cancel reason=%q", cancelledReason) + } + if pc.Has("vol1") { + t.Fatal("should be removed after cancel") + } +} + +func TestPendingCoordinator_Cancel_NoopWhenEmpty(t *testing.T) { + cancelled := false + pc := NewPendingCoordinator(func(pe *PendingExecution, reason string) { + cancelled = true + }) + + pc.Cancel("vol1", "noop") + if cancelled { + t.Fatal("cancel on empty should not invoke callback") + } +} + +func TestPendingCoordinator_Has_And_Peek(t *testing.T) { + pc := NewPendingCoordinator(nil) + + if pc.Has("vol1") { + t.Fatal("empty coordinator should not have vol1") + } + if pc.Peek("vol1") != nil { + t.Fatal("peek on empty should return nil") + } + + pc.Store("vol1", &PendingExecution{VolumeID: "vol1", CatchUpTarget: 42}) + + if !pc.Has("vol1") { + t.Fatal("should have vol1 after store") + } + pe := pc.Peek("vol1") + if pe == nil || pe.CatchUpTarget != 42 { + t.Fatal("peek should return stored execution") + } + + // Peek does not consume. + if !pc.Has("vol1") { + t.Fatal("peek should not remove") + } +} + +func TestPendingCoordinator_StoreReplaces(t *testing.T) { + pc := NewPendingCoordinator(nil) + pc.Store("vol1", &PendingExecution{VolumeID: "vol1", CatchUpTarget: 10}) + pc.Store("vol1", &PendingExecution{VolumeID: "vol1", CatchUpTarget: 20}) + + pe := pc.TakeCatchUp("vol1", 20) + if pe == nil { + t.Fatal("replaced store should be latest") + } +} + +func TestPendingCoordinator_TakeFromEmpty_ReturnsNil(t *testing.T) { + pc := NewPendingCoordinator(nil) + if pc.TakeCatchUp("vol1", 100) != nil { + t.Fatal("take from empty should return nil") + } + if pc.TakeRebuild("vol1", 100) != nil { + t.Fatal("take rebuild from empty should return nil") + } +}