diff --git a/weed/pb/master.proto b/weed/pb/master.proto
index 1332e5119..7464b8417 100644
--- a/weed/pb/master.proto
+++ b/weed/pb/master.proto
@@ -505,6 +505,7 @@ message BlockVolumeInfoMessage {
int64 scrub_errors = 13;
int64 last_scrub_time = 14;
bool replica_degraded = 15;
+ string durability_mode = 16;
}
message BlockVolumeShortInfoMessage {
@@ -535,6 +536,7 @@ message CreateBlockVolumeRequest {
uint64 size_bytes = 2;
string disk_type = 3;
uint32 replica_factor = 4;
+ string durability_mode = 5;
}
message CreateBlockVolumeResponse {
string volume_id = 1;
@@ -563,6 +565,7 @@ message LookupBlockVolumeResponse {
string replica_server = 5;
uint32 replica_factor = 6;
repeated string replica_servers = 7;
+ string durability_mode = 8;
}
message CreateBlockSnapshotRequest {
diff --git a/weed/pb/master_pb/master.pb.go b/weed/pb/master_pb/master.pb.go
index 0d16c4646..a54845644 100644
--- a/weed/pb/master_pb/master.pb.go
+++ b/weed/pb/master_pb/master.pb.go
@@ -3899,6 +3899,7 @@ type BlockVolumeInfoMessage struct {
ScrubErrors int64 `protobuf:"varint,13,opt,name=scrub_errors,json=scrubErrors,proto3" json:"scrub_errors,omitempty"`
LastScrubTime int64 `protobuf:"varint,14,opt,name=last_scrub_time,json=lastScrubTime,proto3" json:"last_scrub_time,omitempty"`
ReplicaDegraded bool `protobuf:"varint,15,opt,name=replica_degraded,json=replicaDegraded,proto3" json:"replica_degraded,omitempty"`
+ DurabilityMode string `protobuf:"bytes,16,opt,name=durability_mode,json=durabilityMode,proto3" json:"durability_mode,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -4038,6 +4039,13 @@ func (x *BlockVolumeInfoMessage) GetReplicaDegraded() bool {
return false
}
+func (x *BlockVolumeInfoMessage) GetDurabilityMode() string {
+ if x != nil {
+ return x.DurabilityMode
+ }
+ return ""
+}
+
type BlockVolumeShortInfoMessage struct {
state protoimpl.MessageState `protogen:"open.v1"`
Path string `protobuf:"bytes,1,opt,name=path,proto3" json:"path,omitempty"`
@@ -4238,9 +4246,10 @@ type CreateBlockVolumeRequest struct {
Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"`
SizeBytes uint64 `protobuf:"varint,2,opt,name=size_bytes,json=sizeBytes,proto3" json:"size_bytes,omitempty"`
DiskType string `protobuf:"bytes,3,opt,name=disk_type,json=diskType,proto3" json:"disk_type,omitempty"`
- ReplicaFactor uint32 `protobuf:"varint,4,opt,name=replica_factor,json=replicaFactor,proto3" json:"replica_factor,omitempty"`
- unknownFields protoimpl.UnknownFields
- sizeCache protoimpl.SizeCache
+ ReplicaFactor uint32 `protobuf:"varint,4,opt,name=replica_factor,json=replicaFactor,proto3" json:"replica_factor,omitempty"`
+ DurabilityMode string `protobuf:"bytes,5,opt,name=durability_mode,json=durabilityMode,proto3" json:"durability_mode,omitempty"`
+ unknownFields protoimpl.UnknownFields
+ sizeCache protoimpl.SizeCache
}
func (x *CreateBlockVolumeRequest) Reset() {
@@ -4301,6 +4310,13 @@ func (x *CreateBlockVolumeRequest) GetReplicaFactor() uint32 {
return 0
}
+func (x *CreateBlockVolumeRequest) GetDurabilityMode() string {
+ if x != nil {
+ return x.DurabilityMode
+ }
+ return ""
+}
+
type CreateBlockVolumeResponse struct {
state protoimpl.MessageState `protogen:"open.v1"`
VolumeId string `protobuf:"bytes,1,opt,name=volume_id,json=volumeId,proto3" json:"volume_id,omitempty"`
@@ -4526,6 +4542,7 @@ type LookupBlockVolumeResponse struct {
ReplicaServer string `protobuf:"bytes,5,opt,name=replica_server,json=replicaServer,proto3" json:"replica_server,omitempty"`
ReplicaFactor uint32 `protobuf:"varint,6,opt,name=replica_factor,json=replicaFactor,proto3" json:"replica_factor,omitempty"`
ReplicaServers []string `protobuf:"bytes,7,rep,name=replica_servers,json=replicaServers,proto3" json:"replica_servers,omitempty"`
+ DurabilityMode string `protobuf:"bytes,8,opt,name=durability_mode,json=durabilityMode,proto3" json:"durability_mode,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -4609,6 +4626,13 @@ func (x *LookupBlockVolumeResponse) GetReplicaServers() []string {
return nil
}
+func (x *LookupBlockVolumeResponse) GetDurabilityMode() string {
+ if x != nil {
+ return x.DurabilityMode
+ }
+ return ""
+}
+
type CreateBlockSnapshotRequest struct {
state protoimpl.MessageState `protogen:"open.v1"`
VolumeName string `protobuf:"bytes,1,opt,name=volume_name,json=volumeName,proto3" json:"volume_name,omitempty"`
diff --git a/weed/pb/volume_server.proto b/weed/pb/volume_server.proto
index 28e1b6740..5d3b348eb 100644
--- a/weed/pb/volume_server.proto
+++ b/weed/pb/volume_server.proto
@@ -779,6 +779,7 @@ message AllocateBlockVolumeRequest {
string name = 1;
uint64 size_bytes = 2;
string disk_type = 3;
+ string durability_mode = 4;
}
message AllocateBlockVolumeResponse {
string path = 1;
diff --git a/weed/pb/volume_server_pb/volume_server.pb.go b/weed/pb/volume_server_pb/volume_server.pb.go
index 88f77f2e3..8efe1cf70 100644
--- a/weed/pb/volume_server_pb/volume_server.pb.go
+++ b/weed/pb/volume_server_pb/volume_server.pb.go
@@ -6189,9 +6189,10 @@ type AllocateBlockVolumeRequest struct {
state protoimpl.MessageState `protogen:"open.v1"`
Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"`
SizeBytes uint64 `protobuf:"varint,2,opt,name=size_bytes,json=sizeBytes,proto3" json:"size_bytes,omitempty"`
- DiskType string `protobuf:"bytes,3,opt,name=disk_type,json=diskType,proto3" json:"disk_type,omitempty"`
- unknownFields protoimpl.UnknownFields
- sizeCache protoimpl.SizeCache
+ DiskType string `protobuf:"bytes,3,opt,name=disk_type,json=diskType,proto3" json:"disk_type,omitempty"`
+ DurabilityMode string `protobuf:"bytes,4,opt,name=durability_mode,json=durabilityMode,proto3" json:"durability_mode,omitempty"`
+ unknownFields protoimpl.UnknownFields
+ sizeCache protoimpl.SizeCache
}
func (x *AllocateBlockVolumeRequest) Reset() {
@@ -6245,6 +6246,13 @@ func (x *AllocateBlockVolumeRequest) GetDiskType() string {
return ""
}
+func (x *AllocateBlockVolumeRequest) GetDurabilityMode() string {
+ if x != nil {
+ return x.DurabilityMode
+ }
+ return ""
+}
+
type AllocateBlockVolumeResponse struct {
state protoimpl.MessageState `protogen:"open.v1"`
Path string `protobuf:"bytes,1,opt,name=path,proto3" json:"path,omitempty"`
diff --git a/weed/server/integration_block_test.go b/weed/server/integration_block_test.go
index d3a590e85..f5d9aad87 100644
--- a/weed/server/integration_block_test.go
+++ b/weed/server/integration_block_test.go
@@ -28,7 +28,7 @@ func integrationMaster(t *testing.T) *MasterServer {
blockAssignmentQueue: NewBlockAssignmentQueue(),
blockFailover: newBlockFailoverState(),
}
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
IQN: fmt.Sprintf("iqn.2024.test:%s", name),
@@ -382,13 +382,13 @@ func TestIntegration_ReplicaFailureSingleCopy(t *testing.T) {
// Make replica allocation always fail.
callCount := 0
origAllocate := ms.blockVSAllocate
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
callCount++
if callCount > 1 {
// Second call (replica) fails.
return nil, fmt.Errorf("disk full on replica")
}
- return origAllocate(ctx, server, name, sizeBytes, diskType)
+ return origAllocate(ctx, server, name, sizeBytes, diskType, durabilityMode)
}
resp, err := ms.CreateBlockVolume(ctx, &master_pb.CreateBlockVolumeRequest{
diff --git a/weed/server/master_block_failover_test.go b/weed/server/master_block_failover_test.go
index 1b60a2f92..6d6439068 100644
--- a/weed/server/master_block_failover_test.go
+++ b/weed/server/master_block_failover_test.go
@@ -18,7 +18,7 @@ func testMasterServerForFailover(t *testing.T) *MasterServer {
blockAssignmentQueue: NewBlockAssignmentQueue(),
blockFailover: newBlockFailoverState(),
}
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
IQN: fmt.Sprintf("iqn.2024.test:%s", name),
diff --git a/weed/server/master_block_registry.go b/weed/server/master_block_registry.go
index 0187ecf24..1cc4cabd7 100644
--- a/weed/server/master_block_registry.go
+++ b/weed/server/master_block_registry.go
@@ -72,6 +72,9 @@ type BlockVolumeEntry struct {
ReplicaDegraded bool // primary reports degraded replicas
WALHeadLSN uint64 // primary WAL head LSN from heartbeat
+ // CP8-3-1: Durability mode.
+ DurabilityMode string // "best_effort", "sync_all", "sync_quorum"
+
// Lease tracking for failover (CP6-3 F2).
LastLeaseGrant time.Time
LeaseTTL time.Duration
@@ -320,6 +323,10 @@ func (r *BlockVolumeRegistry) UpdateFullHeartbeat(server string, infos []*master
existing.HealthScore = info.HealthScore
existing.ReplicaDegraded = info.ReplicaDegraded
existing.WALHeadLSN = info.WalHeadLsn
+ // F3: only update DurabilityMode when non-empty (prevents older VS from clearing strict mode).
+ if info.DurabilityMode != "" {
+ existing.DurabilityMode = info.DurabilityMode
+ }
// F5: update replica addresses from heartbeat info.
if info.ReplicaDataAddr != "" {
existing.ReplicaDataAddr = info.ReplicaDataAddr
@@ -391,6 +398,7 @@ func (r *BlockVolumeRegistry) UpdateFullHeartbeat(server string, infos []*master
LeaseTTL: 30 * time.Second,
HealthScore: info.HealthScore,
WALHeadLSN: info.WalHeadLsn,
+ DurabilityMode: info.DurabilityMode,
}
if info.ReplicaDataAddr != "" {
entry.ReplicaDataAddr = info.ReplicaDataAddr
diff --git a/weed/server/master_grpc_server_block.go b/weed/server/master_grpc_server_block.go
index bbf04df0b..0a5913d9d 100644
--- a/weed/server/master_grpc_server_block.go
+++ b/weed/server/master_grpc_server_block.go
@@ -33,10 +33,29 @@ func (ms *MasterServer) CreateBlockVolume(ctx context.Context, req *master_pb.Cr
return nil, fmt.Errorf("size_bytes must be > 0")
}
- // Idempotent: if already registered, return existing entry (validate size).
+ // F2: validate durability mode in gRPC path (authoritative, not bypassable).
+ var durMode blockvol.DurabilityMode
+ if req.DurabilityMode != "" {
+ var err error
+ durMode, err = blockvol.ParseDurabilityMode(req.DurabilityMode)
+ if err != nil {
+ return nil, fmt.Errorf("invalid durability_mode: %w", err)
+ }
+ }
+
+ // Cross-validate mode + RF (sync_quorum requires RF >= 3).
+ replicaFactor := 2
+ if req.ReplicaFactor > 0 && req.ReplicaFactor <= 3 {
+ replicaFactor = int(req.ReplicaFactor)
+ }
+ if err := durMode.Validate(replicaFactor); err != nil {
+ return nil, fmt.Errorf("durability_mode %q incompatible with replica_factor %d: %w", req.DurabilityMode, replicaFactor, err)
+ }
+
+ // Idempotent: if already registered, return existing entry (validate size + mode + RF).
if entry, ok := ms.blockRegistry.Lookup(req.Name); ok {
- if entry.SizeBytes < req.SizeBytes {
- return nil, fmt.Errorf("block volume %q exists with size %d (requested %d)", req.Name, entry.SizeBytes, req.SizeBytes)
+ if err := ms.validateIdempotentCreate(entry, req, durMode, replicaFactor); err != nil {
+ return nil, err
}
return ms.createBlockVolumeResponseFromEntry(entry), nil
}
@@ -49,6 +68,9 @@ func (ms *MasterServer) CreateBlockVolume(ctx context.Context, req *master_pb.Cr
// Double-check after acquiring lock (another goroutine may have finished).
if entry, ok := ms.blockRegistry.Lookup(req.Name); ok {
+ if err := ms.validateIdempotentCreate(entry, req, durMode, replicaFactor); err != nil {
+ return nil, err
+ }
return ms.createBlockVolumeResponseFromEntry(entry), nil
}
@@ -71,7 +93,7 @@ func (ms *MasterServer) CreateBlockVolume(ctx context.Context, req *master_pb.Cr
return nil, err
}
- result, err := ms.blockVSAllocate(ctx, server, req.Name, req.SizeBytes, req.DiskType)
+ result, err := ms.blockVSAllocate(ctx, server, req.Name, req.SizeBytes, req.DiskType, req.DurabilityMode)
if err != nil {
lastErr = fmt.Errorf("server %s: %w", server, err)
glog.V(0).Infof("[reqID=%s] CreateBlockVolume %q: attempt %d on %s failed: %v", blockReqID(ctx), req.Name, attempt+1, server, err)
@@ -79,12 +101,6 @@ func (ms *MasterServer) CreateBlockVolume(ctx context.Context, req *master_pb.Cr
continue
}
- // CP8-2: determine replica factor from request (default 2).
- replicaFactor := 2
- if req.ReplicaFactor > 0 && req.ReplicaFactor <= 3 {
- replicaFactor = int(req.ReplicaFactor)
- }
-
entry := &BlockVolumeEntry{
Name: req.Name,
VolumeServer: server,
@@ -96,6 +112,7 @@ func (ms *MasterServer) CreateBlockVolume(ctx context.Context, req *master_pb.Cr
Role: blockvol.RoleToWire(blockvol.RolePrimary),
Status: StatusActive,
ReplicaFactor: replicaFactor,
+ DurabilityMode: durMode.String(),
LeaseTTL: 30 * time.Second,
LastLeaseGrant: time.Now(), // R2-F1: set BEFORE Register to avoid stale-lease race
}
@@ -111,6 +128,15 @@ func (ms *MasterServer) CreateBlockVolume(ctx context.Context, req *master_pb.Cr
if len(entry.Replicas) == 0 && replicaFactor > 1 {
glog.V(0).Infof("[reqID=%s] CreateBlockVolume %q: single-copy mode (replica allocation failed)", blockReqID(ctx), req.Name)
}
+
+ // F1: strict modes require minimum replicas at create time.
+ requiredReplicas := durMode.RequiredReplicas(replicaFactor)
+ if len(entry.Replicas) < requiredReplicas {
+ ms.cleanupPartialCreate(ctx, entry)
+ return nil, fmt.Errorf("durability mode %q requires %d replicas but only %d provisioned",
+ durMode.String(), requiredReplicas, len(entry.Replicas))
+ }
+
// Sync deprecated scalar fields from first replica.
if len(entry.Replicas) > 0 {
r0 := &entry.Replicas[0]
@@ -225,6 +251,10 @@ func (ms *MasterServer) LookupBlockVolume(ctx context.Context, req *master_pb.Lo
if rf == 0 {
rf = 2 // default for pre-CP8-2 entries
}
+ durModeStr := entry.DurabilityMode
+ if durModeStr == "" {
+ durModeStr = "best_effort"
+ }
return &master_pb.LookupBlockVolumeResponse{
VolumeServer: entry.VolumeServer,
IscsiAddr: entry.ISCSIAddr,
@@ -233,6 +263,7 @@ func (ms *MasterServer) LookupBlockVolume(ctx context.Context, req *master_pb.Lo
ReplicaServer: entry.ReplicaServer, // backward compat
ReplicaFactor: uint32(rf),
ReplicaServers: replicaServers,
+ DurabilityMode: durModeStr,
}, nil
}
@@ -240,7 +271,7 @@ func (ms *MasterServer) LookupBlockVolume(ctx context.Context, req *master_pb.Lo
// Returns the replica server address on success, or empty string on failure (F4).
func (ms *MasterServer) tryCreateOneReplica(ctx context.Context, req *master_pb.CreateBlockVolumeRequest, entry *BlockVolumeEntry, primaryResult *blockAllocResult, candidates []string) string {
for _, replicaServerStr := range candidates {
- replicaResult, err := ms.blockVSAllocate(ctx, replicaServerStr, req.Name, req.SizeBytes, req.DiskType)
+ replicaResult, err := ms.blockVSAllocate(ctx, replicaServerStr, req.Name, req.SizeBytes, req.DiskType, req.DurabilityMode)
if err != nil {
glog.V(0).Infof("[reqID=%s] CreateBlockVolume %q: replica on %s failed: %v", blockReqID(ctx), req.Name, replicaServerStr, err)
continue
@@ -381,6 +412,31 @@ func (ms *MasterServer) createBlockVolumeResponseFromEntry(entry *BlockVolumeEnt
}
}
+// validateIdempotentCreate checks that an idempotent create request is consistent
+// with an existing entry. Returns nil if compatible, error on mismatch.
+func (ms *MasterServer) validateIdempotentCreate(entry *BlockVolumeEntry, req *master_pb.CreateBlockVolumeRequest, durMode blockvol.DurabilityMode, replicaFactor int) error {
+ if entry.SizeBytes < req.SizeBytes {
+ return fmt.Errorf("block volume %q exists with size %d (requested %d)", req.Name, entry.SizeBytes, req.SizeBytes)
+ }
+ // Validate durability mode consistency.
+ existingMode := entry.DurabilityMode
+ if existingMode == "" {
+ existingMode = "best_effort"
+ }
+ if durMode.String() != existingMode {
+ return fmt.Errorf("block volume %q exists with durability_mode %q (requested %q)", req.Name, existingMode, durMode.String())
+ }
+ // Validate replica factor consistency.
+ existingRF := entry.ReplicaFactor
+ if existingRF == 0 {
+ existingRF = 2 // default
+ }
+ if replicaFactor != existingRF {
+ return fmt.Errorf("block volume %q exists with replica_factor %d (requested %d)", req.Name, existingRF, replicaFactor)
+ }
+ return nil
+}
+
// replicaServerList returns the list of replica server addresses.
// Order matches Replicas[] (append-order), ensuring ReplicaServers[0] == ReplicaServer (legacy).
func replicaServerList(entry *BlockVolumeEntry) []string {
@@ -404,3 +460,23 @@ func removeServer(servers []string, server string) []string {
}
return result
}
+
+// cleanupPartialCreate removes a partially created block volume (primary + any replicas)
+// when strict durability mode enforcement fails due to insufficient replicas.
+// All operations are best-effort: failures are logged but do not propagate.
+func (ms *MasterServer) cleanupPartialCreate(ctx context.Context, entry *BlockVolumeEntry) {
+ // Delete primary volume.
+ if err := ms.blockVSDelete(ctx, entry.VolumeServer, entry.Name); err != nil {
+ glog.Warningf("[reqID=%s] cleanupPartialCreate %q: delete primary on %s: %v",
+ blockReqID(ctx), entry.Name, entry.VolumeServer, err)
+ }
+ // Delete any successfully created replicas.
+ for _, ri := range entry.Replicas {
+ if err := ms.blockVSDelete(ctx, ri.Server, entry.Name); err != nil {
+ glog.Warningf("[reqID=%s] cleanupPartialCreate %q: delete replica on %s: %v",
+ blockReqID(ctx), entry.Name, ri.Server, err)
+ }
+ }
+ // Remove from registry if somehow registered (shouldn't be at this point).
+ ms.blockRegistry.Unregister(entry.Name)
+}
diff --git a/weed/server/master_grpc_server_block_test.go b/weed/server/master_grpc_server_block_test.go
index c7b4f750c..52c9fd259 100644
--- a/weed/server/master_grpc_server_block_test.go
+++ b/weed/server/master_grpc_server_block_test.go
@@ -19,7 +19,7 @@ func testMasterServer(t *testing.T) *MasterServer {
blockAssignmentQueue: NewBlockAssignmentQueue(),
}
// Default mock: succeed with deterministic values.
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
IQN: fmt.Sprintf("iqn.2024.test:%s", name),
@@ -139,7 +139,7 @@ func TestMaster_CreateVSFailure_Retry(t *testing.T) {
ms.blockRegistry.MarkBlockCapable("vs2:9333")
var callCount atomic.Int32
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
n := callCount.Add(1)
if n == 1 {
return nil, fmt.Errorf("disk full")
@@ -170,7 +170,7 @@ func TestMaster_CreateVSFailure_Cleanup(t *testing.T) {
ms := testMasterServer(t)
ms.blockRegistry.MarkBlockCapable("vs1:9333")
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
return nil, fmt.Errorf("all servers broken")
}
@@ -193,7 +193,7 @@ func TestMaster_CreateConcurrentSameName(t *testing.T) {
ms.blockRegistry.MarkBlockCapable("vs1:9333")
var callCount atomic.Int32
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
callCount.Add(1)
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
@@ -275,7 +275,7 @@ func TestMaster_CreateWithReplica(t *testing.T) {
ms.blockRegistry.MarkBlockCapable("vs2:9333")
var allocServers []string
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
allocServers = append(allocServers, server)
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
@@ -328,7 +328,7 @@ func TestMaster_CreateSingleServer_NoReplica(t *testing.T) {
ms.blockRegistry.MarkBlockCapable("vs1:9333")
var allocCount atomic.Int32
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
allocCount.Add(1)
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
@@ -365,7 +365,7 @@ func TestMaster_CreateReplica_SecondFails_SingleCopy(t *testing.T) {
ms.blockRegistry.MarkBlockCapable("vs2:9333")
var callCount atomic.Int32
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
n := callCount.Add(1)
if n == 2 {
// Replica allocation fails.
@@ -402,7 +402,7 @@ func TestMaster_CreateEnqueuesAssignments(t *testing.T) {
ms.blockRegistry.MarkBlockCapable("vs1:9333")
ms.blockRegistry.MarkBlockCapable("vs2:9333")
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
IQN: fmt.Sprintf("iqn.2024.test:%s", name),
@@ -463,7 +463,7 @@ func TestMaster_LookupReturnsReplicaServer(t *testing.T) {
ms.blockRegistry.MarkBlockCapable("vs1:9333")
ms.blockRegistry.MarkBlockCapable("vs2:9333")
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
IQN: fmt.Sprintf("iqn.2024.test:%s", name),
@@ -670,7 +670,7 @@ func TestMaster_LookupBlockVolume(t *testing.T) {
func testMasterServerRF3(t *testing.T) *MasterServer {
t.Helper()
ms := testMasterServer(t)
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
IQN: fmt.Sprintf("iqn.2024.test:%s", name),
@@ -725,7 +725,7 @@ func TestMaster_CreateRF3_ThreeServers(t *testing.T) {
// RF=3 with only 2 servers: should create 1 replica (partial).
func TestMaster_CreateRF3_TwoServers(t *testing.T) {
ms := testMasterServer(t)
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
IQN: fmt.Sprintf("iqn.2024.test:%s", name),
diff --git a/weed/server/master_server.go b/weed/server/master_server.go
index 35d1766e4..197a7f292 100644
--- a/weed/server/master_server.go
+++ b/weed/server/master_server.go
@@ -98,7 +98,7 @@ type MasterServer struct {
blockRegistry *BlockVolumeRegistry
blockAssignmentQueue *BlockAssignmentQueue
blockFailover *blockFailoverState
- blockVSAllocate func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error)
+ blockVSAllocate func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error)
blockVSDelete func(ctx context.Context, server string, name string) error
blockVSSnapshot func(ctx context.Context, server string, name string, snapID uint32) (int64, uint64, error)
blockVSDeleteSnap func(ctx context.Context, server string, name string, snapID uint32) error
@@ -553,13 +553,14 @@ type blockAllocResult struct {
}
// defaultBlockVSAllocate calls a volume server's AllocateBlockVolume RPC.
-func (ms *MasterServer) defaultBlockVSAllocate(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+func (ms *MasterServer) defaultBlockVSAllocate(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
var result blockAllocResult
err := operation.WithVolumeServerClient(false, pb.ServerAddress(server), ms.grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
resp, rerr := client.AllocateBlockVolume(ctx, &volume_server_pb.AllocateBlockVolumeRequest{
- Name: name,
- SizeBytes: sizeBytes,
- DiskType: diskType,
+ Name: name,
+ SizeBytes: sizeBytes,
+ DiskType: diskType,
+ DurabilityMode: durabilityMode,
})
if rerr != nil {
return rerr
diff --git a/weed/server/master_server_handlers_block.go b/weed/server/master_server_handlers_block.go
index 7d1285724..f8d29f011 100644
--- a/weed/server/master_server_handlers_block.go
+++ b/weed/server/master_server_handlers_block.go
@@ -26,10 +26,20 @@ func (ms *MasterServer) blockVolumeCreateHandler(w http.ResponseWriter, r *http.
replicaPlacement = "000"
}
+ // Pre-validate durability_mode (cosmetic — real validation is in gRPC handler).
+ if req.DurabilityMode != "" {
+ if _, perr := blockvol.ParseDurabilityMode(req.DurabilityMode); perr != nil {
+ writeJsonError(w, r, http.StatusBadRequest, fmt.Errorf("invalid durability_mode: %w", perr))
+ return
+ }
+ }
+
resp, err := ms.CreateBlockVolume(r.Context(), &master_pb.CreateBlockVolumeRequest{
- Name: req.Name,
- SizeBytes: req.SizeBytes,
- DiskType: req.DiskType,
+ Name: req.Name,
+ SizeBytes: req.SizeBytes,
+ DiskType: req.DiskType,
+ DurabilityMode: req.DurabilityMode,
+ ReplicaFactor: uint32(req.ReplicaFactor),
})
if err != nil {
writeJsonError(w, r, http.StatusInternalServerError, err)
@@ -177,6 +187,10 @@ func entryToVolumeInfo(e *BlockVolumeEntry) blockapi.VolumeInfo {
if rf == 0 {
rf = 2 // default
}
+ durMode := e.DurabilityMode
+ if durMode == "" {
+ durMode = "best_effort"
+ }
info := blockapi.VolumeInfo{
Name: e.Name,
VolumeServer: e.VolumeServer,
@@ -195,6 +209,7 @@ func entryToVolumeInfo(e *BlockVolumeEntry) blockapi.VolumeInfo {
ReplicaFactor: rf,
HealthScore: e.HealthScore,
ReplicaDegraded: e.ReplicaDegraded,
+ DurabilityMode: durMode,
}
for _, ri := range e.Replicas {
info.Replicas = append(info.Replicas, blockapi.ReplicaDetail{
diff --git a/weed/server/master_server_handlers_block_test.go b/weed/server/master_server_handlers_block_test.go
index c2f309650..89395ba37 100644
--- a/weed/server/master_server_handlers_block_test.go
+++ b/weed/server/master_server_handlers_block_test.go
@@ -20,7 +20,7 @@ func blockTestServer(t *testing.T) (*MasterServer, *httptest.Server) {
blockRegistry: NewBlockVolumeRegistry(),
blockAssignmentQueue: NewBlockAssignmentQueue(),
}
- ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string) (*blockAllocResult, error) {
+ ms.blockVSAllocate = func(ctx context.Context, server string, name string, sizeBytes uint64, diskType string, durabilityMode string) (*blockAllocResult, error) {
return &blockAllocResult{
Path: fmt.Sprintf("/data/%s.blk", name),
IQN: fmt.Sprintf("iqn.2024.test:%s", name),
diff --git a/weed/server/master_server_handlers_block_ui.go b/weed/server/master_server_handlers_block_ui.go
index be0ca0ae5..5e5c0204a 100644
--- a/weed/server/master_server_handlers_block_ui.go
+++ b/weed/server/master_server_handlers_block_ui.go
@@ -13,7 +13,9 @@ var blockOpsTemplate = template.Must(template.New("blockOps").Parse(blockLayoutH
type blockUIVolume struct {
blockapi.VolumeInfo
- SizeMB uint64
+ SizeMB uint64
+ WALHeadLSN uint64
+ MaxWALLag uint64 // max WAL lag across replicas
}
type blockUIData struct {
@@ -25,6 +27,12 @@ type blockUIData struct {
ActiveCount int
PendingCount int
TotalSizeMB uint64
+ // Observability (CP8-4)
+ BarrierLagLSN uint64
+ PromotionsTotal uint64
+ FailoversTotal uint64
+ RebuildsTotal uint64
+ AssignmentQueueLen int
}
func (ms *MasterServer) buildBlockUIData(tab string) blockUIData {
@@ -35,9 +43,17 @@ func (ms *MasterServer) buildBlockUIData(tab string) blockUIData {
for i, e := range entries {
info := entryToVolumeInfo(e)
mb := info.SizeBytes / (1024 * 1024)
+ var maxLag uint64
+ for _, ri := range e.Replicas {
+ if ri.WALLag > maxLag {
+ maxLag = ri.WALLag
+ }
+ }
volumes[i] = blockUIVolume{
VolumeInfo: info,
SizeMB: mb,
+ WALHeadLSN: e.WALHeadLSN,
+ MaxWALLag: maxLag,
}
totalSizeMB += mb
if e.Status == StatusActive {
@@ -58,13 +74,18 @@ func (ms *MasterServer) buildBlockUIData(tab string) blockUIData {
}
return blockUIData{
- Tab: tab,
- Volumes: volumes,
- Servers: servers,
- TotalVolumes: len(entries),
- ActiveCount: activeCount,
- PendingCount: pendingCount,
- TotalSizeMB: totalSizeMB,
+ Tab: tab,
+ Volumes: volumes,
+ Servers: servers,
+ TotalVolumes: len(entries),
+ ActiveCount: activeCount,
+ PendingCount: pendingCount,
+ TotalSizeMB: totalSizeMB,
+ BarrierLagLSN: ms.blockRegistry.MaxBarrierLagLSN(),
+ PromotionsTotal: ms.blockRegistry.PromotionsTotal.Load(),
+ FailoversTotal: ms.blockRegistry.FailoversTotal.Load(),
+ RebuildsTotal: ms.blockRegistry.RebuildsTotal.Load(),
+ AssignmentQueueLen: ms.blockAssignmentQueue.TotalPending(),
}
}
@@ -130,6 +151,9 @@ const blockLayoutHTML = `
.badge-primary { background: #dfe6e9; color: #2d3436; }
.badge-replica { background: #e8daef; color: #6c3483; }
.empty { color: #b2bec3; font-style: italic; padding: 20px; text-align: center; }
+ .card .value.red { color: #d63031; }
+ .card .value.gray { color: #636e72; }
+ .section-label { font-size: 11px; color: #b2bec3; text-transform: uppercase; letter-spacing: 1px; margin-bottom: 8px; }
@@ -166,6 +190,30 @@ const blockDashContentHTML = `
+Cluster Health
+
+
+
Barrier Lag LSN
+
{{.BarrierLagLSN}}
+
+
+
Promotions
+
{{.PromotionsTotal}}
+
+
+
Failovers
+
{{.FailoversTotal}}
+
+
+
Rebuilds
+
{{.RebuildsTotal}}
+
+
+
Queue Depth
+
{{.AssignmentQueueLen}}
+
+
+
Servers
{{if .Servers}}
@@ -186,19 +234,23 @@ const blockDashContentHTML = `
{{if .Volumes}}