mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 14:02:00 +02:00
Primary ships WAL entries to replica over TCP (data channel), confirms durability via barrier RPC (control channel). SyncCache runs local fsync and replica barrier in parallel via MakeDistributedSync. When replica is unreachable, shipper enters permanent degraded mode and falls back to local-only sync (Phase 3 behavior). Key design: two separate TCP ports (data+control), contiguous LSN enforcement, epoch equality check, WAL-full retry on replica, cond.Wait-based barrier with configurable timeout, BarrierFsyncFailed status code. Close lifecycle: shipper → receiver → drain → committer → flusher → fd. New files: repl_proto.go, wal_shipper.go, replica_apply.go, replica_barrier.go, dist_group_commit.go Modified: blockvol.go, blockvol_test.go 27 dev tests + 21 QA tests = 48 new tests; 889 total (609 engine + 280 iSCSI), all passing. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
47 lines
1.1 KiB
Go
47 lines
1.1 KiB
Go
package blockvol
|
|
|
|
import (
|
|
"sync"
|
|
)
|
|
|
|
// MakeDistributedSync creates a sync function that runs local WAL fsync and
|
|
// replica barrier in parallel. If no replica is configured or the replica is
|
|
// degraded, it falls back to local-only sync (Phase 3 behavior).
|
|
//
|
|
// walSync: the local fsync function (typically fd.Sync)
|
|
// shipper: the WAL shipper to the replica (may be nil)
|
|
// vol: the BlockVol (used to read nextLSN and trigger degradation)
|
|
func MakeDistributedSync(walSync func() error, shipper *WALShipper, vol *BlockVol) func() error {
|
|
return func() error {
|
|
if shipper == nil || shipper.IsDegraded() {
|
|
return walSync()
|
|
}
|
|
|
|
// The highest LSN that needs to be durable is nextLSN-1.
|
|
lsnMax := vol.nextLSN.Load() - 1
|
|
|
|
var localErr, remoteErr error
|
|
var wg sync.WaitGroup
|
|
wg.Add(2)
|
|
go func() {
|
|
defer wg.Done()
|
|
localErr = walSync()
|
|
}()
|
|
go func() {
|
|
defer wg.Done()
|
|
remoteErr = shipper.Barrier(lsnMax)
|
|
}()
|
|
wg.Wait()
|
|
|
|
if localErr != nil {
|
|
return localErr
|
|
}
|
|
if remoteErr != nil {
|
|
// Local succeeded, replica failed — degrade but don't fail the client.
|
|
vol.degradeReplica(remoteErr)
|
|
return nil
|
|
}
|
|
return nil
|
|
}
|
|
}
|