feat: V2 MVP milestone — masterv2 + volumev2 + in-process failover

V2 runtime packages:
- sw-block/runtime/masterv2: identity authority (desired state,
  heartbeat handling, promotion arbitration via SelectPromotionCandidate)
- sw-block/runtime/volumev2: per-volume micro-cluster shell (node,
  orchestrator, control session, iSCSI frontend, takeover gate,
  failover session + driver, replica summary reconstruction)
- sw-block/runtime/purev2: RF1 execution shell (engine + store +
  dispatcher + local boundary observations)
- sw-block/runtime/protocolv2: three-channel separation
  (heartbeat/assignment/query + replica summary)

V2 binaries:
- sw-block/cmd/v2singleblock: single-node RF1 block server
- sw-block/cmd/purev2rf1: minimal RF1 runtime binary

Milestone capabilities:
- RF1 write/read/sync with engine-driven mode projection
- masterv2 ↔ volumev2 heartbeat convergence + assignment reissue
- Promotion query with fresh CommittedLSN/WALHeadLSN evidence
- Replica summary for bounded takeover reconstruction
- Primary-loss reconstruction from peer summaries (fail-closed gate)
- In-process failover driver with session observability
- Local boundary observations feed engine (Committed/Durable/Checkpoint)

Design docs:
- v2-two-loop-protocol.md: identity vs data-control separation
- v2-automata-ownership-map.md: event/command ownership split
- v2-loop1-surface-draft.md: heartbeat/query/assignment field spec
- v2-volumev2-single-node-mvp.md: target layering
- v2-kernel-closure-review.md: per-volume micro-cluster principle
- v2-pure-runtime-rf1-bootstrap.md, v2-capability-map.md,
  v2-proof-and-retest-pyramid.md

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
pingqiuandClaude Opus 4.6 committed 2026-04-05 13:08:02 -07:00
1 parent cf16e53b04
commit b8c6944e3f
30 files changed
+6667

No files matched your search

+110
View File
@@ -0,0 +1,110 @@
package main
import (
"encoding/json"
"flag"
"fmt"
"os"
"time"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/purev2"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol"
)
func main() {
if len(os.Args) < 2 {
usage()
os.Exit(2)
}
switch os.Args[1] {
case "bootstrap":
if err := runBootstrap(os.Args[2:]); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
case "status":
if err := runStatus(os.Args[2:]); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
default:
usage()
os.Exit(2)
}
}
func runBootstrap(args []string) error {
fs := flag.NewFlagSet("bootstrap", flag.ContinueOnError)
path := fs.String("path", "", "block volume path")
sizeBytes := fs.Uint64("size-bytes", 1*1024*1024, "logical volume size in bytes")
blockSize := fs.Uint("block-size", 4096, "block size in bytes")
walSize := fs.Uint64("wal-size", 256*1024, "wal size in bytes")
epoch := fs.Uint64("epoch", 1, "assignment epoch")
leaseMs := fs.Int("lease-ms", 30000, "lease ttl in milliseconds")
if err := fs.Parse(args); err != nil {
return err
}
if *path == "" {
return fmt.Errorf("bootstrap: --path is required")
}
rt := purev2.New(purev2.Config{})
defer rt.Close()
opts := blockvol.CreateOptions{
VolumeSize: *sizeBytes,
BlockSize: uint32(*blockSize),
WALSize: *walSize,
}
if err := rt.BootstrapPrimary(*path, opts, *epoch, time.Duration(*leaseMs)*time.Millisecond); err != nil {
return err
}
snap, err := rt.Snapshot(*path)
if err != nil {
return err
}
return printJSON(snap)
}
func runStatus(args []string) error {
fs := flag.NewFlagSet("status", flag.ContinueOnError)
path := fs.String("path", "", "block volume path")
epoch := fs.Uint64("epoch", 0, "optional epoch to rebuild core projection")
leaseMs := fs.Int("lease-ms", 30000, "lease ttl in milliseconds")
if err := fs.Parse(args); err != nil {
return err
}
if *path == "" {
return fmt.Errorf("status: --path is required")
}
rt := purev2.New(purev2.Config{})
defer rt.Close()
if err := rt.OpenVolume(*path); err != nil {
return err
}
if *epoch > 0 {
if err := rt.ApplyPrimaryAssignment(*path, *epoch, time.Duration(*leaseMs)*time.Millisecond); err != nil {
return err
}
}
snap, err := rt.Snapshot(*path)
if err != nil {
return err
}
return printJSON(snap)
}
func printJSON(v any) error {
enc := json.NewEncoder(os.Stdout)
enc.SetIndent("", " ")
return enc.Encode(v)
}
func usage() {
fmt.Fprintln(os.Stderr, "usage:")
fmt.Fprintln(os.Stderr, " purev2rf1 bootstrap --path <file> [--size-bytes N --block-size N --wal-size N --epoch N]")
fmt.Fprintln(os.Stderr, " purev2rf1 status --path <file> [--epoch N]")
}
+436
View File
@@ -0,0 +1,436 @@
package main
import (
"encoding/hex"
"encoding/json"
"flag"
"fmt"
"io"
"net"
"os"
"time"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/masterv2"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/volumev2"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol/iscsi"
)
type smokeResult struct {
VolumeName string `json:"volume_name"`
Path string `json:"path"`
NodeID string `json:"node_id"`
Epoch uint64 `json:"epoch"`
Role string `json:"role"`
Mode string `json:"mode"`
Reason string `json:"reason"`
Readback string `json:"readback_hex"`
}
type restartSmokeResult struct {
VolumeName string `json:"volume_name"`
Path string `json:"path"`
NodeID string `json:"node_id"`
Epoch uint64 `json:"epoch"`
Role string `json:"role"`
Mode string `json:"mode"`
Reason string `json:"reason"`
InitialWrite string `json:"initial_write_hex"`
PostRestart string `json:"post_restart_hex"`
}
type iscsiSmokeResult struct {
VolumeName string `json:"volume_name"`
Path string `json:"path"`
NodeID string `json:"node_id"`
Epoch uint64 `json:"epoch"`
Role string `json:"role"`
Mode string `json:"mode"`
Reason string `json:"reason"`
IQN string `json:"iqn"`
Address string `json:"address"`
}
type commonFlags struct {
name string
path string
nodeID string
writeText string
sizeBytes uint64
blockSize uint
walSize uint64
}
type singleNodeEnv struct {
master *masterv2.Master
node *volumev2.Node
orchestrator *volumev2.Orchestrator
}
var commandOutput io.Writer = os.Stdout
func main() {
if len(os.Args) < 2 {
usage()
os.Exit(2)
}
switch os.Args[1] {
case "smoke":
if err := runSmoke(os.Args[2:]); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
case "restart-smoke":
if err := runRestartSmoke(os.Args[2:]); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
case "iscsi-smoke":
if err := runISCSISmoke(os.Args[2:]); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
default:
usage()
os.Exit(2)
}
}
func runSmoke(args []string) error {
cfg, err := parseCommonFlags("smoke", args)
if err != nil {
return err
}
if cfg.path == "" {
return fmt.Errorf("smoke: --path is required")
}
env, err := bootstrapSingleNode(cfg)
if err != nil {
return err
}
defer env.close()
payload := paddedPayload([]byte(cfg.writeText), uint32(cfg.blockSize))
if err := env.node.WriteLBA(cfg.name, 0, payload); err != nil {
return err
}
if err := env.node.SyncCache(cfg.name); err != nil {
return err
}
readBack, err := env.node.ReadLBA(cfg.name, 0, uint32(len(payload)))
if err != nil {
return err
}
snap, err := env.node.Snapshot(cfg.name)
if err != nil {
return err
}
result := smokeResult{
VolumeName: cfg.name,
Path: cfg.path,
NodeID: cfg.nodeID,
Epoch: snap.Status.Epoch,
Role: snap.Status.Role.String(),
Readback: hex.EncodeToString(readBack[:len(payload)]),
}
if snap.HasProjection {
result.Mode = string(snap.Projection.Mode.Name)
result.Reason = snap.Projection.Publication.Reason
}
return printJSON(result)
}
func runRestartSmoke(args []string) error {
cfg, err := parseCommonFlags("restart-smoke", args)
if err != nil {
return err
}
if cfg.path == "" {
return fmt.Errorf("restart-smoke: --path is required")
}
payload := paddedPayload([]byte(cfg.writeText), uint32(cfg.blockSize))
master := masterv2.New(masterv2.Config{LeaseTTL: 30 * time.Second})
if err := declarePrimary(master, cfg); err != nil {
return err
}
func() {
node, nodeErr := volumev2.New(volumev2.Config{NodeID: cfg.nodeID})
if nodeErr != nil {
err = nodeErr
return
}
defer node.Close()
session, nodeErr := volumev2.NewInProcessSession(master)
if nodeErr != nil {
err = nodeErr
return
}
orchestrator, nodeErr := volumev2.NewOrchestrator(node, session)
if nodeErr != nil {
err = nodeErr
return
}
if nodeErr = syncTwice(orchestrator); nodeErr != nil {
err = nodeErr
return
}
if nodeErr = node.WriteLBA(cfg.name, 0, payload); nodeErr != nil {
err = nodeErr
return
}
if nodeErr = node.SyncCache(cfg.name); nodeErr != nil {
err = nodeErr
return
}
}()
if err != nil {
return err
}
node, err := volumev2.New(volumev2.Config{NodeID: cfg.nodeID})
if err != nil {
return err
}
defer node.Close()
session, err := volumev2.NewInProcessSession(master)
if err != nil {
return err
}
orchestrator, err := volumev2.NewOrchestrator(node, session)
if err != nil {
return err
}
if err := syncTwice(orchestrator); err != nil {
return err
}
readBack, err := node.ReadLBA(cfg.name, 0, uint32(len(payload)))
if err != nil {
return err
}
snap, err := node.Snapshot(cfg.name)
if err != nil {
return err
}
result := restartSmokeResult{
VolumeName: cfg.name,
Path: cfg.path,
NodeID: cfg.nodeID,
Epoch: snap.Status.Epoch,
Role: snap.Status.Role.String(),
InitialWrite: hex.EncodeToString(payload),
PostRestart: hex.EncodeToString(readBack[:len(payload)]),
}
if snap.HasProjection {
result.Mode = string(snap.Projection.Mode.Name)
result.Reason = snap.Projection.Publication.Reason
}
return printJSON(result)
}
func runISCSISmoke(args []string) error {
fs := flag.NewFlagSet("iscsi-smoke", flag.ContinueOnError)
name := fs.String("name", "single-node-vol", "logical volume name")
path := fs.String("path", "", "block volume file path")
nodeID := fs.String("node", "node-a", "volumev2 node id")
writeText := fs.String("write-text", "v2-single-node-smoke", "payload written at LBA 0")
sizeBytes := fs.Uint64("size-bytes", 1*1024*1024, "logical volume size in bytes")
blockSize := fs.Uint("block-size", 4096, "block size in bytes")
walSize := fs.Uint64("wal-size", 256*1024, "wal size in bytes")
listenAddr := fs.String("listen-addr", "127.0.0.1:0", "iSCSI target listen address")
iqn := fs.String("iqn", "", "target iqn")
if err := fs.Parse(args); err != nil {
return err
}
cfg := commonFlags{
name: *name,
path: *path,
nodeID: *nodeID,
writeText: *writeText,
sizeBytes: *sizeBytes,
blockSize: *blockSize,
walSize: *walSize,
}
if cfg.path == "" {
return fmt.Errorf("iscsi-smoke: --path is required")
}
env, err := bootstrapSingleNode(cfg)
if err != nil {
return err
}
defer env.close()
export, err := env.node.ExportISCSI(cfg.name, *listenAddr, *iqn)
if err != nil {
return err
}
defer export.Close()
if err := performISCSILogin(export.Address(), export.IQN()); err != nil {
return err
}
snap, err := env.node.Snapshot(cfg.name)
if err != nil {
return err
}
result := iscsiSmokeResult{
VolumeName: cfg.name,
Path: cfg.path,
NodeID: cfg.nodeID,
Epoch: snap.Status.Epoch,
Role: snap.Status.Role.String(),
IQN: export.IQN(),
Address: export.Address(),
}
if snap.HasProjection {
result.Mode = string(snap.Projection.Mode.Name)
result.Reason = snap.Projection.Publication.Reason
}
return printJSON(result)
}
func parseCommonFlags(name string, args []string) (commonFlags, error) {
fs := flag.NewFlagSet(name, flag.ContinueOnError)
cfg := commonFlags{}
namePtr := fs.String("name", "single-node-vol", "logical volume name")
pathPtr := fs.String("path", "", "block volume file path")
nodePtr := fs.String("node", "node-a", "volumev2 node id")
writePtr := fs.String("write-text", "v2-single-node-smoke", "payload written at LBA 0")
sizePtr := fs.Uint64("size-bytes", 1*1024*1024, "logical volume size in bytes")
blockPtr := fs.Uint("block-size", 4096, "block size in bytes")
walPtr := fs.Uint64("wal-size", 256*1024, "wal size in bytes")
if err := fs.Parse(args); err != nil {
return cfg, err
}
cfg.name = *namePtr
cfg.path = *pathPtr
cfg.nodeID = *nodePtr
cfg.writeText = *writePtr
cfg.sizeBytes = *sizePtr
cfg.blockSize = *blockPtr
cfg.walSize = *walPtr
return cfg, nil
}
func bootstrapSingleNode(cfg commonFlags) (*singleNodeEnv, error) {
master := masterv2.New(masterv2.Config{LeaseTTL: 30 * time.Second})
if err := declarePrimary(master, cfg); err != nil {
return nil, err
}
node, err := volumev2.New(volumev2.Config{NodeID: cfg.nodeID})
if err != nil {
return nil, err
}
session, err := volumev2.NewInProcessSession(master)
if err != nil {
node.Close()
return nil, err
}
orchestrator, err := volumev2.NewOrchestrator(node, session)
if err != nil {
node.Close()
return nil, err
}
if err := syncTwice(orchestrator); err != nil {
node.Close()
return nil, err
}
return &singleNodeEnv{
master: master,
node: node,
orchestrator: orchestrator,
}, nil
}
func declarePrimary(master *masterv2.Master, cfg commonFlags) error {
return master.DeclarePrimary(masterv2.VolumeSpec{
Name: cfg.name,
Path: cfg.path,
PrimaryNodeID: cfg.nodeID,
CreateOptions: blockvol.CreateOptions{
VolumeSize: cfg.sizeBytes,
BlockSize: uint32(cfg.blockSize),
WALSize: cfg.walSize,
},
})
}
func (env *singleNodeEnv) close() {
if env == nil || env.node == nil {
return
}
env.node.Close()
}
func syncTwice(orchestrator *volumev2.Orchestrator) error {
if err := orchestrator.SyncOnce(); err != nil {
return err
}
return orchestrator.SyncOnce()
}
func performISCSILogin(addr, targetIQN string) error {
conn, err := net.DialTimeout("tcp", addr, 2*time.Second)
if err != nil {
return fmt.Errorf("dial iscsi target %s: %w", addr, err)
}
defer conn.Close()
params := iscsi.NewParams()
params.Set("InitiatorName", "iqn.2026-04.com.seaweedfs:cli.initiator")
params.Set("TargetName", targetIQN)
params.Set("SessionType", "Normal")
loginReq := &iscsi.PDU{}
loginReq.SetOpcode(iscsi.OpLoginReq)
loginReq.SetLoginStages(iscsi.StageSecurityNeg, iscsi.StageFullFeature)
loginReq.SetLoginTransit(true)
loginReq.SetISID([6]byte{0x00, 0x02, 0x3D, 0x00, 0x00, 0x02})
loginReq.SetCmdSN(1)
loginReq.DataSegment = params.Encode()
if err := iscsi.WritePDU(conn, loginReq); err != nil {
return fmt.Errorf("write iscsi login: %w", err)
}
resp, err := iscsi.ReadPDU(conn)
if err != nil {
return fmt.Errorf("read iscsi login: %w", err)
}
if resp.LoginStatusClass() != iscsi.LoginStatusSuccess {
return fmt.Errorf("iscsi login failed: %d/%d", resp.LoginStatusClass(), resp.LoginStatusDetail())
}
return nil
}
func paddedPayload(in []byte, blockSize uint32) []byte {
size := int(blockSize)
if size <= 0 {
size = len(in)
}
if size < len(in) {
size = len(in)
}
out := make([]byte, size)
copy(out, in)
return out
}
func printJSON(v any) error {
enc := json.NewEncoder(commandOutput)
enc.SetIndent("", " ")
return enc.Encode(v)
}
func usage() {
fmt.Fprintln(os.Stderr, "usage:")
fmt.Fprintln(os.Stderr, " v2singleblock smoke --path <file> [--name N --node N --write-text TEXT]")
fmt.Fprintln(os.Stderr, " v2singleblock restart-smoke --path <file> [--name N --node N --write-text TEXT]")
fmt.Fprintln(os.Stderr, " v2singleblock iscsi-smoke --path <file> [--name N --node N --iqn IQN]")
}
+129
View File
@@ -0,0 +1,129 @@
package main
import (
"bytes"
"encoding/json"
"path/filepath"
"strings"
"testing"
)
func TestRunSmoke_ProducesSingleNodeMVPResult(t *testing.T) {
path := filepath.Join(t.TempDir(), "single-node-mvp.blk")
var out bytes.Buffer
oldOutput := commandOutput
commandOutput = &out
defer func() {
commandOutput = oldOutput
}()
runErr := runSmoke([]string{
"--path", path,
"--name", "mvp-vol",
"--node", "node-a",
"--write-text", "hello-v2",
})
if runErr != nil {
t.Fatalf("runSmoke: %v", runErr)
}
var result smokeResult
if err := json.Unmarshal(out.Bytes(), &result); err != nil {
t.Fatalf("unmarshal output %q: %v", out.String(), err)
}
if result.VolumeName != "mvp-vol" {
t.Fatalf("volume_name=%q", result.VolumeName)
}
if result.Role != "primary" {
t.Fatalf("role=%q", result.Role)
}
if result.Mode != "allocated_only" {
t.Fatalf("mode=%q", result.Mode)
}
if !strings.Contains(result.Readback, "68656c6c6f2d7632") {
t.Fatalf("readback_hex=%q", result.Readback)
}
}
func TestRunRestartSmoke_PreservesDataAcrossRestart(t *testing.T) {
path := filepath.Join(t.TempDir(), "single-node-restart.blk")
var out bytes.Buffer
oldOutput := commandOutput
commandOutput = &out
defer func() {
commandOutput = oldOutput
}()
runErr := runRestartSmoke([]string{
"--path", path,
"--name", "restart-vol",
"--node", "node-a",
"--write-text", "hello-restart",
})
if runErr != nil {
t.Fatalf("runRestartSmoke: %v", runErr)
}
var result restartSmokeResult
if err := json.Unmarshal(out.Bytes(), &result); err != nil {
t.Fatalf("unmarshal output %q: %v", out.String(), err)
}
if result.VolumeName != "restart-vol" {
t.Fatalf("volume_name=%q", result.VolumeName)
}
if result.Role != "primary" {
t.Fatalf("role=%q", result.Role)
}
if result.Mode != "allocated_only" {
t.Fatalf("mode=%q", result.Mode)
}
if result.InitialWrite != result.PostRestart {
t.Fatalf("initial=%q post_restart=%q", result.InitialWrite, result.PostRestart)
}
if !strings.Contains(result.PostRestart, "68656c6c6f2d72657374617274") {
t.Fatalf("post_restart_hex=%q", result.PostRestart)
}
}
func TestRunISCSISmoke_ExportsFrontendAndReportsAddress(t *testing.T) {
path := filepath.Join(t.TempDir(), "single-node-iscsi.blk")
var out bytes.Buffer
oldOutput := commandOutput
commandOutput = &out
defer func() {
commandOutput = oldOutput
}()
runErr := runISCSISmoke([]string{
"--path", path,
"--name", "iscsi-vol",
"--node", "node-a",
"--iqn", "iqn.2026-04.com.seaweedfs:test.cli",
})
if runErr != nil {
t.Fatalf("runISCSISmoke: %v", runErr)
}
var result iscsiSmokeResult
if err := json.Unmarshal(out.Bytes(), &result); err != nil {
t.Fatalf("unmarshal output %q: %v", out.String(), err)
}
if result.VolumeName != "iscsi-vol" {
t.Fatalf("volume_name=%q", result.VolumeName)
}
if result.Role != "primary" {
t.Fatalf("role=%q", result.Role)
}
if result.Mode != "allocated_only" {
t.Fatalf("mode=%q", result.Mode)
}
if result.IQN != "iqn.2026-04.com.seaweedfs:test.cli" {
t.Fatalf("iqn=%q", result.IQN)
}
if !strings.Contains(result.Address, "127.0.0.1:") {
t.Fatalf("address=%q", result.Address)
}
}
@@ -0,0 +1,339 @@
# V2 Automata Ownership Map
Date: 2026-04-05
Status: active
## Purpose
This note translates the two-loop protocol into automata ownership.
The goal is to answer a practical question for the next implementation step:
1. which decisions belong to `masterv2`
2. which decisions belong to the primary-side data-control automata
3. which current V2 events and commands already fit that split
4. which surfaces should remain compressed projections only
## Top-Level Split
There are four different layers of meaning:
1. identity authority
2. identity evidence
3. data-control truth
4. outward compressed projection
They must not be merged.
## Operating Model
One volume should be treated as a micro-cluster:
1. `masterv2` is outside the micro-cluster and grants identity authority
2. the selected primary is inside the micro-cluster and owns data-control truth
3. replicas contribute bounded evidence and execution progress to the primary
This means `masterv2` is allowed to authorize takeover, but not to act as the
continuous recovery planner.
```mermaid
flowchart TD
master[masterv2]
assignment[AssignmentChannel]
heartbeat[HeartbeatChannel]
query[PromotionQueryChannel]
primary[PrimaryDataControlAutomata]
replica[ReplicaDataControlAutomata]
projection[OutwardProjection]
master --> assignment
assignment --> primary
assignment --> replica
primary --> projection
replica --> projection
projection --> heartbeat
master --> query
query --> primary
query --> replica
primary --> replica
replica --> primary
```
## Channel Ownership
### 1. Heartbeat
Heartbeat is not a recovery-planning channel.
Heartbeat should answer only:
1. is the node alive
2. what role and epoch has been applied
3. what compressed mode is outwardly visible
4. whether replica receiver readiness is present
Heartbeat should not answer:
1. which LSN to catch up to
2. whether a reconnect gap is recoverable
3. which source to rebuild from
4. per-replica durable progress
### 2. Promotion Query
Promotion query is the only `masterv2` path that may request fresh promotion
evidence.
It should answer only:
1. what is your current `CommittedLSN`
2. what is your current `WALHeadLSN`
3. which `Epoch` do you recognize
4. what `Role` do you believe you have
5. are you eligible to become primary
Promotion query should not answer:
1. replay plan details
2. full session graph
3. retention budget internals
4. detailed catch-up progress
### 3. Assignment
Assignment is the identity authorization channel.
It should answer only:
1. who is primary
2. who is replica
3. which epoch is now active
4. which replica identities and addresses belong to the set
5. whether rebuilding role is assigned
Assignment should not answer:
1. committed boundary
2. durable boundary
3. catch-up target LSN
4. detailed recovery plan
### 3a. Takeover Authorization Versus Takeover Choreography
This distinction should stay explicit:
1. `masterv2` may authorize a replacement primary
2. the replacement primary must reconstruct bounded truth itself
3. the replacement primary must decide whether takeover is safe, degraded, or
rebuild-only
So:
1. takeover authorization belongs to `Assignment`
2. takeover reconstruction belongs to `Loop 2`
3. detailed recovery sequencing belongs to the primary-led automata
### 4. Loop 2 Data Control
Loop 2 is where data-control truth lives.
It should answer:
1. what is committed
2. what is durable
3. whether barrier lineage is healthy
4. whether a replica is in `keepup`, `catchup`, `degraded`, or `needs_rebuild`
5. whether rebuild is currently the only safe path
## Existing V2 Events Mapped To Owners
### Identity-Control Entry Event
These belong to assignment application or identity evidence:
1. `AssignmentDelivered`
2. `RoleApplied`
3. `ReceiverReadyObserved`
Interpretation:
- `AssignmentDelivered` enters from `Assignment`
- `RoleApplied` is local confirmation after identity application
- `ReceiverReadyObserved` is local evidence that may affect compressed outward mode
### Data-Control Events
These belong to the primary-led data-control automata:
1. `ShipperConfiguredObserved`
2. `ShipperConnectedObserved`
3. `DiagnosticShippedAdvanced`
4. `CommittedLSNAdvanced`
5. `BarrierAccepted`
6. `BarrierRejected`
7. `CheckpointAdvanced`
8. `CatchUpPlanned`
9. `RecoveryProgressObserved`
10. `CatchUpCompleted`
11. `NeedsRebuildObserved`
12. `RebuildStarted`
13. `RebuildCommitted`
Interpretation:
- these events must stay inside `Loop 2`
- they may influence outward `Mode`
- they must not be surfaced to `masterv2` as raw automata state
## Existing V2 Commands Mapped To Owners
### Identity Commands
These belong to assignment realization:
1. `ApplyRoleCommand`
2. `StartReceiverCommand`
These are still executed locally, but they exist to realize identity control.
### Data-Control Commands
These belong to the primary-led data-control automata:
1. `ConfigureShipperCommand`
2. `StartRecoveryTaskCommand`
3. `DrainRecoveryTaskCommand`
4. `StartCatchUpCommand`
5. `StartRebuildCommand`
6. `InvalidateSessionCommand`
Interpretation:
- these commands must not be emitted by `masterv2`
- these commands are consequences of data-control truth
### Projection Command
`PublishProjectionCommand` is neither identity truth nor data-control truth.
It is the compressed outward surface emitted from the core.
This command exists so:
1. heartbeat can expose compressed state
2. debug and product surfaces can stay consistent
3. runtime seams do not invent independent meanings
## Logic Judgments
This is the decision table that should guide future changes.
### Judgment: Is this enough for heartbeat?
If a fact is needed only to answer:
1. alive
2. applied role
3. outward mode
then it belongs in heartbeat projection.
If it is needed to answer:
1. which node is safest to promote
2. whether recovery is catch-up or rebuild
3. how far durability actually advanced
then it does not belong in normal heartbeat.
### Judgment: Does master need this continuously?
If a fact affects:
1. assignment ownership
2. lease ownership
3. fencing
4. failover authorization
then `masterv2` may need it.
If a fact affects:
1. replay planning
2. retention decisions
3. barrier lineage
4. rebuild execution
then `masterv2` must not own it continuously.
### Judgment: Is master authorizing or choreographing?
If the action is:
1. choose a legal owner
2. fence stale owners
3. publish epoch and replica identity
then it belongs to `masterv2`.
If the action is:
1. reconstruct bounded truth from summaries
2. decide `keepup` vs `catchup`
3. select rebuild source and progress target
4. gate takeover on degraded or ambiguous lineage
then it belongs to the selected primary, not `masterv2`.
### Judgment: Does this belong to promotion query?
A fact belongs to promotion query when all are true:
1. it is needed only during promotion arbitration
2. stale cached heartbeat would be unsafe
3. the fact is still smaller than full recovery truth
Examples:
1. `CommittedLSN`
2. `WALHeadLSN`
3. `Epoch`
4. promotion eligibility
### Judgment: Does this belong to Loop 2?
A fact belongs to `Loop 2` when it changes at write/barrier/reconnect scale or
when it directly affects recovery planning.
Examples:
1. `ReplicaFlushedLSN`
2. `CatchUpTarget`
3. retention floor
4. session invalidation
5. recovery progress
## What Should Change Next
Before any major implementation, these changes should be reflected in code
structure:
1. `masterv2` types should split heartbeat, query, and assignment surfaces
2. `volumev2` control session should stop pretending heartbeat and promotion
evidence are the same message shape
3. engine-facing runtime seams should explicitly tag which observations are
identity evidence and which are data-control observations
4. tests should separately prove:
- heartbeat convergence
- promotion query freshness
- data-control progress closure
## Non-Goals
This note does not freeze every field of the future wire protocol.
It only fixes:
1. ownership
2. judgment boundaries
3. which existing V2 automata pieces are already valid
+474
View File
@@ -0,0 +1,474 @@
# V2 Capability Map
Date: 2026-04-05
Status: active
Purpose: define the V2 capability expansion map that drives feature closure, test closure, and the transition from bounded scenario debugging to systematic product validation
## Why This Document Exists
If `V2` is a real system line, it needs more than:
1. accepted protocol truths
2. passing point fixes
3. a few successful scenarios
It also needs one explicit map that answers:
1. what product capabilities exist in the V2 line
2. in what order those capabilities should close
3. what "done" means for each capability
4. which tests prove the capability
5. which proofs are V2-owned versus runtime-specific
This document is that map.
It complements:
1. `v2-protocol-truths.md` for stable semantic rules
2. `v2-product-completion-overview.md` for product-level completion status
3. `v2-phase-development-plan.md` for active execution sequencing
4. `v2_scenarios.md` for scenario backlog and historical failure sources
## How To Use This Map
For any new feature, bug fix, or test expansion, ask:
1. which capability tier does this belong to
2. which closure claim does it strengthen
3. which proof tier should carry it
4. whether it is V2-owned truth or current-runtime integration
This prevents three common failures:
1. growing V2 by random scenario accumulation
2. confusing `weed` integration success with V2 semantic completion
3. re-testing everything from zero when the runtime boundary changes later
## Core Method
The map uses three linked ideas:
### 1. Capability expansion
V2 should expand from:
1. single-volume correctness
2. bounded RF=2 replication
3. failover and rejoin
4. multi-replica behavior
5. lifecycle operations
6. control-plane and operations closure
7. CSI and product-surface closure
### 2. Completion definition
A capability is not "done" because code exists.
It is only closed when all of these are true:
1. semantic rule is explicit
2. runtime path exists
3. observability exists
4. focused tests prove the rule
5. one product-level scenario proves the real path
### 3. Proof layering
Each capability should be proven across four proof tiers:
1. `Core semantic`
- pure V2 truth
- fastest feedback
- should remain reusable if runtime changes
2. `Seam / adapter`
- queue, heartbeat, registry, proto, assignment, bridge ownership
- catches most integrated bugs cheaply
3. `Integrated runtime`
- real `weed` path today
- smaller number of high-value scenarios
4. `Soak / benchmark / adversarial`
- slow, broad, or disturbance-heavy validation
- not the daily development loop
## Capability Tiers
## Tier 0: Semantic Foundation
Goal:
1. make V2 the source of truth for replication semantics
Main closure claims:
1. epoch and lineage are authoritative
2. committed truth is explicit
3. catch-up versus rebuild boundary is explicit
4. stale authority fails closed
5. replica identity is stable across endpoint change
Done means:
1. truths are explicit in `v2-protocol-truths.md`
2. engine events and commands preserve those truths
3. core tests cover replay, stale events, fencing, and recovery choice
Primary proof tiers:
1. core semantic
2. seam only where identity/transport adaptation matters
Typical tests:
1. event -> projection -> command tests
2. recovery-choice tests
3. stale session / stale epoch rejection
4. stable `ReplicaID` versus mutable endpoint tests
## Tier 1: Single-Volume Base Capability
Goal:
1. prove one volume is correct before adding replication
Capabilities:
1. create/delete
2. single-node read/write
3. restart durability
4. publication correctness
5. bounded observability
Done means:
1. RF=1 write/read survives restart
2. publication reflects the true serving node
3. explicit health/publication state is observable
Primary proof tiers:
1. core semantic for boundaries
2. integrated runtime for real read/write/restart
Typical scenarios:
1. create -> write -> restart -> read
2. publication remains coherent after restart
## Tier 2: RF=2 Replication Base
Goal:
1. close the smallest useful HA replication unit
Capabilities:
1. primary/replica assignment
2. receiver readiness
3. shipper configuration
4. barrier semantics
5. explicit `publish_healthy`
6. explicit `degraded`
7. explicit `needs_rebuild`
Done means:
1. replica membership reaches the primary truthfully
2. `sync_all` cannot succeed vacuously with zero shippers
3. publication health depends on real closure, not optimistic state
4. RF=2 replicated write/read works on the integrated path
Primary proof tiers:
1. core semantic
2. seam
3. one integrated replicated IO scenario
Typical tests:
1. assignment-delivered membership tests
2. `RoleApplied`, `ReceiverReady`, `ShipperConfigured` closure tests
3. barrier strictness tests
4. replicated checksum scenarios
## Tier 3: RF=2 Recovery And Failover
Goal:
1. turn RF=2 replication into a fault-tolerant runtime path
Capabilities:
1. manual promote
2. auto failover
3. old primary fencing
4. old primary rejoin
5. catch-up-first reconnect
6. rebuild fallback
7. data continuity after failover
Done means:
1. promotion bumps epoch and fences stale authority
2. promoted primary regains replica membership after rejoin
3. reconnect chooses catch-up or rebuild explicitly
4. failover preserves committed data
5. one data-verified integrated scenario exists for each supported failover path
Primary proof tiers:
1. seam
2. integrated runtime
3. soak/adversarial for disturbance variants
Current note:
1. manual promote on the integrated `weed` path has now closed with data continuity verification
2. this tier remains broader than one passing scenario and still requires systematic matrix expansion
Typical scenarios:
1. kill primary -> promote replica -> restart old primary -> data verified
2. lease-expiry auto failover
3. rejoin with address change
4. rebuild fallback when catch-up path is unavailable
## Tier 4: Multi-Replica Runtime (`RF>=3`)
Goal:
1. extend the model from one replica to a replica set
Capabilities:
1. multi-replica membership
2. multi-shipper convergence
3. strict `sync_all`
4. `sync_quorum`
5. partial failure tolerance
6. replacement and rebuild target choice
Done means:
1. primary ownership and closure remain replica-scoped, not scalar-only
2. quorum/all durability rules hold under mixed replica states
3. failover and rejoin do not collapse back to RF=2-only assumptions
Primary proof tiers:
1. core semantic
2. seam
3. targeted integrated RF=3 scenarios
Typical tests:
1. multi-replica assignment closure
2. quorum durability tests
3. partial-failure promotion eligibility tests
4. RF=3 disturbance scenarios
## Tier 5: Lifecycle Capability
Goal:
1. prove that product operations remain correct under replication and recovery
Capabilities:
1. expand
2. truncate
3. snapshot
4. snapshot export/import
5. clone/restore style flows where supported
Done means:
1. lifecycle operations preserve V2 recovery truth
2. lifecycle metadata does not bypass fencing or recovery boundaries
3. lifecycle operations continue to hold under restart/failover
Primary proof tiers:
1. core semantic for boundary rules
2. seam where command ownership matters
3. integrated scenarios for user-visible lifecycle behavior
Typical scenarios:
1. snapshot then failover
2. expand under replicated volume
3. truncate under degraded or catch-up conditions
## Tier 6: Control Plane And Operations
Goal:
1. make the system diagnosable and operationally trustworthy
Capabilities:
1. heartbeat convergence
2. assignment queue correctness
3. registry truth coherence
4. publication truth coherence
5. debug surfaces
6. metrics and operator diagnosis
7. restart and disturbance policy clarity
Done means:
1. the control plane reports the same truth the runtime acts on
2. major failure classes are diagnosable from bounded logs/debug state
3. restart/rejoin behavior is policy-shaped, not accidental
Primary proof tiers:
1. seam
2. integrated runtime
3. soak for repeated disturbance
Typical tests:
1. registry/publication coherence tests
2. assignment queue confirm/refresh tests
3. reconnect/restart diagnosis tests
4. bounded failover observability tests
## Tier 7: Product Surfaces (`CSI`, `iSCSI`, `NVMe`)
Goal:
1. project V2 storage truth through real product interfaces
Capabilities:
1. volume create/publish through `CSI`
2. node stage/node publish
3. failover-visible remount or reconnect behavior
4. expansion through product surface
5. snapshot through product surface
6. front-end publication coherence
Done means:
1. product surfaces do not hide or weaken V2 truth
2. frontend publication follows actual authority after failover
3. product workflows survive supported restart/failover envelopes
Primary proof tiers:
1. seam
2. integrated runtime
3. slower end-to-end scenario pack
Typical scenarios:
1. CSI create/publish/write/failover/read
2. CSI expand under replicated volume
3. snapshot + restore + failover
## Tier 8: Launch Envelope
Goal:
1. convert bounded capability proof into a bounded support statement
Capabilities:
1. supported topology matrix
2. supported disturbance matrix
3. known unsupported branches
4. pilot stop conditions
5. rollout review evidence
Done means:
1. supported claims are explicit
2. unsupported areas are explicit
3. pilot and rollout review use the same capability map and proof layers
Primary proof tiers:
1. integrated runtime
2. soak / perf / operational review
## Capability Map Summary
| Tier | Scope | What closes here | Main proof emphasis |
|------|-------|------------------|---------------------|
| 0 | Semantic foundation | truth rules and fail-closed boundaries | core semantic |
| 1 | Single-volume base | RF=1 correctness and restart durability | core + integrated |
| 2 | RF=2 replication | receiver/shipper/barrier/publication closure | core + seam + one integrated path |
| 3 | RF=2 recovery/failover | promote, rejoin, catch-up, rebuild, data continuity | seam + integrated |
| 4 | RF>=3 runtime | multi-replica membership and durability semantics | core + seam + targeted integrated |
| 5 | Lifecycle | snapshot/expand/truncate under replication truth | mixed by feature |
| 6 | Control/ops | registry/heartbeat/publication/diagnosis closure | seam + integrated |
| 7 | Product surfaces | CSI and frontend projection of V2 truth | integrated |
| 8 | Launch envelope | bounded support and rollout claims | integrated + soak |
## Test Expansion Strategy From This Map
This map should drive testing in a faster order than "one expensive scenario at a time."
### Fast lane
Run on most code changes:
1. core semantic tests for the touched rule
2. seam tests for ingress/egress/control delivery
3. one focused scenario only if the change crosses a real product seam
### Medium lane
Run on milestone closure for a tier:
1. representative integrated scenarios for that tier
2. checksum or historical-read validation where data continuity matters
### Slow lane
Run on nightly or bounded review:
1. disturbance matrix
2. soak
3. benchmark
4. larger product-surface packs
## What Must Stay Runtime-Agnostic
To avoid re-testing everything from zero when `weed` ownership shrinks later,
these proof categories must stay V2-owned:
1. assignment semantics
2. role/epoch/fencing semantics
3. recovery-choice semantics
4. publication closure semantics
5. data continuity contracts
The current `weed` path remains valuable as:
1. the present integrated runtime
2. one proof backend for product-level behavior
It must not become the only place where V2 truth is tested.
## Immediate Next Use
This map should be used to produce:
1. one capability-to-test taxonomy
2. one current coverage matrix marking which tiers are:
- `strong`
- `bounded`
- `partial`
- `not yet closed`
3. one reduced high-value integrated scenario pack aligned to tiers rather than ad hoc bug history
## Current Practical Reading
For near-term work, read in this order:
1. `v2-protocol-truths.md`
2. `v2-capability-map.md`
3. `v2-product-completion-overview.md`
4. `v2-phase-development-plan.md`
5. `v2_scenarios.md`
+143
View File
@@ -0,0 +1,143 @@
# V2 Kernel Closure Review
Date: 2026-04-05
Status: active
## Question
The goal is not to prove whether iSCSI itself can be implemented. The reusable
`blockvol` + frontend code already shows that.
The real question is whether the current kernel split can grow into a product:
1. Is the brain owned by V2 semantics?
2. Is the control plane owned by V2 messages and convergence?
3. Is the data plane attached as an execution/backend service instead of a truth
owner?
## Current Answer
The current `masterv2 + volumev2 + purev2` shape is viable as a product kernel
because ownership is split in the right direction.
### Brain
Owner:
- `sw-block/engine/replication/`
What it owns:
1. semantic state
2. event ingestion
3. command intent
4. outward projection
What it must not own:
1. backend I/O
2. transport lifecycle
3. frontend serving details
### Control Plane
Owner:
- `sw-block/runtime/masterv2/`
- `sw-block/runtime/volumev2/control_session.go`
- `sw-block/runtime/volumev2/orchestrator.go`
What it owns:
1. desired declaration
2. heartbeat observation
3. assignment emission
4. assignment apply loop
5. convergence/idempotence
What it must not own:
1. WAL/extent execution
2. frontend protocol implementation
### Data Plane
Owner:
- `sw-block/runtime/purev2/`
- `sw-block/runtime/volumev2/frontend.go`
- reused `weed/storage/blockvol/*`
What it owns:
1. create/open
2. read/write/flush
3. restart durability
4. frontend export such as iSCSI
What it must not own:
1. role truth
2. publication truth
3. assignment policy
## Closure Proofs
Two small closure proofs are enough for the current stage.
### 1. Control-plane closure
Scenario:
1. `masterv2` declares one RF1 primary
2. `volumev2` heartbeats with no local role yet
3. `masterv2` emits an assignment
4. `volumev2` applies it through the V2 path
5. a later heartbeat converges to quiet state
6. if desired state changes, assignment is reissued once and converges again
Why it matters:
- this proves the new head is not piggybacking on `weed/server` loops
### 2. Data-plane closure
Scenario:
1. `volumev2` exports a named volume through iSCSI
2. a client logs in and issues SCSI write/read
3. data is verified through the frontend and local backend view
Why it matters:
- this proves the kernel can host a real frontend while keeping truth ownership
outside the frontend/backend code
## Product Meaning
If these two closures stay true while features expand, then the architecture can
scale toward:
1. RF1 productized single-node block service
2. RF2/RF3 replication as additional control/data workflows
3. failover and rebuild without moving semantic truth back into backend code
4. CSI on top of a clearer runtime contract
Another way to state the same result:
1. `masterv2` behaves like an external identity authority
2. each `volumev2` instance behaves like a per-volume micro-cluster shell
3. the selected primary inside that shell owns data-control truth and recovery
choreography
## Main Risk
The main risk is not iSCSI or local I/O. The main risk is semantic leakage:
1. adding more backend-state shortcuts into control decisions
2. letting frontend/backend code redefine publication truth
3. rebuilding `weed/server` ownership inside `volumev2`
4. letting `masterv2` grow from identity authority into a centralized recovery
planner
As long as those three are resisted, the kernel can keep expanding cleanly.
+253
View File
@@ -0,0 +1,253 @@
# V2 Loop1 Surface Draft
Date: 2026-04-05
Status: active
## Purpose
This note turns the two-loop design into a code-facing draft for the current
`masterv2` and `volumev2` packages.
It does not implement the new protocol yet.
It defines the smallest surface refactor that should happen first.
## Goal
Replace the current mixed `heartbeat -> assignments` MVP surface with three
separate `Loop 1` surfaces:
1. periodic heartbeat
2. promotion query
3. assignment
## Current State
Today, `masterv2` uses one heartbeat type and one assignment type:
- `sw-block/runtime/masterv2/master.go`
- `sw-block/runtime/volumev2/control_session.go`
- `sw-block/runtime/volumev2/volume.go`
The current heartbeat still mixes:
1. applied identity
2. outward projection naming
3. fields that should later become failover evidence
## Target Surfaces
### 1. Heartbeat
Purpose:
1. liveness
2. applied identity
3. outward compressed mode
Recommended fields:
- `NodeID`
- `ReportedAt`
- per-volume:
- `Name`
- `Path`
- `Epoch`
- `Role`
- `Mode`
- `ModeReason`
- passive `CommittedLSN` cache
- `RoleApplied`
- `ReplicaReady`
Not included:
- per-replica progress
- catch-up target
- detailed recovery phase
If `CommittedLSN` is present here, it is only a convenience cache for
`masterv2`. Fresh promotion judgment still uses the query channel.
### 2. Promotion Query
Purpose:
1. fresh promotion evidence
2. failover arbitration
Recommended request:
- `VolumeName`
- `ExpectedEpoch`
Recommended response:
- `VolumeName`
- `NodeID`
- `Epoch`
- `Role`
- `CommittedLSN`
- `WALHeadLSN`
- `ReceiverReady`
- `Eligible`
- `Reason`
Rule:
- `CommittedLSN` is the primary selection key
- `WALHeadLSN` is only a tiebreaker
### 3. Assignment
Purpose:
1. authorize role
2. fence stale owners
3. deliver member identity
Recommended fields:
- `Name`
- `Path`
- `NodeID`
- `Epoch`
- `LeaseTTL`
- `Role`
- `ReplicaSet`
- `CreateOptions`
`ReplicaSet` should carry identity plus transport addresses, not progress.
## Field Migration From Current MVP
### Current `masterv2.VolumeHeartbeat`
Today:
- `Name`
- `Path`
- `Epoch`
- `Role`
- `ProjectionMode`
- `PublicationReason`
- `RoleApplied`
Should become:
- `Name`
- `Path`
- `Epoch`
- `Role`
- `Mode`
- `ModeReason`
- `RoleApplied`
- `ReplicaReady`
Interpretation change:
- `ProjectionMode` should be renamed to `Mode`
- `PublicationReason` should stop pretending to be a generic reason field
- `ReplicaReady` should be explicit on the identity surface
### Current `masterv2.VolumeView`
Today:
- `ObservedEpoch`
- `ObservedRole`
- `ProjectionMode`
- `PublicationReason`
- `RoleApplied`
Should become:
- `ObservedEpoch`
- `ObservedRole`
- `Mode`
- `ModeReason`
- `RoleApplied`
- `ReplicaReady`
Optional later:
- bounded cached failover evidence for debugging only
### Current `volumev2.Node.Heartbeat()`
Current source:
- `snap.Status`
- `snap.Projection`
Recommended extraction:
- `Mode` from `snap.Projection.Mode.Name`
- `ModeReason` from `snap.Projection.Mode.Reason`
- passive `CommittedLSN` cache from local status snapshot
- `RoleApplied` from `snap.Projection.Readiness.RoleApplied`
- `ReplicaReady` from `snap.Projection.Readiness.ReplicaReady`
Not from heartbeat:
- `CatchUpTarget`
- `RecoveryProgress`
`CommittedLSN` may appear in heartbeat as a passive cache only.
Fresh promotion authority still belongs to the promotion-query channel.
`CatchUpTarget` and `RecoveryProgress` belong to Loop 2.
## Code Refactor Order
### Step 1
Add a small shared contract package for `Loop 1` types, for example under:
- `sw-block/runtime/protocolv2/`
Start with:
1. `Heartbeat`
2. `Assignment`
3. `PromotionQueryRequest`
4. `PromotionQueryResponse`
### Step 2
Make `masterv2` use the shared `Loop 1` contract types instead of local ad hoc
duplicates.
### Step 3
Make `volumev2` heartbeat generation write to the narrowed heartbeat shape.
### Step 4
Add a promotion-query interface without implementing full failover yet.
The first code slice only needs:
1. request and response types
2. local evidence extraction helper
3. one focused test proving query returns fresh local evidence
## Test Guidance
The first code refactor should keep existing control-loop tests and add one new
test:
1. existing heartbeat-assignment convergence should still pass
2. new promotion-query test should prove:
- query is separate from heartbeat
- fresh state is returned at call time
- heartbeat `CommittedLSN` is only a passive cache
- `WALHeadLSN` is not part of periodic heartbeat
## Non-Goals
This draft does not yet define:
1. full Loop 2 message schema
2. full new primary truth reconstruction choreography
3. quorum-specific `CommittedLSN` algorithm
Those come after the `Loop 1` surface is cleanly split.
@@ -0,0 +1,186 @@
# V2 Proof And Retest Pyramid
Date: 2026-04-05
Status: active
## Purpose
This note defines how the pure V2 runtime should accumulate proof so that most
closure stays in fast, reusable tests instead of expensive mixed scenarios.
It also defines where reuse is allowed to keep narrow regression coverage and
where truth-owner changes require full V2 retesting.
## Proof Pyramid
### Layer 1: core semantic tests
Owner:
- `sw-block/engine/replication/`
This layer proves:
1. assignment semantics
2. role-application semantics
3. readiness/publication semantics
4. stale/replay/idempotence behavior
5. recovery and fencing meaning
Rule:
- if a behavior changes V2 protocol truth, this layer is mandatory
### Layer 2: component and seam tests
Owner:
- `sw-block/runtime/purev2/`
- narrow bridge/dispatcher seams reused from `weed/`
This layer proves:
1. static assignment ingestion
2. dispatcher to backend binding
3. local role application feedback into the core
4. local restart/open/replay behavior
5. debug/projection cache consistency
Rule:
- this is the default daily development surface for the new runtime
### Layer 3: minimal integrated runtime tests
Owner:
- pure runtime binary/package smoke
This layer proves only:
1. create -> assign -> write -> read
2. restart -> reopen -> read
3. status/debug snapshot visibility
Rule:
- keep this pack tiny and fast
- do not pull failover, multi-node, or CSI into this layer
### Layer 4: compatibility and oracle checks
Owner:
- current `weed` integrated path
This layer remains useful for:
1. regression oracle coverage
2. parity comparison
3. late acceptance confidence
Rule:
- it is not the daily semantic development surface for pure V2
## Reuse Versus Full Retest
### Reuse with narrow regression
The following areas remain execution muscles and should keep focused coverage:
1. WAL append and flush mechanics
2. checkpoint and dirty-map mechanics
3. extent install mechanics
4. local read/write I/O
5. transport/frontend plumbing when semantic meaning is unchanged
Typical proof:
1. component tests
2. mechanical regression tests
3. one narrow smoke if the seam changed
### Full V2 retest required
The following areas change truth ownership and therefore require explicit V2 proof:
1. assignment meaning
2. publication meaning
3. readiness closure
4. fencing and lineage meaning
5. recovery classification
6. operator-visible health and degraded semantics
Typical proof:
1. core semantic tests first
2. pure-runtime component tests second
3. only then one small integrated confirmation
## Stage Gates
### Stage A: RF1 pure shell
Must be green before any RF2 work starts:
1. create/open works
2. static primary assignment works
3. local write/read works
4. restart durability works
5. debug/projection snapshot works
### Stage B: RF2 replication base
May start only after Stage A is closed.
This stage adds:
1. replica membership truth
2. receiver/shipper wiring
3. barrier semantics
4. degraded versus publish closure
5. catch-up and rebuild base semantics
### Stage C: failover and rejoin
May start only after RF2 replication base is closed.
This stage adds:
1. manual promote
2. auto failover
3. rejoin
4. recovery ownership handoff
### Stage D: product surfaces
May start only after failover/rejoin closure exists.
This stage adds:
1. CSI
2. operator APIs
3. external readiness and health surfaces
4. acceptance and soak packs
## Daily Working Rule
When a new bug appears, classify it first:
1. core semantic bug
2. pure-runtime seam bug
3. execution-muscle bug
4. compatibility-only oracle bug
Then choose the cheapest proof tier that still matches the truth owner.
If the answer is "mixed scenario first", the classification is probably still
too vague.
## Related References
- `v2-pure-runtime-rf1-bootstrap.md`
- `v2-capability-map.md`
- `v2-reuse-replacement-boundary.md`
- `v2-legacy-runtime-exit-criteria.md`
@@ -0,0 +1,148 @@
# Pure V2 RF1 Bootstrap
Date: 2026-04-05
Status: active
## Purpose
This note turns the bootstrap plan into a concrete runtime boundary that can be
implemented and tested without going through `weed/server`.
The goal of the first executable slice is not feature completeness.
The goal is to establish one small, closed, V2-owned runtime that can accumulate
semantic truth and local execution proof without mixed-runtime distortion.
## Implemented Boundary
### Pure V2 shell
The first pure runtime shell lives in:
- `sw-block/runtime/purev2/`
It owns:
- local process/runtime lifecycle
- local block volume registration
- V2 core engine ownership
- local projection/debug cache
- static RF1 assignment injection
It does not own:
- master heartbeat loops
- assignment queues
- failover control
- recovery task orchestration
- CSI or product-facing control APIs
### Reused execution muscles
The pure shell reuses mechanics instead of re-implementing them:
- `weed/storage/store_blockvol.go`
- `weed/storage/blockvol/blockvol.go`
- `weed/storage/blockvol/v2bridge/command_bindings.go`
- `weed/server/blockcmd/`
The intended rule is:
- `sw-block/runtime/purev2` owns local runtime closure
- `sw-block/engine/replication` owns semantic authority
- reused `weed/` pieces execute concrete backend actions only
### Small executable entrypoint
The first operator-facing entrypoint lives in:
- `sw-block/cmd/purev2rf1/`
Current commands:
1. `bootstrap`
2. `status`
This is intentionally small.
It exists to make the first slice executable and inspectable, not to define the
final product surface.
## RF1 First Slice Contract
The first slice is closed only if all of the following are true.
### Included
1. create one local block volume
2. open an existing local block volume
3. inject one static RF1 primary assignment
4. apply local role through V2 command dispatch
5. perform real local read/write through `blockvol`
6. survive restart and preserve data
7. expose explicit projection/debug state
### Explicitly excluded
1. replica membership truth
2. receiver/shipper wiring as a required path
3. catch-up or rebuild orchestration
4. manual or auto failover
5. CSI
If a new requirement needs any excluded item, it belongs to a later tier and
must not be forced into the RF1 shell.
## Runtime Shape
The implemented ownership split is:
```text
purev2 runtime
-> V2 core engine
-> dispatcher
-> command bindings
-> blockvol store/backend
-> projection/debug snapshot
```
The critical closure path is:
```text
static assignment
-> core ApplyEvent(AssignmentDelivered)
-> dispatcher executes apply_role
-> runtime feeds back RoleApplied
-> projection is cached explicitly
-> local debug snapshot becomes inspectable
```
## Current Behavioral Meaning
For the current engine semantics, an RF1 primary with zero replicas remains:
1. locally writable as a block volume
2. explicitly visible in debug/projection state
3. projected as `allocated_only`, not `publish_healthy`
That is acceptable for this first bootstrap slice because:
1. the purpose is shell closure, not final RF1 publication semantics
2. publication meaning remains explicit instead of being guessed from local role
3. RF1 publication policy can evolve later without changing the shell boundary
## Immediate Engineering Rules
While the pure runtime remains in the RF1 stage:
1. new shell work goes into `sw-block/runtime/purev2/`
2. pure runtime must not import `weed/server/volume_server_block.go`
3. new proof should prefer unit/component tests in the pure shell
4. mixed `weed` scenario results remain oracle coverage, not semantic authority
## Related References
- `v2-capability-map.md`
- `v2-reuse-replacement-boundary.md`
- `v2-legacy-runtime-exit-criteria.md`
- `v2-protocol-truths.md`
+414
View File
@@ -0,0 +1,414 @@
# V2 Two-Loop Protocol
Date: 2026-04-05
Status: active
## Purpose
This note fixes the protocol boundary for the next V2 step.
The goal is not to finalize every wire field before implementation.
The goal is to make the ownership boundary stable enough that automata,
constraints, and runtime packages can be reorganized without mixing identity
control and replication consensus again.
## Core Rule
The protocol is split into two loops:
1. `Loop 1`: identity control
2. `Loop 2`: data control
These loops must not be collapsed into one heartbeat or one state owner.
## Authority Principle
Each volume should be treated as a small distributed cluster:
1. `masterv2` is the identity authority outside the cluster
2. the selected primary is the data-control authority inside the cluster
3. replicas report bounded facts to the primary, not full truth to `masterv2`
This means:
1. `masterv2` decides who is allowed to own the role
2. the new primary decides how takeover, catch-up, and rebuild proceed
3. `masterv2` may query bounded facts for arbitration, but it does not choreograph
data recovery step by step
## Loop 1: Identity Control
Owner:
- `masterv2 <-> volumev2`
Frequency:
- low
- heartbeat scale
- assignment scale
- promotion-query scale
## Three Control Channels
Within `Loop 1`, the control plane should be split into three different
channels. They must not be collapsed into one message type.
### 1. Heartbeat
Direction:
- `volumev2 -> masterv2`
Frequency:
- periodic
- lightweight
Purpose:
1. liveness detection
2. compressed outward mode
3. confirmation that assignment was applied
Heartbeat should carry only:
1. `NodeID`
2. per-volume `Mode`
3. applied `Epoch`
4. applied `Role`
5. `RoleApplied`
6. `ReplicaReady`
7. optional passive `CommittedLSN` cache
Heartbeat should not carry:
1. per-replica progress
2. catch-up targets
3. rebuild detail
4. full failover evidence
If `CommittedLSN` is carried in heartbeat, it is only a passive cache. The
authoritative failover-time value still comes from promotion query.
### 2. Promotion Query
Direction:
- `masterv2 -> candidate volumev2`
- candidate `volumev2 -> masterv2`
Frequency:
- on demand
- only during failover or promotion arbitration
Purpose:
1. obtain fresh failover evidence
2. avoid treating stale heartbeat cache as authority
Candidate response should include:
1. `CommittedLSN`
2. `WALHeadLSN` as a weaker tiebreaker
3. `Epoch`
4. `Role`
5. `ReceiverReady`
6. bounded eligibility reason if not promotable
The promotion query is where fresh identity-loop evidence is collected.
It is not a replacement for the data-control loop, and it is not a continuous
replication-progress feed.
If heartbeat also carries `CommittedLSN`, promotion query still wins whenever
fresh arbitration is required.
### 3. Assignment
Direction:
- `masterv2 -> volumev2`
Frequency:
- on demand
- role change or membership change
Purpose:
1. authorize role ownership
2. fence stale owners
3. deliver replica-set identity
### Master To Volume
`masterv2 -> volumev2` carries only:
1. `Epoch`
2. `Role`
3. `LeaseTTL`
4. `ReplicaSet` identities and addresses
It does not carry:
1. per-replica progress
2. catch-up target history
3. detailed rebuild plan
### Volume To Master
`volumev2 -> masterv2` carries only bounded identity evidence:
1. applied `Epoch`
2. applied `Role`
3. outward `Mode`
4. `RoleApplied`
5. `ReplicaReady`
This channel must remain lightweight. Fresh failover evidence belongs to the
promotion-query channel, not the periodic heartbeat.
## Loop 2: Data Control
Owner:
- `primary engine <-> replica engine`
Frequency:
- high
- write scale
- barrier scale
- reconnect scale
This is where replication consensus lives.
### Primary To Replica
Primary-side data-control messages should cover:
1. WAL entry stream
2. barrier request with `Epoch` and target LSN
3. reconnect or resume handshake
4. rebuild/catch-up execution requests when needed
### Replica To Primary
Replica-side data-control messages should cover:
1. `FlushedLSN`
2. bounded status such as `ok`, `epoch_mismatch`, `timeout`, `fsync_failed`
3. reconnect gap evidence
4. coarse local recovery state
### Primary-Owned Per-Replica State
The primary brain should own:
1. `ReplicaFlushedLSN`
2. `ShippedLSN` as diagnostic only
3. replica `State`
4. `CatchUpTarget`
5. `RetentionFloor`
6. `LastContactTime`
## Role Of Masterv2
`masterv2` authorizes:
1. who is primary
2. who is replica
3. which epoch is active
4. when stale owners must be fenced
`masterv2` must not decide:
1. replay from LSN `X` to `Y`
2. whether the next action is `keepup` or `catchup`
3. how rebuild is executed
4. continuous commit progress
`masterv2` may query candidates for fresh promotion evidence, but it still does
not become the owner of replication history.
## Reconstruction And Takeover
Promotion and reconstruction are related, but they do not have the same owner.
### What Masterv2 Leads
`masterv2` leads:
1. failover detection
2. epoch fencing
3. candidate query for fresh promotion evidence
4. primary selection
5. assignment of the new primary role
### What The New Primary Leads
The selected replacement primary leads:
1. local assignment realization
2. collection of self and peer replica summaries
3. bounded truth reconstruction
4. fail-closed takeover gating
5. follow-on `keepup`, `catchup`, or `rebuild` orchestration
### Rule
`masterv2` may say:
1. "you are now the authorized primary candidate for epoch `E`"
2. "these are the members of the replica set"
But `masterv2` must not say:
1. "replay from LSN `X` to `Y`"
2. "use replica `R` as the rebuild source"
3. "enter `catchup` before `rebuild`"
4. "the cluster is safe because my last cached view looked healthy"
## Role Of Primary Brain
The primary brain discovers:
1. current replica state
2. gap or retention situation
3. barrier success or failure
4. whether the volume is `keepup`, `catchup`, `degraded`, or `needs_rebuild`
The primary brain decides:
1. keep shipping
2. start catch-up
3. escalate to `needs_rebuild`
4. start rebuild after role/assignment allows it
## Distributed State-Machine Rules
### 1. Different Nodes Have Different Views
Each node must distinguish:
1. local execution truth
2. last observed peer truth
3. cluster identity truth from `masterv2`
Do not collapse these into one blob.
### 2. All Peer Observations Are Epoch-Scoped
Any peer observation that affects recovery must be tied to:
1. `Epoch`
2. session or generation token
Old-epoch observations must be ignored or fail closed.
### 3. New Primary Reconstructs Truth
After failover, the new primary must rebuild its own data-control truth from:
1. local state
2. peer summaries
3. reconnect handshakes
It must not trust `masterv2` as a cache of full recovery history.
It may use `masterv2` only as the source of authorization and replica identity.
### 4. Outward Mode Is Compressed Evidence
`allocated_only`, `bootstrap_pending`, `publish_healthy`, `degraded`, and
`needs_rebuild` are public meanings, not the full internal recovery automaton.
### 5. Ambiguity Fails Closed
When barrier lineage, progress lineage, or epoch lineage is unclear, the system
must prefer:
1. `degraded`
2. `needs_rebuild`
3. no promotion without enough eligibility evidence
## Constraint Migration
Most of the last week's V2 work remains valid. The important change is where
each constraint belongs.
### Keep As-Is
These constraints still stand:
1. epoch fencing
2. one active session per replica per epoch
3. `catchup` and `rebuild` are different paths
4. fail closed on ambiguous recovery truth
5. semantics first, adapters later
### Move To Loop 1
These belong to identity control:
1. assignment application
2. role ownership
3. lease ownership
4. stable replica identity and addressing
5. compressed heartbeat evidence
6. on-demand promotion query for fresh evidence
### Move To Loop 2
These belong to data control:
1. committed/durable/checkpoint boundaries
2. barrier result meaning
3. replica progress
4. keepup/catchup/rebuild progression
5. retention-floor and catch-up targeting
## Existing V2 Seeds To Reuse
The current V2 code already has the right seeds for the primary-led loop:
1. `sw-block/engine/replication/state.go`
2. `sw-block/engine/replication/event.go`
3. `sw-block/engine/replication/command.go`
4. `sw-block/engine/replication/sender.go`
5. `sw-block/engine/replication/session.go`
6. `sw-block/engine/replication/registry.go`
The current MVP already has the right seeds for the identity loop:
1. `sw-block/runtime/masterv2/master.go`
2. `sw-block/runtime/volumev2/control_session.go`
3. `sw-block/runtime/volumev2/orchestrator.go`
## Immediate Next Step
Before deeper implementation, the codebase should next define:
1. the minimal `Loop 1` contract types in code, split into heartbeat, promotion
query, and assignment
2. the minimal `Loop 2` progress and reconnect contract draft
3. the automata ownership map showing which engine events and commands belong
to identity control versus data control
## Promotion Logic
Promotion should use fresh on-demand evidence, not stale heartbeat cache.
Recommended judgment order:
1. fence the old primary by epoch
2. query all surviving candidates
3. reject any candidate with wrong epoch, wrong role lineage, or not-ready
receiver state
4. choose the candidate with highest `CommittedLSN`
5. use `WALHeadLSN` only as a tiebreaker for equally committed candidates
6. assign new primary role at a new epoch
This keeps the durability boundary centered on `CommittedLSN`, which is the
last LSN that satisfied the configured durability mode such as `sync_all` or
`sync_quorum`.
@@ -0,0 +1,151 @@
# V2 VolumeV2 Single-Node MVP
Date: 2026-04-05
Status: active
## Purpose
This note defines the target shape for a single-node `volumev2` MVP that can
ship as a normal block service before HA/failover exists.
The core idea is:
1. `masterv2` is fully new control ownership
2. `volumev2` is a new shell and brain host
3. `blockvol` and related backend mechanics remain reusable muscles
## Target Layering
`volumev2` should be strengthened around four layers.
### 1. Engine
Owner:
- `sw-block/engine/replication/`
Responsibility:
1. state
2. event ingestion
3. command emission
4. outward projection
Rule:
- semantic truth lives here
- no backend I/O or network ownership
### 2. Engine Interface
Owner:
- command/event vocabulary between control/runtime and backend execution
Responsibility:
1. assignment -> event translation
2. observation -> event translation
3. command -> execution dispatch contract
Rule:
- runtime shell may not mutate engine truth directly
### 3. Control Plane
Owner:
- `masterv2 <-> volumev2` coordination
Responsibility:
1. node identity
2. registration and heartbeat
3. assignment receipt
4. state reporting
5. future recovery-control vocabulary (`keepup`, `catchup`, `rebuild`)
Rule:
- control plane carries protocol messages
- it does not own local data execution
### 4. Data Plane
Owner:
- local storage and serving mechanics
Responsibility:
1. WAL/extent management
2. read/write/flush
3. background workers
4. receiver/shipper mechanics
5. NVMe/iSCSI/frontend serving
Rule:
- data plane knows how to execute
- it does not define publication or role semantics
## Single-Node MVP Contract
The first ship-capable `volumev2` slice should include:
1. `masterv2` declaration of one RF1 primary volume
2. `volumev2` control session to fetch assignments
3. local create/open through reused `blockvol`
4. local primary assignment application through the V2 engine
5. local read/write plus restart durability
6. debug/status snapshot
7. one small executable entrypoint for smoke usage
The first slice explicitly excludes:
1. failover
2. RF2 replication
3. catch-up/rebuild ownership
4. CSI
## Why This Is Enough
This is enough to prove:
1. the `masterv2 + volumev2` head is viable
2. `volumev2` can host V2 semantics while reusing V1 muscles
3. a useful non-HA block service can exist before HA complexity is added
## Module Shape
Recommended package split:
1. `sw-block/runtime/masterv2/`
2. `sw-block/runtime/volumev2/`
3. `sw-block/runtime/purev2/`
4. `sw-block/engine/replication/`
5. `sw-block/bridge/blockvol/`
Within `volumev2`, strengthen toward:
1. `control_session.go`
2. `orchestrator.go`
3. `node.go`
4. later: `heartbeat.go`, `frontend.go`, `workers.go`
## Stage Gate
`volumev2` may be treated as a single-node MVP only when:
1. assignment sync is repeatable and idempotent
2. local IO is data-verified
3. restart/open path is proven
4. status/debug state is explicit
5. no `weed/server` lifecycle owner is required
## Related References
- `v2-pure-runtime-rf1-bootstrap.md`
- `v2-proof-and-retest-pyramid.md`
- `v2-capability-map.md`
+264
View File
@@ -0,0 +1,264 @@
package masterv2
import (
"fmt"
"slices"
"sync"
"time"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/protocolv2"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol"
)
// Config defines the minimal masterv2 control-loop settings.
type Config struct {
LeaseTTL time.Duration
}
// VolumeSpec is one master-owned desired volume intent.
type VolumeSpec struct {
Name string
Path string
PrimaryNodeID string
CreateOptions blockvol.CreateOptions
}
// Assignment is the minimal masterv2 -> volumev2 control message.
type Assignment = protocolv2.Assignment
// VolumeHeartbeat is the minimal volumev2 -> masterv2 periodic identity observation.
type VolumeHeartbeat = protocolv2.VolumeHeartbeat
// NodeHeartbeat is the minimal per-node heartbeat used by the POC control loop.
type NodeHeartbeat = protocolv2.NodeHeartbeat
// PromotionQueryRequest asks one candidate for fresh failover evidence.
type PromotionQueryRequest = protocolv2.PromotionQueryRequest
// PromotionQueryResponse returns fresh candidate evidence at query time.
type PromotionQueryResponse = protocolv2.PromotionQueryResponse
// VolumeView is the master-side observed state for one desired volume.
type VolumeView struct {
Name string
Path string
PrimaryNodeID string
DesiredEpoch uint64
ObservedEpoch uint64
ObservedRole string
Mode string
ModeReason string
CommittedLSN uint64
RoleApplied bool
ReplicaReady bool
LastHeartbeatAt time.Time
}
type desiredVolume struct {
spec VolumeSpec
epoch uint64
}
// Master is a small in-process V2 control plane POC.
// It owns desired state and emits assignments on heartbeats.
type Master struct {
mu sync.RWMutex
cfg Config
desired map[string]desiredVolume
views map[string]VolumeView
}
// New creates the minimal masterv2 POC.
func New(cfg Config) *Master {
if cfg.LeaseTTL <= 0 {
cfg.LeaseTTL = 30 * time.Second
}
return &Master{
cfg: cfg,
desired: make(map[string]desiredVolume),
views: make(map[string]VolumeView),
}
}
// DeclarePrimary records one desired RF1 primary placement.
func (m *Master) DeclarePrimary(spec VolumeSpec) error {
if spec.Name == "" {
return fmt.Errorf("masterv2: volume name is required")
}
if spec.Path == "" {
return fmt.Errorf("masterv2: volume path is required")
}
if spec.PrimaryNodeID == "" {
return fmt.Errorf("masterv2: primary node id is required")
}
m.mu.Lock()
defer m.mu.Unlock()
entry, ok := m.desired[spec.Name]
if !ok {
entry = desiredVolume{epoch: 1}
} else if entry.spec.Path != spec.Path || entry.spec.PrimaryNodeID != spec.PrimaryNodeID || entry.spec.CreateOptions != spec.CreateOptions {
entry.epoch++
}
entry.spec = spec
m.desired[spec.Name] = entry
view := m.views[spec.Name]
view.Name = spec.Name
view.Path = spec.Path
view.PrimaryNodeID = spec.PrimaryNodeID
view.DesiredEpoch = entry.epoch
m.views[spec.Name] = view
return nil
}
// HandleHeartbeat consumes one node heartbeat and returns any assignments the node should apply.
func (m *Master) HandleHeartbeat(hb NodeHeartbeat) ([]Assignment, error) {
if hb.NodeID == "" {
return nil, fmt.Errorf("masterv2: heartbeat node id is required")
}
if hb.ReportedAt.IsZero() {
hb.ReportedAt = time.Now()
}
m.mu.Lock()
defer m.mu.Unlock()
byName := make(map[string]VolumeHeartbeat, len(hb.Volumes))
for _, vol := range hb.Volumes {
byName[vol.Name] = vol
view := m.views[vol.Name]
view.Name = vol.Name
view.Path = vol.Path
view.ObservedEpoch = vol.Epoch
view.ObservedRole = vol.Role
view.Mode = vol.Mode
view.ModeReason = vol.ModeReason
view.CommittedLSN = vol.CommittedLSN
view.RoleApplied = vol.RoleApplied
view.ReplicaReady = vol.ReplicaReady
view.LastHeartbeatAt = hb.ReportedAt
if desired, ok := m.desired[vol.Name]; ok {
view.PrimaryNodeID = desired.spec.PrimaryNodeID
view.DesiredEpoch = desired.epoch
}
m.views[vol.Name] = view
}
var assignments []Assignment
names := make([]string, 0, len(m.desired))
for name := range m.desired {
names = append(names, name)
}
slices.Sort(names)
for _, name := range names {
desired := m.desired[name]
if desired.spec.PrimaryNodeID != hb.NodeID {
continue
}
reported, ok := byName[name]
if ok && reported.Path == desired.spec.Path && reported.Epoch == desired.epoch && reported.Role == "primary" && reported.RoleApplied {
continue
}
assignments = append(assignments, m.primaryAssignment(desired))
}
return assignments, nil
}
// Volume returns the latest master-side view for one volume.
func (m *Master) Volume(name string) (VolumeView, bool) {
m.mu.RLock()
defer m.mu.RUnlock()
view, ok := m.views[name]
return view, ok
}
// SelectPromotionCandidate chooses the best eligible candidate from fresh
// promotion-query responses. Selection is durability-first: highest
// CommittedLSN wins, then WALHeadLSN as a weaker tiebreaker.
func (m *Master) SelectPromotionCandidate(responses []PromotionQueryResponse) (PromotionQueryResponse, error) {
candidates := make([]PromotionQueryResponse, 0, len(responses))
for _, resp := range responses {
if resp.Eligible {
candidates = append(candidates, resp)
}
}
if len(candidates) == 0 {
return PromotionQueryResponse{}, fmt.Errorf("masterv2: no eligible promotion candidates")
}
slices.SortStableFunc(candidates, func(a, b PromotionQueryResponse) int {
if a.CommittedLSN != b.CommittedLSN {
if a.CommittedLSN > b.CommittedLSN {
return -1
}
return 1
}
if a.WALHeadLSN != b.WALHeadLSN {
if a.WALHeadLSN > b.WALHeadLSN {
return -1
}
return 1
}
switch {
case a.NodeID < b.NodeID:
return -1
case a.NodeID > b.NodeID:
return 1
default:
return 0
}
})
return candidates[0], nil
}
// AuthorizePromotion selects the best candidate from fresh promotion evidence,
// advances desired ownership when needed, and returns the assignment the chosen
// node should apply. This is authorization only; takeover reconstruction and
// activation remain the new primary's responsibility.
func (m *Master) AuthorizePromotion(volumeName string, responses []PromotionQueryResponse) (Assignment, error) {
if volumeName == "" {
return Assignment{}, fmt.Errorf("masterv2: volume name is required")
}
selected, err := m.SelectPromotionCandidate(responses)
if err != nil {
return Assignment{}, err
}
if selected.VolumeName != "" && selected.VolumeName != volumeName {
return Assignment{}, fmt.Errorf("masterv2: promotion candidate volume %q does not match %q", selected.VolumeName, volumeName)
}
m.mu.Lock()
defer m.mu.Unlock()
desired, ok := m.desired[volumeName]
if !ok {
return Assignment{}, fmt.Errorf("masterv2: unknown volume %q", volumeName)
}
if desired.spec.PrimaryNodeID != selected.NodeID {
desired.spec.PrimaryNodeID = selected.NodeID
desired.epoch++
m.desired[volumeName] = desired
}
view := m.views[volumeName]
view.Name = desired.spec.Name
view.Path = desired.spec.Path
view.PrimaryNodeID = desired.spec.PrimaryNodeID
view.DesiredEpoch = desired.epoch
m.views[volumeName] = view
return m.primaryAssignment(desired), nil
}
func (m *Master) primaryAssignment(desired desiredVolume) Assignment {
return Assignment{
Name: desired.spec.Name,
Path: desired.spec.Path,
NodeID: desired.spec.PrimaryNodeID,
Epoch: desired.epoch,
LeaseTTL: m.cfg.LeaseTTL,
CreateOptions: desired.spec.CreateOptions,
Role: "primary",
}
}
+138
View File
@@ -0,0 +1,138 @@
package masterv2
import "testing"
func TestSelectPromotionCandidate_PrefersHighestCommittedLSN(t *testing.T) {
master := New(Config{})
selected, err := master.SelectPromotionCandidate([]PromotionQueryResponse{
{
VolumeName: "vol-a",
NodeID: "node-b",
CommittedLSN: 12,
WALHeadLSN: 20,
Eligible: true,
},
{
VolumeName: "vol-a",
NodeID: "node-a",
CommittedLSN: 15,
WALHeadLSN: 16,
Eligible: true,
},
})
if err != nil {
t.Fatalf("select candidate: %v", err)
}
if selected.NodeID != "node-a" {
t.Fatalf("selected node=%q, want node-a", selected.NodeID)
}
}
func TestSelectPromotionCandidate_UsesWalHeadAsTiebreaker(t *testing.T) {
master := New(Config{})
selected, err := master.SelectPromotionCandidate([]PromotionQueryResponse{
{
VolumeName: "vol-a",
NodeID: "node-a",
CommittedLSN: 15,
WALHeadLSN: 16,
Eligible: true,
},
{
VolumeName: "vol-a",
NodeID: "node-b",
CommittedLSN: 15,
WALHeadLSN: 19,
Eligible: true,
},
})
if err != nil {
t.Fatalf("select candidate: %v", err)
}
if selected.NodeID != "node-b" {
t.Fatalf("selected node=%q, want node-b", selected.NodeID)
}
}
func TestSelectPromotionCandidate_RejectsWhenNoEligibleCandidates(t *testing.T) {
master := New(Config{})
_, err := master.SelectPromotionCandidate([]PromotionQueryResponse{
{NodeID: "node-a", Eligible: false, Reason: "needs_rebuild"},
{NodeID: "node-b", Eligible: false, Reason: "epoch_mismatch"},
})
if err == nil {
t.Fatal("expected error when no eligible candidates")
}
}
func TestAuthorizePromotion_ReassignsPrimaryAndAdvancesEpoch(t *testing.T) {
master := New(Config{})
if err := master.DeclarePrimary(VolumeSpec{
Name: "vol-a",
Path: "/tmp/vol-a.blk",
PrimaryNodeID: "node-a",
}); err != nil {
t.Fatalf("declare primary: %v", err)
}
assign, err := master.AuthorizePromotion("vol-a", []PromotionQueryResponse{
{
VolumeName: "vol-a",
NodeID: "node-b",
CommittedLSN: 18,
WALHeadLSN: 19,
Eligible: true,
},
{
VolumeName: "vol-a",
NodeID: "node-c",
CommittedLSN: 17,
WALHeadLSN: 20,
Eligible: true,
},
})
if err != nil {
t.Fatalf("authorize promotion: %v", err)
}
if assign.NodeID != "node-b" {
t.Fatalf("assignment node=%q, want node-b", assign.NodeID)
}
if assign.Epoch != 2 {
t.Fatalf("assignment epoch=%d, want 2", assign.Epoch)
}
view, ok := master.Volume("vol-a")
if !ok {
t.Fatal("master view missing")
}
if view.PrimaryNodeID != "node-b" {
t.Fatalf("view primary=%q, want node-b", view.PrimaryNodeID)
}
if view.DesiredEpoch != 2 {
t.Fatalf("view desired epoch=%d, want 2", view.DesiredEpoch)
}
}
func TestAuthorizePromotion_RejectsMismatchedVolume(t *testing.T) {
master := New(Config{})
if err := master.DeclarePrimary(VolumeSpec{
Name: "vol-a",
Path: "/tmp/vol-a.blk",
PrimaryNodeID: "node-a",
}); err != nil {
t.Fatalf("declare primary: %v", err)
}
_, err := master.AuthorizePromotion("vol-a", []PromotionQueryResponse{
{
VolumeName: "vol-b",
NodeID: "node-b",
CommittedLSN: 18,
WALHeadLSN: 19,
Eligible: true,
},
})
if err == nil {
t.Fatal("expected mismatched volume error")
}
}
+68
View File
@@ -0,0 +1,68 @@
package protocolv2
import (
"time"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol"
)
// ReplicaMember identifies one replica in the low-frequency identity loop.
// It carries identity and transport only, never replication progress.
type ReplicaMember struct {
NodeID string
DataAddr string
CtrlAddr string
}
// Assignment is the Loop 1 assignment surface from masterv2 to volumev2.
type Assignment struct {
Name string
Path string
NodeID string
Epoch uint64
LeaseTTL time.Duration
CreateOptions blockvol.CreateOptions
Role string
ReplicaSet []ReplicaMember
}
// VolumeHeartbeat is the periodic identity-loop heartbeat surface.
// It stays small and carries only applied identity, compressed outward mode,
// and a passive committed-boundary cache.
type VolumeHeartbeat struct {
Name string
Path string
Epoch uint64
Role string
Mode string
ModeReason string
CommittedLSN uint64
RoleApplied bool
ReplicaReady bool
}
// NodeHeartbeat is the low-frequency periodic heartbeat from a volume node.
type NodeHeartbeat struct {
NodeID string
ReportedAt time.Time
Volumes []VolumeHeartbeat
}
// PromotionQueryRequest asks one candidate for fresh failover evidence.
type PromotionQueryRequest struct {
VolumeName string
ExpectedEpoch uint64
}
// PromotionQueryResponse returns fresh candidate evidence at query time.
type PromotionQueryResponse struct {
VolumeName string
NodeID string
Epoch uint64
Role string
CommittedLSN uint64
WALHeadLSN uint64
ReceiverReady bool
Eligible bool
Reason string
}
+33
View File
@@ -0,0 +1,33 @@
package protocolv2
// ReplicaSummaryRequest asks one node for a bounded takeover/reconstruction
// summary for a specific volume. This is richer than promotion evidence, but
// still smaller than full internal engine/session state.
type ReplicaSummaryRequest struct {
VolumeName string
ExpectedEpoch uint64
}
// ReplicaSummaryResponse is the bounded summary a future primary can use to
// reconstruct recovery truth. It preserves distinct boundary semantics without
// exposing raw shipper/session internals.
type ReplicaSummaryResponse struct {
VolumeName string
NodeID string
Epoch uint64
Role string
Mode string
ModeReason string
RoleApplied bool
ReceiverReady bool
CommittedLSN uint64
DurableLSN uint64
CheckpointLSN uint64
TargetLSN uint64
AchievedLSN uint64
RecoveryPhase string
LastBarrierOK bool
LastBarrierReason string
Eligible bool
Reason string
}
+4
View File
@@ -0,0 +1,4 @@
// Package purev2 provides a small pure-V2 runtime shell that keeps
// semantic authority in the V2 engine while reusing blockvol execution
// mechanics. The initial slice is intentionally RF1/single-node only.
package purev2
+349
View File
@@ -0,0 +1,349 @@
package purev2
import (
"fmt"
"os"
"path/filepath"
"sync"
"time"
engine "github.com/seaweedfs/seaweedfs/sw-block/engine/replication"
"github.com/seaweedfs/seaweedfs/weed/server/blockcmd"
"github.com/seaweedfs/seaweedfs/weed/storage"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol/v2bridge"
)
const defaultListenAddr = "127.0.0.1:3260"
// Config defines the small pure V2 runtime shell configuration.
type Config struct {
ListenAddr string
AdvertisedHost string
DiskType string
}
// VolumeDebugSnapshot is the bounded outward runtime view for one volume.
type VolumeDebugSnapshot struct {
Path string
Info blockvol.VolumeInfo
Status blockvol.BlockVolumeStatus
Projection engine.PublicationProjection
HasProjection bool
CoreState engine.VolumeState
HasCoreState bool
ExecutedCommands []string
}
// Runtime is a small pure-V2 process shell for RF1 local execution.
// It intentionally excludes master heartbeat, failover, and product-surface loops.
type Runtime struct {
store *storage.BlockVolumeStore
bindings *v2bridge.CommandBindings
core *engine.CoreEngine
dispatcher *blockcmd.Dispatcher
config Config
mu sync.RWMutex
projections map[string]engine.PublicationProjection
executed map[string][]string
}
// New creates a new RF1-oriented pure V2 runtime shell.
func New(cfg Config) *Runtime {
if cfg.ListenAddr == "" {
cfg.ListenAddr = defaultListenAddr
}
rt := &Runtime{
store: storage.NewBlockVolumeStore(),
core: engine.NewCoreEngine(),
config: cfg,
projections: make(map[string]engine.PublicationProjection),
executed: make(map[string][]string),
}
rt.bindings = v2bridge.NewCommandBindings(rt.store, cfg.ListenAddr, cfg.AdvertisedHost)
rt.dispatcher = blockcmd.NewDispatcher(
blockcmd.NewServiceOps(runtimeBackend{runtime: rt}, nil, rt, nil),
blockcmd.NewHostEffects(rt.recordCommand, rt.emitCoreEvent, rt, rt),
)
return rt
}
// CreateVolume creates a new local block volume and registers it in the runtime store.
func (rt *Runtime) CreateVolume(path string, opts blockvol.CreateOptions, cfgs ...blockvol.BlockVolConfig) error {
if rt == nil {
return fmt.Errorf("purev2: runtime is nil")
}
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
return fmt.Errorf("purev2: create parent dir for %s: %w", path, err)
}
vol, err := blockvol.CreateBlockVol(path, opts, cfgs...)
if err != nil {
return fmt.Errorf("purev2: create volume %s: %w", path, err)
}
if err := vol.Close(); err != nil {
return fmt.Errorf("purev2: close created volume %s: %w", path, err)
}
_, err = rt.store.AddBlockVolume(path, rt.config.DiskType, cfgs...)
if err != nil {
return fmt.Errorf("purev2: register created volume %s: %w", path, err)
}
return nil
}
// OpenVolume opens an existing volume into the runtime store.
func (rt *Runtime) OpenVolume(path string, cfgs ...blockvol.BlockVolConfig) error {
if rt == nil {
return fmt.Errorf("purev2: runtime is nil")
}
if _, err := rt.store.AddBlockVolume(path, rt.config.DiskType, cfgs...); err != nil {
return fmt.Errorf("purev2: open volume %s: %w", path, err)
}
return nil
}
// ApplyPrimaryAssignment injects a static RF1 primary assignment into the pure runtime.
func (rt *Runtime) ApplyPrimaryAssignment(path string, epoch uint64, leaseTTL time.Duration) error {
return rt.applyAssignment(blockvol.BlockVolumeAssignment{
Path: path,
Epoch: epoch,
Role: blockvol.RoleToWire(blockvol.RolePrimary),
LeaseTtlMs: blockvol.LeaseTTLToWire(leaseTTL),
})
}
// BootstrapPrimary creates a volume if needed and applies a primary RF1 assignment.
func (rt *Runtime) BootstrapPrimary(path string, opts blockvol.CreateOptions, epoch uint64, leaseTTL time.Duration, cfgs ...blockvol.BlockVolConfig) error {
if _, ok := rt.store.GetBlockVolume(path); !ok {
if _, err := os.Stat(path); err == nil {
if err := rt.OpenVolume(path, cfgs...); err != nil {
return err
}
} else {
if err := rt.CreateVolume(path, opts, cfgs...); err != nil {
return err
}
}
}
return rt.ApplyPrimaryAssignment(path, epoch, leaseTTL)
}
// WriteLBA writes data at the given logical block address.
func (rt *Runtime) WriteLBA(path string, lba uint64, data []byte) error {
var status blockvol.V2StatusSnapshot
err := rt.store.WithVolume(path, func(vol *blockvol.BlockVol) error {
if err := vol.WriteLBA(lba, data); err != nil {
return err
}
status = vol.StatusSnapshot()
return nil
})
if err != nil {
return err
}
rt.observeLocalBoundaries(path, status)
return nil
}
// ReadLBA reads data at the given logical block address.
func (rt *Runtime) ReadLBA(path string, lba uint64, length uint32) ([]byte, error) {
var out []byte
err := rt.store.WithVolume(path, func(vol *blockvol.BlockVol) error {
data, err := vol.ReadLBA(lba, length)
if err != nil {
return err
}
out = data
return nil
})
return out, err
}
// SyncCache forces durable local flush through the reused blockvol backend.
func (rt *Runtime) SyncCache(path string) error {
var (
status blockvol.V2StatusSnapshot
attempt bool
)
err := rt.store.WithVolume(path, func(vol *blockvol.BlockVol) error {
attempt = true
if err := vol.SyncCache(); err != nil {
return err
}
status = vol.StatusSnapshot()
return nil
})
if err != nil {
if attempt {
rt.emitCoreEvent(engine.BarrierRejected{ID: path, Reason: err.Error()})
}
return err
}
rt.observeLocalBoundaries(path, status)
rt.emitCoreEvent(engine.BarrierAccepted{ID: path, FlushedLSN: status.CommittedLSN})
rt.emitCoreEvent(engine.CheckpointAdvanced{ID: path, CheckpointLSN: status.CheckpointLSN})
return nil
}
// Snapshot returns the bounded runtime and core view for one volume.
func (rt *Runtime) Snapshot(path string) (VolumeDebugSnapshot, error) {
if rt == nil {
return VolumeDebugSnapshot{}, fmt.Errorf("purev2: runtime is nil")
}
var snap VolumeDebugSnapshot
err := rt.store.WithVolume(path, func(vol *blockvol.BlockVol) error {
snap.Path = path
snap.Info = vol.Info()
snap.Status = vol.Status()
return nil
})
if err != nil {
return VolumeDebugSnapshot{}, err
}
if proj, ok := rt.Projection(path); ok {
snap.HasProjection = true
snap.Projection = proj
}
if st, ok := rt.core.State(path); ok {
snap.HasCoreState = true
snap.CoreState = st
}
rt.mu.RLock()
snap.ExecutedCommands = append([]string(nil), rt.executed[path]...)
rt.mu.RUnlock()
return snap, nil
}
// WithVolume exposes controlled access to one underlying block volume.
// It exists so higher-level runtimes can attach frontend adapters without
// importing weed/server lifecycle code.
func (rt *Runtime) WithVolume(path string, fn func(*blockvol.BlockVol) error) error {
if rt == nil {
return fmt.Errorf("purev2: runtime is nil")
}
return rt.store.WithVolume(path, fn)
}
// Projection implements blockcmd.ProjectionReader.
func (rt *Runtime) Projection(volumeID string) (engine.PublicationProjection, bool) {
rt.mu.RLock()
defer rt.mu.RUnlock()
proj, ok := rt.projections[volumeID]
return proj, ok
}
// StoreProjection implements blockcmd.ProjectionCacheWriter.
func (rt *Runtime) StoreProjection(volumeID string, projection engine.PublicationProjection) {
rt.mu.Lock()
defer rt.mu.Unlock()
rt.projections[volumeID] = projection
}
// Close closes all registered volumes.
func (rt *Runtime) Close() {
if rt == nil {
return
}
rt.store.Close()
}
func (rt *Runtime) applyAssignment(a blockvol.BlockVolumeAssignment) error {
ev, ok := assignmentEvent(a)
if !ok {
return nil
}
result := rt.core.ApplyEvent(ev)
return rt.dispatcher.Run(result.Commands, &a)
}
func (rt *Runtime) emitCoreEvent(ev engine.Event) {
result := rt.core.ApplyEvent(ev)
_ = rt.dispatcher.Run(result.Commands, nil)
}
func (rt *Runtime) recordCommand(volumeID, name string) {
rt.mu.Lock()
defer rt.mu.Unlock()
rt.executed[volumeID] = append(rt.executed[volumeID], name)
}
func (rt *Runtime) observeLocalBoundaries(path string, status blockvol.V2StatusSnapshot) {
rt.emitCoreEvent(engine.CommittedLSNAdvanced{
ID: path,
CommittedLSN: status.CommittedLSN,
})
}
func assignmentEvent(a blockvol.BlockVolumeAssignment) (engine.AssignmentDelivered, bool) {
ev := engine.AssignmentDelivered{
ID: a.Path,
Epoch: a.Epoch,
}
switch blockvol.RoleFromWire(a.Role) {
case blockvol.RolePrimary:
ev.Role = engine.RolePrimary
return ev, true
case blockvol.RoleReplica:
ev.Role = engine.RoleReplica
return ev, true
case blockvol.RoleRebuilding:
ev.Role = engine.RoleReplica
return ev, true
default:
return engine.AssignmentDelivered{}, false
}
}
type runtimeBackend struct {
runtime *Runtime
}
func (ops runtimeBackend) ApplyRole(assignment blockvol.BlockVolumeAssignment) (bool, error) {
if ops.runtime == nil || ops.runtime.bindings == nil {
return false, nil
}
if err := ops.runtime.bindings.ApplyRole(assignment); err != nil {
return false, err
}
ops.runtime.emitCoreEvent(engine.RoleApplied{ID: assignment.Path})
return true, nil
}
func (ops runtimeBackend) StartReceiver(assignment blockvol.BlockVolumeAssignment) (bool, error) {
if ops.runtime == nil || ops.runtime.bindings == nil {
return false, nil
}
if assignment.ReplicaDataAddr == "" || assignment.ReplicaCtrlAddr == "" {
return false, nil
}
if _, err := ops.runtime.bindings.StartReceiver(assignment.Path, assignment.ReplicaDataAddr, assignment.ReplicaCtrlAddr); err != nil {
return false, err
}
return true, nil
}
func (ops runtimeBackend) ConfigureShipper(volumeID string, replicas []engine.ReplicaAssignment) (bool, bool, error) {
if ops.runtime == nil || ops.runtime.bindings == nil || len(replicas) == 0 {
return false, false, nil
}
addrs := make([]blockvol.ReplicaAddr, 0, len(replicas))
for _, replica := range replicas {
if replica.Endpoint.DataAddr == "" || replica.Endpoint.CtrlAddr == "" {
continue
}
addrs = append(addrs, blockvol.ReplicaAddr{
ServerID: replica.ReplicaID,
DataAddr: replica.Endpoint.DataAddr,
CtrlAddr: replica.Endpoint.CtrlAddr,
})
}
if len(addrs) == 0 {
return false, false, nil
}
if _, err := ops.runtime.bindings.ConfigurePrimaryReplication(volumeID, addrs); err != nil {
return false, false, err
}
return true, ops.runtime.bindings.IsPrimaryShipperConnected(volumeID), nil
}
+154
View File
@@ -0,0 +1,154 @@
package purev2
import (
"bytes"
"path/filepath"
"reflect"
"testing"
"time"
engine "github.com/seaweedfs/seaweedfs/sw-block/engine/replication"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol"
)
func TestRuntime_RF1BootstrapCreateAssignSnapshot(t *testing.T) {
tempDir := t.TempDir()
rt := New(Config{})
defer rt.Close()
path := filepath.Join(tempDir, "rf1-primary.blk")
if err := rt.BootstrapPrimary(path, testCreateOptions(), 1, 30*time.Second); err != nil {
t.Fatalf("bootstrap primary: %v", err)
}
snap, err := rt.Snapshot(path)
if err != nil {
t.Fatalf("snapshot: %v", err)
}
if snap.Status.Role != blockvol.RolePrimary {
t.Fatalf("role=%v", snap.Status.Role)
}
if snap.Status.Epoch != 1 {
t.Fatalf("epoch=%d", snap.Status.Epoch)
}
if !snap.HasProjection {
t.Fatal("expected published projection")
}
if snap.Projection.Role != engine.RolePrimary {
t.Fatalf("projection role=%v", snap.Projection.Role)
}
if snap.Projection.Mode.Name != engine.ModeAllocatedOnly {
t.Fatalf("projection mode=%s", snap.Projection.Mode.Name)
}
if snap.Projection.Publication.Reason != "allocated_only" {
t.Fatalf("publication reason=%q", snap.Projection.Publication.Reason)
}
if !snap.Projection.Readiness.RoleApplied {
t.Fatal("role_applied should be true after apply_role")
}
if !reflect.DeepEqual(snap.ExecutedCommands, []string{"apply_role"}) {
t.Fatalf("executed=%v", snap.ExecutedCommands)
}
}
func TestRuntime_RF1RestartPreservesLocalData(t *testing.T) {
tempDir := t.TempDir()
path := filepath.Join(tempDir, "rf1-restart.blk")
payload := bytes.Repeat([]byte{0x5A}, 4096)
func() {
rt := New(Config{})
defer rt.Close()
if err := rt.BootstrapPrimary(path, testCreateOptions(), 1, 30*time.Second); err != nil {
t.Fatalf("bootstrap primary: %v", err)
}
if err := rt.WriteLBA(path, 0, payload); err != nil {
t.Fatalf("write: %v", err)
}
if err := rt.SyncCache(path); err != nil {
t.Fatalf("sync cache: %v", err)
}
}()
rt := New(Config{})
defer rt.Close()
if err := rt.OpenVolume(path); err != nil {
t.Fatalf("open volume: %v", err)
}
readBack, err := rt.ReadLBA(path, 0, uint32(len(payload)))
if err != nil {
t.Fatalf("read after restart: %v", err)
}
if !bytes.Equal(readBack, payload) {
t.Fatal("readback mismatch after restart")
}
if err := rt.ApplyPrimaryAssignment(path, 1, 30*time.Second); err != nil {
t.Fatalf("re-apply primary: %v", err)
}
snap, err := rt.Snapshot(path)
if err != nil {
t.Fatalf("snapshot after reopen: %v", err)
}
if !snap.HasCoreState {
t.Fatal("expected core state after re-apply")
}
if snap.CoreState.Readiness.RoleApplied != true {
t.Fatal("role_applied should remain true after reopen assignment")
}
}
func TestRuntime_LocalBoundaryObservationsAdvanceCoreState(t *testing.T) {
tempDir := t.TempDir()
rt := New(Config{})
defer rt.Close()
path := filepath.Join(tempDir, "rf1-boundaries.blk")
if err := rt.BootstrapPrimary(path, testCreateOptions(), 1, 30*time.Second); err != nil {
t.Fatalf("bootstrap primary: %v", err)
}
payload := bytes.Repeat([]byte{0x33}, 4096)
if err := rt.WriteLBA(path, 0, payload); err != nil {
t.Fatalf("write: %v", err)
}
snap, err := rt.Snapshot(path)
if err != nil {
t.Fatalf("snapshot after write: %v", err)
}
if !snap.HasCoreState {
t.Fatal("expected core state after write")
}
if snap.CoreState.Boundary.CommittedLSN == 0 {
t.Fatalf("committed_lsn=%d, want > 0", snap.CoreState.Boundary.CommittedLSN)
}
if snap.CoreState.Boundary.DurableLSN != 0 {
t.Fatalf("durable_lsn=%d before sync, want 0", snap.CoreState.Boundary.DurableLSN)
}
if err := rt.SyncCache(path); err != nil {
t.Fatalf("sync cache: %v", err)
}
snap, err = rt.Snapshot(path)
if err != nil {
t.Fatalf("snapshot after sync: %v", err)
}
if snap.CoreState.Boundary.DurableLSN < snap.CoreState.Boundary.CommittedLSN {
t.Fatalf("durable_lsn=%d committed_lsn=%d", snap.CoreState.Boundary.DurableLSN, snap.CoreState.Boundary.CommittedLSN)
}
if snap.CoreState.Boundary.CheckpointLSN > snap.CoreState.Boundary.DurableLSN {
t.Fatalf("checkpoint_lsn=%d durable_lsn=%d", snap.CoreState.Boundary.CheckpointLSN, snap.CoreState.Boundary.DurableLSN)
}
}
func testCreateOptions() blockvol.CreateOptions {
return blockvol.CreateOptions{
VolumeSize: 1 * 1024 * 1024,
BlockSize: 4096,
WALSize: 256 * 1024,
}
}
@@ -0,0 +1,45 @@
package volumev2
import (
"fmt"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/masterv2"
)
// ControlSession is the minimal control-plane contract between volumev2 and masterv2.
type ControlSession interface {
Heartbeat(masterv2.NodeHeartbeat) ([]masterv2.Assignment, error)
}
// PromotionEvidenceSource is the on-demand Loop 1 query surface used during
// failover arbitration.
type PromotionEvidenceSource interface {
QueryPromotionEvidence(masterv2.PromotionQueryRequest) (masterv2.PromotionQueryResponse, error)
}
// InProcessSession is the first in-process control-plane adapter used by the MVP.
type InProcessSession struct {
master *masterv2.Master
}
// NewInProcessSession creates a control session backed by one in-process masterv2.
func NewInProcessSession(master *masterv2.Master) (*InProcessSession, error) {
if master == nil {
return nil, fmt.Errorf("volumev2: master is nil")
}
return &InProcessSession{master: master}, nil
}
// Heartbeat sends one periodic heartbeat to masterv2 and returns the assignments to apply.
func (s *InProcessSession) Heartbeat(hb masterv2.NodeHeartbeat) ([]masterv2.Assignment, error) {
if s == nil || s.master == nil {
return nil, fmt.Errorf("volumev2: control session is nil")
}
return s.master.HandleHeartbeat(hb)
}
// Sync is retained as a narrow compatibility shim for existing tests while the
// three-channel Loop 1 surface settles.
func (s *InProcessSession) Sync(hb masterv2.NodeHeartbeat) ([]masterv2.Assignment, error) {
return s.Heartbeat(hb)
}
+83
View File
@@ -0,0 +1,83 @@
package volumev2
import (
"fmt"
"time"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/purev2"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol"
)
// DataPlane is the minimal single-node execution contract for volumev2.
// It keeps backend mechanics replaceable while volumev2 owns the control shell.
type DataPlane interface {
BootstrapPrimary(path string, opts blockvol.CreateOptions, epoch uint64, leaseTTL time.Duration) error
WriteLBA(path string, lba uint64, data []byte) error
ReadLBA(path string, lba uint64, length uint32) ([]byte, error)
SyncCache(path string) error
Snapshot(path string) (purev2.VolumeDebugSnapshot, error)
WithVolume(path string, fn func(*blockvol.BlockVol) error) error
Close()
}
// PureRuntimeDataPlane adapts the current purev2 runtime as a volumev2 data plane.
type PureRuntimeDataPlane struct {
runtime *purev2.Runtime
}
// NewPureRuntimeDataPlane wraps one purev2 runtime behind the DataPlane contract.
func NewPureRuntimeDataPlane(runtime *purev2.Runtime) (*PureRuntimeDataPlane, error) {
if runtime == nil {
return nil, fmt.Errorf("volumev2: pure runtime is nil")
}
return &PureRuntimeDataPlane{runtime: runtime}, nil
}
func (dp *PureRuntimeDataPlane) BootstrapPrimary(path string, opts blockvol.CreateOptions, epoch uint64, leaseTTL time.Duration) error {
if dp == nil || dp.runtime == nil {
return fmt.Errorf("volumev2: data plane is nil")
}
return dp.runtime.BootstrapPrimary(path, opts, epoch, leaseTTL)
}
func (dp *PureRuntimeDataPlane) WriteLBA(path string, lba uint64, data []byte) error {
if dp == nil || dp.runtime == nil {
return fmt.Errorf("volumev2: data plane is nil")
}
return dp.runtime.WriteLBA(path, lba, data)
}
func (dp *PureRuntimeDataPlane) ReadLBA(path string, lba uint64, length uint32) ([]byte, error) {
if dp == nil || dp.runtime == nil {
return nil, fmt.Errorf("volumev2: data plane is nil")
}
return dp.runtime.ReadLBA(path, lba, length)
}
func (dp *PureRuntimeDataPlane) SyncCache(path string) error {
if dp == nil || dp.runtime == nil {
return fmt.Errorf("volumev2: data plane is nil")
}
return dp.runtime.SyncCache(path)
}
func (dp *PureRuntimeDataPlane) Snapshot(path string) (purev2.VolumeDebugSnapshot, error) {
if dp == nil || dp.runtime == nil {
return purev2.VolumeDebugSnapshot{}, fmt.Errorf("volumev2: data plane is nil")
}
return dp.runtime.Snapshot(path)
}
func (dp *PureRuntimeDataPlane) WithVolume(path string, fn func(*blockvol.BlockVol) error) error {
if dp == nil || dp.runtime == nil {
return fmt.Errorf("volumev2: data plane is nil")
}
return dp.runtime.WithVolume(path, fn)
}
func (dp *PureRuntimeDataPlane) Close() {
if dp == nil || dp.runtime == nil {
return
}
dp.runtime.Close()
}
+305
View File
@@ -0,0 +1,305 @@
package volumev2
import (
"fmt"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/masterv2"
)
// FailoverParticipant is the minimal surface needed to execute one failover
// flow across Loop 1 authorization and Loop 2 takeover preparation.
type FailoverParticipant interface {
QueryPromotionEvidence(masterv2.PromotionQueryRequest) (masterv2.PromotionQueryResponse, error)
QueryReplicaSummarySource
PreparePrimaryTakeover(PrimaryTakeoverPlan) (ReconstructedPrimaryTruth, error)
GatePrimaryActivation(volumeName string, truth ReconstructedPrimaryTruth) error
}
// QueryReplicaSummarySource aliases the peer summary surface so the failover
// helper can reuse the existing bounded takeover contract.
type QueryReplicaSummarySource interface {
ReplicaSummarySource
}
// FailoverResult captures the outputs of one authorized and prepared failover.
type FailoverResult struct {
Candidate masterv2.PromotionQueryResponse
Assignment masterv2.Assignment
Truth ReconstructedPrimaryTruth
}
// FailoverStage is the coarse external progress marker for one failover
// session. It is intended for orchestration and debugging, not semantics.
type FailoverStage string
const (
FailoverStageNew FailoverStage = "new"
FailoverStageEvidenceCollected FailoverStage = "evidence_collected"
FailoverStageAuthorized FailoverStage = "authorized"
FailoverStagePrepared FailoverStage = "prepared"
FailoverStageActivated FailoverStage = "activated"
FailoverStageFailed FailoverStage = "failed"
)
// FailoverSnapshot is a read-only summary of the session's current observable
// state for drivers, tests, and debug surfaces.
type FailoverSnapshot struct {
VolumeName string
ExpectedEpoch uint64
Stage FailoverStage
LastError string
ResponseCount int
SelectedNodeID string
Result FailoverResult
}
// FailoverSession is a thin orchestration object that exposes the narrow
// failover stages explicitly so higher-level drivers can stop after
// authorization, inspect intermediate results, or run the whole sequence.
type FailoverSession struct {
master *masterv2.Master
volumeName string
expectedEpoch uint64
participants []FailoverParticipant
responses []masterv2.PromotionQueryResponse
byNode map[string]FailoverParticipant
result FailoverResult
stage FailoverStage
lastErr error
}
// NewFailoverSession validates the narrow failover inputs and returns a
// stepwise orchestration session.
func NewFailoverSession(master *masterv2.Master, volumeName string, expectedEpoch uint64, participants []FailoverParticipant) (*FailoverSession, error) {
if master == nil {
return nil, fmt.Errorf("volumev2: master is nil")
}
if volumeName == "" {
return nil, fmt.Errorf("volumev2: volume name is required")
}
if len(participants) == 0 {
return nil, fmt.Errorf("volumev2: failover participants are required")
}
return &FailoverSession{
master: master,
volumeName: volumeName,
expectedEpoch: expectedEpoch,
participants: participants,
stage: FailoverStageNew,
}, nil
}
// CollectPromotionEvidence gathers fresh promotion responses from all
// configured participants.
func (s *FailoverSession) CollectPromotionEvidence() ([]masterv2.PromotionQueryResponse, error) {
if s == nil {
return nil, fmt.Errorf("volumev2: failover session is nil")
}
responses := make([]masterv2.PromotionQueryResponse, 0, len(s.participants))
byNode := make(map[string]FailoverParticipant, len(s.participants))
for _, participant := range s.participants {
if participant == nil {
continue
}
resp, err := participant.QueryPromotionEvidence(masterv2.PromotionQueryRequest{
VolumeName: s.volumeName,
ExpectedEpoch: s.expectedEpoch,
})
if err != nil {
return nil, s.failf("volumev2: promotion evidence %s: %w", s.volumeName, err)
}
if resp.NodeID == "" {
return nil, s.failf("volumev2: promotion evidence for %s missing node id", s.volumeName)
}
responses = append(responses, resp)
byNode[resp.NodeID] = participant
}
if len(responses) == 0 {
return nil, s.failf("volumev2: no failover evidence collected for %s", s.volumeName)
}
s.responses = responses
s.byNode = byNode
s.stage = FailoverStageEvidenceCollected
s.lastErr = nil
return append([]masterv2.PromotionQueryResponse(nil), responses...), nil
}
// Authorize asks masterv2 to pick and authorize the new primary assignment from
// the collected fresh promotion evidence.
func (s *FailoverSession) Authorize() (masterv2.Assignment, error) {
if s == nil {
return masterv2.Assignment{}, fmt.Errorf("volumev2: failover session is nil")
}
if len(s.responses) == 0 {
if _, err := s.CollectPromotionEvidence(); err != nil {
return masterv2.Assignment{}, err
}
}
assignment, err := s.master.AuthorizePromotion(s.volumeName, s.responses)
if err != nil {
return masterv2.Assignment{}, s.fail(err)
}
s.result.Assignment = assignment
for _, resp := range s.responses {
if resp.NodeID == assignment.NodeID {
s.result.Candidate = resp
break
}
}
s.stage = FailoverStageAuthorized
s.lastErr = nil
return assignment, nil
}
// PrepareTakeover runs the selected node's bounded takeover reconstruction.
func (s *FailoverSession) PrepareTakeover() (ReconstructedPrimaryTruth, error) {
if s == nil {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: failover session is nil")
}
if s.result.Assignment.NodeID == "" {
if _, err := s.Authorize(); err != nil {
return ReconstructedPrimaryTruth{}, err
}
}
selected, ok := s.byNode[s.result.Assignment.NodeID]
if !ok {
return ReconstructedPrimaryTruth{}, s.failf("volumev2: authorized node %q missing participant", s.result.Assignment.NodeID)
}
peers := make([]ReplicaSummarySource, 0, len(s.byNode)-1)
for nodeID, participant := range s.byNode {
if nodeID == s.result.Assignment.NodeID {
continue
}
peers = append(peers, participant)
}
truth, err := selected.PreparePrimaryTakeover(PrimaryTakeoverPlan{
Assignment: s.result.Assignment,
Peers: peers,
})
s.result.Truth = truth
if err != nil {
return truth, s.fail(err)
}
s.stage = FailoverStagePrepared
s.lastErr = nil
return truth, nil
}
// Activate gates the selected primary on the reconstructed takeover truth.
func (s *FailoverSession) Activate() error {
if s == nil {
return fmt.Errorf("volumev2: failover session is nil")
}
if s.result.Assignment.NodeID == "" {
if _, err := s.Authorize(); err != nil {
return err
}
}
if s.result.Truth.PrimaryNodeID == "" {
if _, err := s.PrepareTakeover(); err != nil {
return err
}
}
selected, ok := s.byNode[s.result.Assignment.NodeID]
if !ok {
return s.failf("volumev2: authorized node %q missing participant", s.result.Assignment.NodeID)
}
if err := selected.GatePrimaryActivation(s.volumeName, s.result.Truth); err != nil {
return s.fail(err)
}
s.stage = FailoverStageActivated
s.lastErr = nil
return nil
}
// Result returns the latest collected candidate, assignment, and reconstructed
// takeover truth known to the session.
func (s *FailoverSession) Result() FailoverResult {
if s == nil {
return FailoverResult{}
}
return s.result
}
// Stage returns the current coarse failover stage.
func (s *FailoverSession) Stage() FailoverStage {
if s == nil {
return FailoverStageFailed
}
return s.stage
}
// LastError returns the last stage error observed by the session.
func (s *FailoverSession) LastError() error {
if s == nil {
return fmt.Errorf("volumev2: failover session is nil")
}
return s.lastErr
}
// Snapshot returns a stable read-only view of the session's externally useful
// state.
func (s *FailoverSession) Snapshot() FailoverSnapshot {
if s == nil {
return FailoverSnapshot{Stage: FailoverStageFailed, LastError: "volumev2: failover session is nil"}
}
lastErr := ""
if s.lastErr != nil {
lastErr = s.lastErr.Error()
}
return FailoverSnapshot{
VolumeName: s.volumeName,
ExpectedEpoch: s.expectedEpoch,
Stage: s.stage,
LastError: lastErr,
ResponseCount: len(s.responses),
SelectedNodeID: s.result.Assignment.NodeID,
Result: s.result,
}
}
// Run executes the full failover path from fresh evidence collection through
// activation gating.
func (s *FailoverSession) Run() (FailoverResult, error) {
if s == nil {
return FailoverResult{}, fmt.Errorf("volumev2: failover session is nil")
}
if _, err := s.CollectPromotionEvidence(); err != nil {
return s.Result(), err
}
if _, err := s.Authorize(); err != nil {
return s.Result(), err
}
if _, err := s.PrepareTakeover(); err != nil {
return s.Result(), err
}
if err := s.Activate(); err != nil {
return s.Result(), err
}
return s.Result(), nil
}
// ExecuteFailoverFlow runs the narrow failover path:
// fresh promotion evidence -> master authorization -> takeover preparation ->
// activation gate. It intentionally does not choreograph catch-up or rebuild.
func ExecuteFailoverFlow(master *masterv2.Master, volumeName string, expectedEpoch uint64, participants []FailoverParticipant) (FailoverResult, error) {
session, err := NewFailoverSession(master, volumeName, expectedEpoch, participants)
if err != nil {
return FailoverResult{}, err
}
return session.Run()
}
func (s *FailoverSession) fail(err error) error {
if s == nil {
return err
}
s.stage = FailoverStageFailed
s.lastErr = err
return err
}
func (s *FailoverSession) failf(format string, args ...any) error {
return s.fail(fmt.Errorf(format, args...))
}
@@ -0,0 +1,122 @@
package volumev2
import (
"fmt"
"slices"
"sync"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/masterv2"
)
// InProcessFailoverDriver is the first thin driver that wires one in-process
// masterv2 instance to a set of failover-capable participants. It owns no
// recovery logic; it only resolves participants and constructs sessions.
type InProcessFailoverDriver struct {
master *masterv2.Master
mu sync.RWMutex
participants map[string]FailoverParticipant
}
// NewInProcessFailoverDriver creates a driver for one in-process masterv2.
func NewInProcessFailoverDriver(master *masterv2.Master) (*InProcessFailoverDriver, error) {
if master == nil {
return nil, fmt.Errorf("volumev2: master is nil")
}
return &InProcessFailoverDriver{
master: master,
participants: make(map[string]FailoverParticipant),
}, nil
}
// RegisterParticipant binds one stable node id to one failover participant.
func (d *InProcessFailoverDriver) RegisterParticipant(nodeID string, participant FailoverParticipant) error {
if d == nil {
return fmt.Errorf("volumev2: failover driver is nil")
}
if nodeID == "" {
return fmt.Errorf("volumev2: participant node id is required")
}
if participant == nil {
return fmt.Errorf("volumev2: participant %q is nil", nodeID)
}
d.mu.Lock()
defer d.mu.Unlock()
d.participants[nodeID] = participant
return nil
}
// UnregisterParticipant removes one node from the in-process driver.
func (d *InProcessFailoverDriver) UnregisterParticipant(nodeID string) {
if d == nil || nodeID == "" {
return
}
d.mu.Lock()
defer d.mu.Unlock()
delete(d.participants, nodeID)
}
// ParticipantNodeIDs returns the currently registered node ids in stable order.
func (d *InProcessFailoverDriver) ParticipantNodeIDs() []string {
if d == nil {
return nil
}
d.mu.RLock()
defer d.mu.RUnlock()
nodeIDs := make([]string, 0, len(d.participants))
for nodeID := range d.participants {
nodeIDs = append(nodeIDs, nodeID)
}
slices.Sort(nodeIDs)
return nodeIDs
}
// NewSession constructs a FailoverSession using either the requested node ids
// or all currently registered participants when none are specified.
func (d *InProcessFailoverDriver) NewSession(volumeName string, expectedEpoch uint64, nodeIDs ...string) (*FailoverSession, error) {
if d == nil {
return nil, fmt.Errorf("volumev2: failover driver is nil")
}
participants, err := d.resolveParticipants(nodeIDs)
if err != nil {
return nil, err
}
return NewFailoverSession(d.master, volumeName, expectedEpoch, participants)
}
// Execute runs one failover using the resolved participant set.
func (d *InProcessFailoverDriver) Execute(volumeName string, expectedEpoch uint64, nodeIDs ...string) (FailoverResult, error) {
session, err := d.NewSession(volumeName, expectedEpoch, nodeIDs...)
if err != nil {
return FailoverResult{}, err
}
return session.Run()
}
func (d *InProcessFailoverDriver) resolveParticipants(nodeIDs []string) ([]FailoverParticipant, error) {
d.mu.RLock()
defer d.mu.RUnlock()
if len(d.participants) == 0 {
return nil, fmt.Errorf("volumev2: no failover participants registered")
}
resolvedIDs := nodeIDs
if len(resolvedIDs) == 0 {
resolvedIDs = make([]string, 0, len(d.participants))
for nodeID := range d.participants {
resolvedIDs = append(resolvedIDs, nodeID)
}
slices.Sort(resolvedIDs)
}
participants := make([]FailoverParticipant, 0, len(resolvedIDs))
for _, nodeID := range resolvedIDs {
participant, ok := d.participants[nodeID]
if !ok {
return nil, fmt.Errorf("volumev2: unknown failover participant %q", nodeID)
}
participants = append(participants, participant)
}
return participants, nil
}
+92
View File
@@ -0,0 +1,92 @@
package volumev2
import (
"fmt"
"io"
"log"
"net"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol/iscsi"
)
// ISCSITargetExport is a small handle for one running iSCSI frontend export.
type ISCSITargetExport struct {
iqn string
addr string
server *iscsi.TargetServer
}
// IQN returns the exported target name.
func (e *ISCSITargetExport) IQN() string {
if e == nil {
return ""
}
return e.iqn
}
// Address returns the current listen address.
func (e *ISCSITargetExport) Address() string {
if e == nil {
return ""
}
return e.addr
}
// Close stops the iSCSI target server.
func (e *ISCSITargetExport) Close() error {
if e == nil || e.server == nil {
return nil
}
return e.server.Close()
}
// ExportISCSI starts a small iSCSI target server for one named volume.
// This gives the single-node MVP a real block frontend without depending on weed/server.
func (n *Node) ExportISCSI(name, listenAddr, iqn string) (*ISCSITargetExport, error) {
if n == nil {
return nil, fmt.Errorf("volumev2: node is nil")
}
if listenAddr == "" {
listenAddr = "127.0.0.1:0"
}
if iqn == "" {
iqn = "iqn.2026-04.com.seaweedfs:v2." + name
}
path, err := n.pathFor(name)
if err != nil {
return nil, err
}
var dev iscsi.BlockDevice
if err := n.dataPlane.WithVolume(path, func(vol *blockvol.BlockVol) error {
dev = blockvol.NewBlockVolAdapter(vol)
return nil
}); err != nil {
return nil, err
}
if dev == nil {
return nil, fmt.Errorf("volumev2: no block device for %q", name)
}
cfg := iscsi.DefaultTargetConfig()
cfg.TargetName = iqn
logger := log.New(io.Discard, "", 0)
server := iscsi.NewTargetServer(listenAddr, cfg, logger)
server.AddVolume(iqn, dev)
ln, err := net.Listen("tcp", listenAddr)
if err != nil {
return nil, fmt.Errorf("volumev2: iscsi listen %s: %w", listenAddr, err)
}
server.SetPortalAddr(ln.Addr().String() + ",1")
go func() {
_ = server.Serve(ln)
}()
return &ISCSITargetExport{
iqn: iqn,
addr: ln.Addr().String(),
server: server,
}, nil
}
+233
View File
@@ -0,0 +1,233 @@
package volumev2
import (
"bytes"
"encoding/binary"
"net"
"path/filepath"
"testing"
"time"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/masterv2"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol/iscsi"
)
func TestVolumeV2_ExportISCSI_StartsFrontendAndAcceptsLogin(t *testing.T) {
master := masterv2.New(masterv2.Config{})
node, err := New(Config{NodeID: "node-a"})
if err != nil {
t.Fatalf("new node: %v", err)
}
defer node.Close()
session, err := NewInProcessSession(master)
if err != nil {
t.Fatalf("new session: %v", err)
}
orchestrator, err := NewOrchestrator(node, session)
if err != nil {
t.Fatalf("new orchestrator: %v", err)
}
path := filepath.Join(t.TempDir(), "frontend-vol.blk")
if err := master.DeclarePrimary(masterv2.VolumeSpec{
Name: "frontend-vol",
Path: path,
PrimaryNodeID: "node-a",
CreateOptions: testCreateOptions(),
}); err != nil {
t.Fatalf("declare primary: %v", err)
}
if err := orchestrator.SyncOnce(); err != nil {
t.Fatalf("sync 1: %v", err)
}
if err := orchestrator.SyncOnce(); err != nil {
t.Fatalf("sync 2: %v", err)
}
export, err := node.ExportISCSI("frontend-vol", "127.0.0.1:0", "iqn.2026-04.com.seaweedfs:test.frontend-vol")
if err != nil {
t.Fatalf("export iscsi: %v", err)
}
defer export.Close()
conn, err := net.DialTimeout("tcp", export.Address(), 2*time.Second)
if err != nil {
t.Fatalf("dial target: %v", err)
}
defer conn.Close()
params := iscsi.NewParams()
params.Set("InitiatorName", "iqn.2026-04.com.seaweedfs:initiator.test")
params.Set("TargetName", export.IQN())
params.Set("SessionType", "Normal")
loginReq := &iscsi.PDU{}
loginReq.SetOpcode(iscsi.OpLoginReq)
loginReq.SetLoginStages(iscsi.StageSecurityNeg, iscsi.StageFullFeature)
loginReq.SetLoginTransit(true)
loginReq.SetISID([6]byte{0x00, 0x02, 0x3D, 0x00, 0x00, 0x01})
loginReq.SetCmdSN(1)
loginReq.DataSegment = params.Encode()
if err := iscsi.WritePDU(conn, loginReq); err != nil {
t.Fatalf("write login req: %v", err)
}
resp, err := iscsi.ReadPDU(conn)
if err != nil {
t.Fatalf("read login resp: %v", err)
}
if resp.LoginStatusClass() != iscsi.LoginStatusSuccess {
t.Fatalf("login failed: %d/%d", resp.LoginStatusClass(), resp.LoginStatusDetail())
}
}
func TestKernelDataPlaneClosure_ISCSIWriteReadVerify(t *testing.T) {
master := masterv2.New(masterv2.Config{})
node, err := New(Config{NodeID: "node-a"})
if err != nil {
t.Fatalf("new node: %v", err)
}
defer node.Close()
session, err := NewInProcessSession(master)
if err != nil {
t.Fatalf("new session: %v", err)
}
orchestrator, err := NewOrchestrator(node, session)
if err != nil {
t.Fatalf("new orchestrator: %v", err)
}
path := filepath.Join(t.TempDir(), "data-plane-vol.blk")
if err := master.DeclarePrimary(masterv2.VolumeSpec{
Name: "data-plane-vol",
Path: path,
PrimaryNodeID: "node-a",
CreateOptions: testCreateOptions(),
}); err != nil {
t.Fatalf("declare primary: %v", err)
}
if err := orchestrator.SyncOnce(); err != nil {
t.Fatalf("sync 1: %v", err)
}
if err := orchestrator.SyncOnce(); err != nil {
t.Fatalf("sync 2: %v", err)
}
export, err := node.ExportISCSI("data-plane-vol", "127.0.0.1:0", "iqn.2026-04.com.seaweedfs:test.data-plane-vol")
if err != nil {
t.Fatalf("export iscsi: %v", err)
}
defer export.Close()
conn := mustLoginISCSI(t, export.Address(), export.IQN())
defer conn.Close()
writeData := make([]byte, 4096)
for i := range writeData {
writeData[i] = byte((i * 7) % 251)
}
var writeCDB [16]byte
writeCDB[0] = iscsi.ScsiWrite10
binary.BigEndian.PutUint32(writeCDB[2:6], 0)
binary.BigEndian.PutUint16(writeCDB[7:9], 1)
resp := sendSCSICmd(t, conn, writeCDB, 2, false, true, writeData, uint32(len(writeData)))
if resp.SCSIStatus() != iscsi.SCSIStatusGood {
t.Fatalf("iscsi write failed: status=%d", resp.SCSIStatus())
}
var syncCDB [16]byte
syncCDB[0] = iscsi.ScsiSyncCache10
resp = sendSCSICmd(t, conn, syncCDB, 3, false, false, nil, 0)
if resp.SCSIStatus() != iscsi.SCSIStatusGood {
t.Fatalf("iscsi sync cache failed: status=%d", resp.SCSIStatus())
}
var readCDB [16]byte
readCDB[0] = iscsi.ScsiRead10
binary.BigEndian.PutUint32(readCDB[2:6], 0)
binary.BigEndian.PutUint16(readCDB[7:9], 1)
resp = sendSCSICmd(t, conn, readCDB, 4, true, false, nil, uint32(len(writeData)))
if resp.Opcode() != iscsi.OpSCSIDataIn {
t.Fatalf("expected Data-In, got %s", iscsi.OpcodeName(resp.Opcode()))
}
if !bytes.Equal(resp.DataSegment, writeData) {
t.Fatal("iscsi readback mismatch")
}
readBack, err := node.ReadLBA("data-plane-vol", 0, uint32(len(writeData)))
if err != nil {
t.Fatalf("backend read: %v", err)
}
if !bytes.Equal(readBack, writeData) {
t.Fatal("backend readback mismatch")
}
}
func mustLoginISCSI(t *testing.T, addr, iqn string) net.Conn {
t.Helper()
conn, err := net.DialTimeout("tcp", addr, 2*time.Second)
if err != nil {
t.Fatalf("dial target: %v", err)
}
params := iscsi.NewParams()
params.Set("InitiatorName", "iqn.2026-04.com.seaweedfs:initiator.test")
params.Set("TargetName", iqn)
params.Set("SessionType", "Normal")
loginReq := &iscsi.PDU{}
loginReq.SetOpcode(iscsi.OpLoginReq)
loginReq.SetLoginStages(iscsi.StageSecurityNeg, iscsi.StageFullFeature)
loginReq.SetLoginTransit(true)
loginReq.SetISID([6]byte{0x00, 0x02, 0x3D, 0x00, 0x00, 0x01})
loginReq.SetCmdSN(1)
loginReq.DataSegment = params.Encode()
if err := iscsi.WritePDU(conn, loginReq); err != nil {
conn.Close()
t.Fatalf("write login req: %v", err)
}
resp, err := iscsi.ReadPDU(conn)
if err != nil {
conn.Close()
t.Fatalf("read login resp: %v", err)
}
if resp.LoginStatusClass() != iscsi.LoginStatusSuccess {
conn.Close()
t.Fatalf("login failed: %d/%d", resp.LoginStatusClass(), resp.LoginStatusDetail())
}
return conn
}
func sendSCSICmd(t *testing.T, conn net.Conn, cdb [16]byte, cmdSN uint32, read bool, write bool, dataOut []byte, expLen uint32) *iscsi.PDU {
t.Helper()
cmd := &iscsi.PDU{}
cmd.SetOpcode(iscsi.OpSCSICmd)
flags := uint8(iscsi.FlagF)
if read {
flags |= iscsi.FlagR
}
if write {
flags |= iscsi.FlagW
}
cmd.SetOpSpecific1(flags)
cmd.SetInitiatorTaskTag(cmdSN)
cmd.SetExpectedDataTransferLength(expLen)
cmd.SetCmdSN(cmdSN)
cmd.SetCDB(cdb)
if dataOut != nil {
cmd.DataSegment = dataOut
}
if err := iscsi.WritePDU(conn, cmd); err != nil {
t.Fatalf("write scsi cmd: %v", err)
}
resp, err := iscsi.ReadPDU(conn)
if err != nil {
t.Fatalf("read scsi resp: %v", err)
}
return resp
}
+37
View File
@@ -0,0 +1,37 @@
package volumev2
import "fmt"
// Orchestrator closes the MVP control loop:
// heartbeat -> assignments -> local apply.
type Orchestrator struct {
node *Node
session ControlSession
}
// NewOrchestrator creates a small volumev2 control/data orchestrator.
func NewOrchestrator(node *Node, session ControlSession) (*Orchestrator, error) {
if node == nil {
return nil, fmt.Errorf("volumev2: node is nil")
}
if session == nil {
return nil, fmt.Errorf("volumev2: control session is nil")
}
return &Orchestrator{node: node, session: session}, nil
}
// SyncOnce reports one heartbeat and applies any returned assignments.
func (o *Orchestrator) SyncOnce() error {
if o == nil || o.node == nil || o.session == nil {
return fmt.Errorf("volumev2: orchestrator is not initialized")
}
hb, err := o.node.Heartbeat()
if err != nil {
return err
}
assignments, err := o.session.Heartbeat(hb)
if err != nil {
return err
}
return o.node.ApplyAssignments(assignments)
}
File diff suppressed because it is too large. Load diff
+150
View File
@@ -0,0 +1,150 @@
package volumev2
import (
"fmt"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/protocolv2"
)
// ReconstructedPrimaryTruth is the bounded truth a newly selected primary can
// derive from replica summaries before resuming data-control ownership.
type ReconstructedPrimaryTruth struct {
VolumeName string
PrimaryNodeID string
Epoch uint64
CommittedLSN uint64
DurableLSN uint64
CheckpointLSN uint64
TargetLSN uint64
AchievedLSN uint64
RecoveryPhase string
ReplicaCount int
Degraded bool
NeedsRebuild bool
Reason string
}
// ReconstructPrimaryTruth derives a bounded recovery view for a newly chosen
// primary from the latest replica summaries. It intentionally stays smaller
// than the full internal engine/session graph and fail-closes on ambiguous
// epoch or recovery signals.
func ReconstructPrimaryTruth(primaryNodeID string, summaries []protocolv2.ReplicaSummaryResponse) (ReconstructedPrimaryTruth, error) {
if primaryNodeID == "" {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: primary node id is required")
}
if len(summaries) == 0 {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: replica summaries are required")
}
var (
selected protocolv2.ReplicaSummaryResponse
foundSelected bool
recoveryObserved bool
aggregateTarget uint64
aggregateAchieved uint64
)
for _, summary := range summaries {
if summary.NodeID == primaryNodeID {
selected = summary
foundSelected = true
break
}
}
if !foundSelected {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: selected primary %q missing from summaries", primaryNodeID)
}
if !selected.Eligible {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: selected primary %q is not eligible: %s", primaryNodeID, selected.Reason)
}
result := ReconstructedPrimaryTruth{
VolumeName: selected.VolumeName,
PrimaryNodeID: selected.NodeID,
Epoch: selected.Epoch,
CommittedLSN: selected.CommittedLSN,
DurableLSN: selected.DurableLSN,
CheckpointLSN: selected.CheckpointLSN,
TargetLSN: selected.TargetLSN,
AchievedLSN: selected.AchievedLSN,
RecoveryPhase: selected.RecoveryPhase,
}
for _, summary := range summaries {
if summary.VolumeName != selected.VolumeName {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: mixed volume summaries %q and %q", selected.VolumeName, summary.VolumeName)
}
if summary.Mode == "needs_rebuild" || summary.Reason == "needs_rebuild" {
result.NeedsRebuild = true
}
if !summary.LastBarrierOK && summary.LastBarrierReason != "" {
result.Degraded = true
if result.Reason == "" {
result.Reason = summary.LastBarrierReason
}
}
if summary.Epoch != selected.Epoch {
result.Degraded = true
result.Reason = "peer_epoch_mismatch"
continue
}
result.ReplicaCount++
if summary.CommittedLSN > result.CommittedLSN {
result.Degraded = true
result.Reason = "selected_not_most_recent"
}
if isRecoveryPhase(summary.RecoveryPhase) {
if !recoveryObserved {
aggregateTarget = summary.TargetLSN
aggregateAchieved = summary.AchievedLSN
recoveryObserved = true
} else {
if summary.TargetLSN > aggregateTarget {
aggregateTarget = summary.TargetLSN
}
if summary.AchievedLSN < aggregateAchieved {
aggregateAchieved = summary.AchievedLSN
}
}
result.RecoveryPhase = mergeRecoveryPhase(result.RecoveryPhase, summary.RecoveryPhase)
}
}
if recoveryObserved {
result.TargetLSN = aggregateTarget
result.AchievedLSN = aggregateAchieved
}
if result.NeedsRebuild {
result.RecoveryPhase = "needs_rebuild"
}
return result, nil
}
func isRecoveryPhase(phase string) bool {
switch phase {
case "catching_up", "rebuilding", "needs_rebuild":
return true
default:
return false
}
}
func mergeRecoveryPhase(current, next string) string {
if recoveryRank(next) > recoveryRank(current) {
return next
}
return current
}
func recoveryRank(phase string) int {
switch phase {
case "needs_rebuild":
return 3
case "rebuilding":
return 2
case "catching_up":
return 1
default:
return 0
}
}
@@ -0,0 +1,138 @@
package volumev2
import (
"testing"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/masterv2"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/protocolv2"
)
func TestPrimaryLoss_ReconstructsBoundedTruthFromReplicaSummaries(t *testing.T) {
master := masterv2.New(masterv2.Config{})
candidate, err := master.SelectPromotionCandidate([]masterv2.PromotionQueryResponse{
{
VolumeName: "vol-a",
NodeID: "node-b",
CommittedLSN: 15,
WALHeadLSN: 18,
Eligible: true,
},
{
VolumeName: "vol-a",
NodeID: "node-c",
CommittedLSN: 12,
WALHeadLSN: 13,
Eligible: true,
},
})
if err != nil {
t.Fatalf("select candidate: %v", err)
}
truth, err := ReconstructPrimaryTruth(candidate.NodeID, []protocolv2.ReplicaSummaryResponse{
{
VolumeName: "vol-a",
NodeID: "node-b",
Epoch: 4,
Role: "replica",
Mode: "replica_ready",
CommittedLSN: 15,
DurableLSN: 15,
CheckpointLSN: 10,
TargetLSN: 40,
AchievedLSN: 30,
RecoveryPhase: "catching_up",
LastBarrierOK: true,
Eligible: true,
},
{
VolumeName: "vol-a",
NodeID: "node-c",
Epoch: 4,
Role: "replica",
Mode: "replica_ready",
CommittedLSN: 12,
DurableLSN: 12,
CheckpointLSN: 8,
TargetLSN: 40,
AchievedLSN: 20,
RecoveryPhase: "catching_up",
LastBarrierOK: true,
Eligible: true,
},
})
if err != nil {
t.Fatalf("reconstruct truth: %v", err)
}
if truth.PrimaryNodeID != "node-b" {
t.Fatalf("primary_node=%q, want node-b", truth.PrimaryNodeID)
}
if truth.CommittedLSN != 15 {
t.Fatalf("committed_lsn=%d, want 15", truth.CommittedLSN)
}
if truth.DurableLSN != 15 {
t.Fatalf("durable_lsn=%d, want 15", truth.DurableLSN)
}
if truth.CheckpointLSN != 10 {
t.Fatalf("checkpoint_lsn=%d, want 10", truth.CheckpointLSN)
}
if truth.TargetLSN != 40 {
t.Fatalf("target_lsn=%d, want 40", truth.TargetLSN)
}
if truth.AchievedLSN != 20 {
t.Fatalf("achieved_lsn=%d, want 20", truth.AchievedLSN)
}
if truth.RecoveryPhase != "catching_up" {
t.Fatalf("recovery_phase=%q, want catching_up", truth.RecoveryPhase)
}
if truth.Degraded {
t.Fatalf("unexpected degraded truth: %+v", truth)
}
}
func TestPrimaryLoss_ReconstructionFailsClosedOnMismatchAndNeedsRebuild(t *testing.T) {
truth, err := ReconstructPrimaryTruth("node-b", []protocolv2.ReplicaSummaryResponse{
{
VolumeName: "vol-a",
NodeID: "node-b",
Epoch: 5,
Role: "replica",
Mode: "replica_ready",
CommittedLSN: 15,
DurableLSN: 15,
CheckpointLSN: 10,
RecoveryPhase: "idle",
LastBarrierOK: true,
Eligible: true,
LastBarrierReason: "",
},
{
VolumeName: "vol-a",
NodeID: "node-c",
Epoch: 4,
Role: "replica",
Mode: "needs_rebuild",
CommittedLSN: 12,
DurableLSN: 10,
CheckpointLSN: 8,
RecoveryPhase: "needs_rebuild",
LastBarrierOK: false,
LastBarrierReason: "timeout",
Eligible: false,
Reason: "needs_rebuild",
},
})
if err != nil {
t.Fatalf("reconstruct truth: %v", err)
}
if !truth.Degraded {
t.Fatalf("expected degraded truth: %+v", truth)
}
if !truth.NeedsRebuild {
t.Fatalf("expected needs_rebuild truth: %+v", truth)
}
if truth.RecoveryPhase != "needs_rebuild" {
t.Fatalf("recovery_phase=%q, want needs_rebuild", truth.RecoveryPhase)
}
}
+115
View File
@@ -0,0 +1,115 @@
package volumev2
import (
"fmt"
"slices"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/masterv2"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/protocolv2"
)
// ReplicaSummarySource is the bounded Loop 2 query surface a replacement
// primary uses to reconstruct takeover truth from peers.
type ReplicaSummarySource interface {
QueryReplicaSummary(protocolv2.ReplicaSummaryRequest) (protocolv2.ReplicaSummaryResponse, error)
}
// PrimaryTakeoverPlan is the minimal input needed for a selected replacement
// primary to prepare takeover locally.
type PrimaryTakeoverPlan struct {
Assignment masterv2.Assignment
Peers []ReplicaSummarySource
}
// ReconstructTakeoverTruth lets the selected replacement primary gather its own
// bounded summary plus peer summaries before it resumes data-control ownership.
// Peers should exclude the current node; duplicate node IDs are ignored.
func (n *Node) ReconstructTakeoverTruth(volumeName string, expectedEpoch uint64, peers []ReplicaSummarySource) (ReconstructedPrimaryTruth, error) {
if n == nil {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: node is nil")
}
req := protocolv2.ReplicaSummaryRequest{
VolumeName: volumeName,
ExpectedEpoch: expectedEpoch,
}
self, err := n.QueryReplicaSummary(req)
if err != nil {
return ReconstructedPrimaryTruth{}, err
}
byNode := map[string]protocolv2.ReplicaSummaryResponse{
self.NodeID: self,
}
for _, peer := range peers {
if peer == nil {
continue
}
summary, err := peer.QueryReplicaSummary(req)
if err != nil {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: peer replica summary %s: %w", volumeName, err)
}
if summary.NodeID == "" || summary.NodeID == n.id {
continue
}
byNode[summary.NodeID] = summary
}
nodeIDs := make([]string, 0, len(byNode))
for nodeID := range byNode {
nodeIDs = append(nodeIDs, nodeID)
}
slices.Sort(nodeIDs)
summaries := make([]protocolv2.ReplicaSummaryResponse, 0, len(nodeIDs))
for _, nodeID := range nodeIDs {
summaries = append(summaries, byNode[nodeID])
}
return ReconstructPrimaryTruth(n.id, summaries)
}
// PreparePrimaryTakeover applies the local primary assignment and reconstructs
// bounded takeover truth from self and peers. It does not decide whether the
// new primary is allowed to activate data control yet.
func (n *Node) PreparePrimaryTakeover(plan PrimaryTakeoverPlan) (ReconstructedPrimaryTruth, error) {
if n == nil {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: node is nil")
}
a := plan.Assignment
if a.Role != "primary" {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: unsupported takeover role %q", a.Role)
}
if a.NodeID != "" && a.NodeID != n.id {
return ReconstructedPrimaryTruth{}, fmt.Errorf("volumev2: takeover assignment targets %q, node is %q", a.NodeID, n.id)
}
if err := n.ApplyAssignments([]masterv2.Assignment{a}); err != nil {
return ReconstructedPrimaryTruth{}, err
}
return n.ReconstructTakeoverTruth(a.Name, a.Epoch, plan.Peers)
}
// GatePrimaryActivation fail-closes activation when the reconstructed truth
// says takeover is degraded, ambiguous, or rebuild-only.
func (n *Node) GatePrimaryActivation(volumeName string, truth ReconstructedPrimaryTruth) error {
if n == nil {
return fmt.Errorf("volumev2: node is nil")
}
if truth.NeedsRebuild {
return fmt.Errorf("volumev2: takeover gated for %s: needs rebuild", volumeName)
}
if truth.Degraded {
return fmt.Errorf("volumev2: takeover gated for %s: %s", volumeName, truth.Reason)
}
return nil
}
// ApplyPrimaryTakeover is a narrow compatibility wrapper that prepares
// takeover truth and then gates activation.
func (n *Node) ApplyPrimaryTakeover(plan PrimaryTakeoverPlan) (ReconstructedPrimaryTruth, error) {
truth, err := n.PreparePrimaryTakeover(plan)
if err != nil {
return ReconstructedPrimaryTruth{}, err
}
if err := n.GatePrimaryActivation(plan.Assignment.Name, truth); err != nil {
return truth, err
}
return truth, nil
}
+312
View File
@@ -0,0 +1,312 @@
package volumev2
import (
"fmt"
"sync"
"time"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/masterv2"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/protocolv2"
"github.com/seaweedfs/seaweedfs/sw-block/runtime/purev2"
"github.com/seaweedfs/seaweedfs/weed/storage/blockvol"
)
// Config defines the minimal volumev2 runtime identity.
type Config struct {
NodeID string
DataPlane DataPlane
Runtime *purev2.Runtime
}
type volumeBinding struct {
name string
path string
}
// Node is the minimal volumev2 runtime POC.
// It owns volume identity, local execution, and heartbeat reporting.
type Node struct {
id string
dataPlane DataPlane
mu sync.RWMutex
volumes map[string]volumeBinding
}
// New creates a minimal volumev2 node runtime.
func New(cfg Config) (*Node, error) {
if cfg.NodeID == "" {
return nil, fmt.Errorf("volumev2: node id is required")
}
dp := cfg.DataPlane
if dp == nil {
rt := cfg.Runtime
if rt == nil {
rt = purev2.New(purev2.Config{})
}
var err error
dp, err = NewPureRuntimeDataPlane(rt)
if err != nil {
return nil, err
}
}
return &Node{
id: cfg.NodeID,
dataPlane: dp,
volumes: make(map[string]volumeBinding),
}, nil
}
// NodeID returns the stable runtime identity.
func (n *Node) NodeID() string {
if n == nil {
return ""
}
return n.id
}
// ApplyAssignments applies the control messages emitted by masterv2.
func (n *Node) ApplyAssignments(assignments []masterv2.Assignment) error {
if n == nil {
return fmt.Errorf("volumev2: node is nil")
}
for _, assignment := range assignments {
if assignment.NodeID != "" && assignment.NodeID != n.id {
continue
}
if assignment.Role != "primary" {
return fmt.Errorf("volumev2: unsupported role %q", assignment.Role)
}
if err := n.dataPlane.BootstrapPrimary(
assignment.Path,
assignment.CreateOptions,
assignment.Epoch,
assignment.LeaseTTL,
); err != nil {
return fmt.Errorf("volumev2: apply assignment %s: %w", assignment.Name, err)
}
n.mu.Lock()
n.volumes[assignment.Name] = volumeBinding{name: assignment.Name, path: assignment.Path}
n.mu.Unlock()
}
return nil
}
// Heartbeat reports the minimal local runtime state back to masterv2.
func (n *Node) Heartbeat() (masterv2.NodeHeartbeat, error) {
if n == nil {
return masterv2.NodeHeartbeat{}, fmt.Errorf("volumev2: node is nil")
}
n.mu.RLock()
bindings := make([]volumeBinding, 0, len(n.volumes))
for _, binding := range n.volumes {
bindings = append(bindings, binding)
}
n.mu.RUnlock()
report := masterv2.NodeHeartbeat{
NodeID: n.id,
ReportedAt: time.Now(),
Volumes: make([]masterv2.VolumeHeartbeat, 0, len(bindings)),
}
for _, binding := range bindings {
snap, err := n.dataPlane.Snapshot(binding.path)
if err != nil {
return masterv2.NodeHeartbeat{}, fmt.Errorf("volumev2: snapshot %s: %w", binding.name, err)
}
var committedLSN uint64
if err := n.dataPlane.WithVolume(binding.path, func(vol *blockvol.BlockVol) error {
committedLSN = vol.StatusSnapshot().CommittedLSN
return nil
}); err != nil {
return masterv2.NodeHeartbeat{}, fmt.Errorf("volumev2: committed snapshot %s: %w", binding.name, err)
}
vol := masterv2.VolumeHeartbeat{
Name: binding.name,
Path: binding.path,
Epoch: snap.Status.Epoch,
Role: snap.Status.Role.String(),
CommittedLSN: committedLSN,
}
if snap.HasProjection {
vol.Mode = string(snap.Projection.Mode.Name)
vol.ModeReason = snap.Projection.Mode.Reason
vol.RoleApplied = snap.Projection.Readiness.RoleApplied
vol.ReplicaReady = snap.Projection.Readiness.ReplicaReady
}
report.Volumes = append(report.Volumes, vol)
}
return report, nil
}
// QueryPromotionEvidence returns fresh failover evidence outside the periodic
// heartbeat path. This keeps promotion arbitration separate from liveness.
func (n *Node) QueryPromotionEvidence(req masterv2.PromotionQueryRequest) (masterv2.PromotionQueryResponse, error) {
path, err := n.pathFor(req.VolumeName)
if err != nil {
return masterv2.PromotionQueryResponse{}, err
}
snap, err := n.dataPlane.Snapshot(path)
if err != nil {
return masterv2.PromotionQueryResponse{}, fmt.Errorf("volumev2: snapshot %s: %w", req.VolumeName, err)
}
resp := masterv2.PromotionQueryResponse{
VolumeName: req.VolumeName,
NodeID: n.id,
Epoch: snap.Status.Epoch,
Role: snap.Status.Role.String(),
}
if err := n.dataPlane.WithVolume(path, func(vol *blockvol.BlockVol) error {
status := vol.StatusSnapshot()
resp.CommittedLSN = status.CommittedLSN
resp.WALHeadLSN = status.WALHeadLSN
return nil
}); err != nil {
return masterv2.PromotionQueryResponse{}, fmt.Errorf("volumev2: promotion evidence %s: %w", req.VolumeName, err)
}
if snap.HasProjection {
resp.ReceiverReady = snap.Projection.Readiness.ReplicaReady
switch {
case req.ExpectedEpoch != 0 && resp.Epoch != req.ExpectedEpoch:
resp.Reason = "epoch_mismatch"
case !snap.Projection.Readiness.RoleApplied:
resp.Reason = "role_not_applied"
case snap.Projection.Mode.Name == "needs_rebuild":
resp.Reason = "needs_rebuild"
default:
resp.Eligible = true
}
} else if req.ExpectedEpoch != 0 && resp.Epoch != req.ExpectedEpoch {
resp.Reason = "epoch_mismatch"
} else {
resp.Eligible = true
}
return resp, nil
}
// QueryReplicaSummary returns a bounded takeover/reconstruction summary for one
// volume. It preserves distinct LSN semantics while avoiding raw internal
// shipper/session detail.
func (n *Node) QueryReplicaSummary(req protocolv2.ReplicaSummaryRequest) (protocolv2.ReplicaSummaryResponse, error) {
path, err := n.pathFor(req.VolumeName)
if err != nil {
return protocolv2.ReplicaSummaryResponse{}, err
}
snap, err := n.dataPlane.Snapshot(path)
if err != nil {
return protocolv2.ReplicaSummaryResponse{}, fmt.Errorf("volumev2: snapshot %s: %w", req.VolumeName, err)
}
resp := protocolv2.ReplicaSummaryResponse{
VolumeName: req.VolumeName,
NodeID: n.id,
Epoch: snap.Status.Epoch,
Role: snap.Status.Role.String(),
}
var status blockvol.V2StatusSnapshot
if err := n.dataPlane.WithVolume(path, func(vol *blockvol.BlockVol) error {
status = vol.StatusSnapshot()
return nil
}); err != nil {
return protocolv2.ReplicaSummaryResponse{}, fmt.Errorf("volumev2: replica summary %s: %w", req.VolumeName, err)
}
resp.CommittedLSN = status.CommittedLSN
resp.CheckpointLSN = status.CheckpointLSN
if snap.HasProjection {
resp.Mode = string(snap.Projection.Mode.Name)
resp.ModeReason = snap.Projection.Mode.Reason
resp.RoleApplied = snap.Projection.Readiness.RoleApplied
resp.ReceiverReady = snap.Projection.Readiness.ReplicaReady
resp.DurableLSN = snap.Projection.Boundary.DurableLSN
resp.TargetLSN = snap.Projection.Boundary.TargetLSN
resp.AchievedLSN = snap.Projection.Boundary.AchievedLSN
resp.LastBarrierOK = snap.Projection.Boundary.LastBarrierOK
resp.LastBarrierReason = snap.Projection.Boundary.LastBarrierReason
}
if snap.HasCoreState {
resp.RecoveryPhase = string(snap.CoreState.Recovery.Phase)
if snap.CoreState.Boundary.DurableLSN > resp.DurableLSN {
resp.DurableLSN = snap.CoreState.Boundary.DurableLSN
}
if snap.CoreState.Boundary.TargetLSN > resp.TargetLSN {
resp.TargetLSN = snap.CoreState.Boundary.TargetLSN
}
if snap.CoreState.Boundary.AchievedLSN > resp.AchievedLSN {
resp.AchievedLSN = snap.CoreState.Boundary.AchievedLSN
}
if snap.CoreState.Boundary.CheckpointLSN > resp.CheckpointLSN {
resp.CheckpointLSN = snap.CoreState.Boundary.CheckpointLSN
}
if snap.CoreState.Boundary.LastBarrierReason != "" {
resp.LastBarrierReason = snap.CoreState.Boundary.LastBarrierReason
}
resp.LastBarrierOK = snap.CoreState.Boundary.LastBarrierOK
}
switch {
case req.ExpectedEpoch != 0 && resp.Epoch != req.ExpectedEpoch:
resp.Reason = "epoch_mismatch"
case !resp.RoleApplied:
resp.Reason = "role_not_applied"
case resp.Mode == "needs_rebuild":
resp.Reason = "needs_rebuild"
default:
resp.Eligible = true
}
return resp, nil
}
// WriteLBA writes data to one named local volume.
func (n *Node) WriteLBA(name string, lba uint64, data []byte) error {
path, err := n.pathFor(name)
if err != nil {
return err
}
return n.dataPlane.WriteLBA(path, lba, data)
}
// ReadLBA reads data from one named local volume.
func (n *Node) ReadLBA(name string, lba uint64, length uint32) ([]byte, error) {
path, err := n.pathFor(name)
if err != nil {
return nil, err
}
return n.dataPlane.ReadLBA(path, lba, length)
}
// SyncCache flushes one named local volume.
func (n *Node) SyncCache(name string) error {
path, err := n.pathFor(name)
if err != nil {
return err
}
return n.dataPlane.SyncCache(path)
}
// Snapshot returns the local debug snapshot for one named volume.
func (n *Node) Snapshot(name string) (purev2.VolumeDebugSnapshot, error) {
path, err := n.pathFor(name)
if err != nil {
return purev2.VolumeDebugSnapshot{}, err
}
return n.dataPlane.Snapshot(path)
}
// Close shuts down the underlying pure runtime.
func (n *Node) Close() {
if n == nil || n.dataPlane == nil {
return
}
n.dataPlane.Close()
}
func (n *Node) pathFor(name string) (string, error) {
n.mu.RLock()
defer n.mu.RUnlock()
binding, ok := n.volumes[name]
if !ok {
return "", fmt.Errorf("volumev2: unknown volume %q", name)
}
return binding.path, nil
}