mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-24 17:04:30 +00:00
feat: Task H — PendingCoordinator extracted to sw-block/engine/replication/runtime
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) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.6
parent
519c849946
commit
6fea93e821
@@ -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]
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user