diff --git a/seaweed-volume/proto/master.proto b/seaweed-volume/proto/master.proto index e9279f911..642d87afe 100644 --- a/seaweed-volume/proto/master.proto +++ b/seaweed-volume/proto/master.proto @@ -64,6 +64,9 @@ service Seaweed { message DiskTag { uint32 disk_id = 1; repeated string tags = 2; + // Physical disk descriptor, reported for every location including empty ones. + string type = 3; + int64 max_volume_count = 4; } message Heartbeat { @@ -206,6 +209,17 @@ message ClusterNodeUpdate { message KeepConnectedResponse { VolumeLocation volume_location = 1; ClusterNodeUpdate cluster_node_update = 2; + LockRingUpdate lock_ring_update = 3; +} + +// LockRingUpdate is sent by the master to all filers when the lock ring +// membership changes. The master batches rapid changes (e.g., node drop + join) +// and sends the complete member list atomically, avoiding intermediate ring +// states that would cause unnecessary lock churn. +message LockRingUpdate { + string filer_group = 1; + repeated string servers = 2; + int64 version = 3; } message LookupVolumeRequest { @@ -308,8 +322,15 @@ message DiskInfo { repeated VolumeInformationMessage volume_infos = 6; repeated VolumeEcShardInformationMessage ec_shard_infos = 7; int64 remote_volume_count = 8; + // On a per-physical-disk DiskInfo (from SplitByPhysicalDisk) this is the disk's + // identity; on the type-keyed aggregate it is only a representative fallback + // (the first volume's disk id). uint32 disk_id = 9; repeated string tags = 10; + // Max volume count for every physical disk of this type, keyed by disk id, + // including disks with no volumes or EC shards; recovers empty disks that + // carry no per-volume/per-shard records. + map max_volume_count_by_disk = 11; } message DataNodeInfo { string id = 1; @@ -360,11 +381,13 @@ message VacuumVolumeResponse { } message DisableVacuumRequest { + bool by_plugin = 1; } message DisableVacuumResponse { } message EnableVacuumRequest { + bool by_plugin = 1; } message EnableVacuumResponse { } diff --git a/seaweed-volume/src/server/heartbeat.rs b/seaweed-volume/src/server/heartbeat.rs index 00070960f..39ef24461 100644 --- a/seaweed-volume/src/server/heartbeat.rs +++ b/seaweed-volume/src/server/heartbeat.rs @@ -625,7 +625,8 @@ async fn send_deregister_heartbeat( ) { let empty = { let store = state.store.read().unwrap(); - let (location_uuids, disk_tags) = collect_location_metadata(&store); + // Deregister: no effective max computed, fall back to configured max. + let (location_uuids, disk_tags) = collect_location_metadata(&store, &[]); master_pb::Heartbeat { id: store.id.clone(), ip: config.ip.clone(), @@ -747,7 +748,10 @@ fn collect_heartbeat( ) } -fn collect_location_metadata(store: &Store) -> (Vec, Vec) { +fn collect_location_metadata( + store: &Store, + disk_max_by_id: &[i32], +) -> (Vec, Vec) { let location_uuids = store .locations .iter() @@ -757,9 +761,17 @@ fn collect_location_metadata(store: &Store) -> (Vec, Vec max_volume_count_by_disk = 11; } message DataNodeInfo { string id = 1; diff --git a/weed/pb/master_pb/master.pb.go b/weed/pb/master_pb/master.pb.go index efdb66380..4ad4f4990 100644 --- a/weed/pb/master_pb/master.pb.go +++ b/weed/pb/master_pb/master.pb.go @@ -1,7 +1,7 @@ // Code generated by protoc-gen-go. DO NOT EDIT. // versions: // protoc-gen-go v1.36.6 -// protoc v6.33.4 +// protoc v7.35.0 // source: master.proto package master_pb @@ -23,11 +23,14 @@ const ( ) type DiskTag struct { - state protoimpl.MessageState `protogen:"open.v1"` - DiskId uint32 `protobuf:"varint,1,opt,name=disk_id,json=diskId,proto3" json:"disk_id,omitempty"` - Tags []string `protobuf:"bytes,2,rep,name=tags,proto3" json:"tags,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + state protoimpl.MessageState `protogen:"open.v1"` + DiskId uint32 `protobuf:"varint,1,opt,name=disk_id,json=diskId,proto3" json:"disk_id,omitempty"` + Tags []string `protobuf:"bytes,2,rep,name=tags,proto3" json:"tags,omitempty"` + // Physical disk descriptor, reported for every location including empty ones. + Type string `protobuf:"bytes,3,opt,name=type,proto3" json:"type,omitempty"` + MaxVolumeCount int64 `protobuf:"varint,4,opt,name=max_volume_count,json=maxVolumeCount,proto3" json:"max_volume_count,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *DiskTag) Reset() { @@ -74,6 +77,20 @@ func (x *DiskTag) GetTags() []string { return nil } +func (x *DiskTag) GetType() string { + if x != nil { + return x.Type + } + return "" +} + +func (x *DiskTag) GetMaxVolumeCount() int64 { + if x != nil { + return x.MaxVolumeCount + } + return 0 +} + type Heartbeat struct { state protoimpl.MessageState `protogen:"open.v1"` Ip string `protobuf:"bytes,1,opt,name=ip,proto3" json:"ip,omitempty"` @@ -2121,10 +2138,17 @@ type DiskInfo struct { VolumeInfos []*VolumeInformationMessage `protobuf:"bytes,6,rep,name=volume_infos,json=volumeInfos,proto3" json:"volume_infos,omitempty"` EcShardInfos []*VolumeEcShardInformationMessage `protobuf:"bytes,7,rep,name=ec_shard_infos,json=ecShardInfos,proto3" json:"ec_shard_infos,omitempty"` RemoteVolumeCount int64 `protobuf:"varint,8,opt,name=remote_volume_count,json=remoteVolumeCount,proto3" json:"remote_volume_count,omitempty"` - DiskId uint32 `protobuf:"varint,9,opt,name=disk_id,json=diskId,proto3" json:"disk_id,omitempty"` - Tags []string `protobuf:"bytes,10,rep,name=tags,proto3" json:"tags,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + // On a per-physical-disk DiskInfo (from SplitByPhysicalDisk) this is the disk's + // identity; on the type-keyed aggregate it is only a representative fallback + // (the first volume's disk id). + DiskId uint32 `protobuf:"varint,9,opt,name=disk_id,json=diskId,proto3" json:"disk_id,omitempty"` + Tags []string `protobuf:"bytes,10,rep,name=tags,proto3" json:"tags,omitempty"` + // Max volume count for every physical disk of this type, keyed by disk id, + // including disks with no volumes or EC shards; recovers empty disks that + // carry no per-volume/per-shard records. + MaxVolumeCountByDisk map[uint32]int64 `protobuf:"bytes,11,rep,name=max_volume_count_by_disk,json=maxVolumeCountByDisk,proto3" json:"max_volume_count_by_disk,omitempty" protobuf_key:"varint,1,opt,name=key" protobuf_val:"varint,2,opt,name=value"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *DiskInfo) Reset() { @@ -2227,6 +2251,13 @@ func (x *DiskInfo) GetTags() []string { return nil } +func (x *DiskInfo) GetMaxVolumeCountByDisk() map[uint32]int64 { + if x != nil { + return x.MaxVolumeCountByDisk + } + return nil +} + type DataNodeInfo struct { state protoimpl.MessageState `protogen:"open.v1"` Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"` @@ -4186,7 +4217,7 @@ type LookupEcVolumeResponse_EcShardIdLocation struct { func (x *LookupEcVolumeResponse_EcShardIdLocation) Reset() { *x = LookupEcVolumeResponse_EcShardIdLocation{} - mi := &file_master_proto_msgTypes[71] + mi := &file_master_proto_msgTypes[72] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4198,7 +4229,7 @@ func (x *LookupEcVolumeResponse_EcShardIdLocation) String() string { func (*LookupEcVolumeResponse_EcShardIdLocation) ProtoMessage() {} func (x *LookupEcVolumeResponse_EcShardIdLocation) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[71] + mi := &file_master_proto_msgTypes[72] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4241,7 +4272,7 @@ type ListClusterNodesResponse_ClusterNode struct { func (x *ListClusterNodesResponse_ClusterNode) Reset() { *x = ListClusterNodesResponse_ClusterNode{} - mi := &file_master_proto_msgTypes[72] + mi := &file_master_proto_msgTypes[73] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4253,7 +4284,7 @@ func (x *ListClusterNodesResponse_ClusterNode) String() string { func (*ListClusterNodesResponse_ClusterNode) ProtoMessage() {} func (x *ListClusterNodesResponse_ClusterNode) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[72] + mi := &file_master_proto_msgTypes[73] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4316,7 +4347,7 @@ type RaftListClusterServersResponse_ClusterServers struct { func (x *RaftListClusterServersResponse_ClusterServers) Reset() { *x = RaftListClusterServersResponse_ClusterServers{} - mi := &file_master_proto_msgTypes[73] + mi := &file_master_proto_msgTypes[74] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4328,7 +4359,7 @@ func (x *RaftListClusterServersResponse_ClusterServers) String() string { func (*RaftListClusterServersResponse_ClusterServers) ProtoMessage() {} func (x *RaftListClusterServersResponse_ClusterServers) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[73] + mi := &file_master_proto_msgTypes[74] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4376,10 +4407,12 @@ var File_master_proto protoreflect.FileDescriptor const file_master_proto_rawDesc = "" + "\n" + - "\fmaster.proto\x12\tmaster_pb\x1a\x13volume_server.proto\"6\n" + + "\fmaster.proto\x12\tmaster_pb\x1a\x13volume_server.proto\"t\n" + "\aDiskTag\x12\x17\n" + "\adisk_id\x18\x01 \x01(\rR\x06diskId\x12\x12\n" + - "\x04tags\x18\x02 \x03(\tR\x04tags\"\xbc\b\n" + + "\x04tags\x18\x02 \x03(\tR\x04tags\x12\x12\n" + + "\x04type\x18\x03 \x01(\tR\x04type\x12(\n" + + "\x10max_volume_count\x18\x04 \x01(\x03R\x0emaxVolumeCount\"\xbc\b\n" + "\tHeartbeat\x12\x0e\n" + "\x02ip\x18\x01 \x01(\tR\x02ip\x12\x12\n" + "\x04port\x18\x02 \x01(\rR\x04port\x12\x1d\n" + @@ -4600,7 +4633,7 @@ const file_master_proto_rawDesc = "" + "\vcollections\x18\x01 \x03(\v2\x15.master_pb.CollectionR\vcollections\"-\n" + "\x17CollectionDeleteRequest\x12\x12\n" + "\x04name\x18\x01 \x01(\tR\x04name\"\x1a\n" + - "\x18CollectionDeleteResponse\"\xbe\x03\n" + + "\x18CollectionDeleteResponse\"\xee\x04\n" + "\bDiskInfo\x12\x12\n" + "\x04type\x18\x01 \x01(\tR\x04type\x12!\n" + "\fvolume_count\x18\x02 \x01(\x03R\vvolumeCount\x12(\n" + @@ -4612,7 +4645,11 @@ const file_master_proto_rawDesc = "" + "\x13remote_volume_count\x18\b \x01(\x03R\x11remoteVolumeCount\x12\x17\n" + "\adisk_id\x18\t \x01(\rR\x06diskId\x12\x12\n" + "\x04tags\x18\n" + - " \x03(\tR\x04tags\"\xee\x01\n" + + " \x03(\tR\x04tags\x12e\n" + + "\x18max_volume_count_by_disk\x18\v \x03(\v2-.master_pb.DiskInfo.MaxVolumeCountByDiskEntryR\x14maxVolumeCountByDisk\x1aG\n" + + "\x19MaxVolumeCountByDiskEntry\x12\x10\n" + + "\x03key\x18\x01 \x01(\rR\x03key\x12\x14\n" + + "\x05value\x18\x02 \x01(\x03R\x05value:\x028\x01\"\xee\x01\n" + "\fDataNodeInfo\x12\x0e\n" + "\x02id\x18\x01 \x01(\tR\x02id\x12D\n" + "\tdiskInfos\x18\x02 \x03(\v2&.master_pb.DataNodeInfo.DiskInfosEntryR\tdiskInfos\x12\x1b\n" + @@ -4800,7 +4837,7 @@ func file_master_proto_rawDescGZIP() []byte { return file_master_proto_rawDescData } -var file_master_proto_msgTypes = make([]protoimpl.MessageInfo, 74) +var file_master_proto_msgTypes = make([]protoimpl.MessageInfo, 75) var file_master_proto_goTypes = []any{ (*DiskTag)(nil), // 0: master_pb.DiskTag (*Heartbeat)(nil), // 1: master_pb.Heartbeat @@ -4869,14 +4906,15 @@ var file_master_proto_goTypes = []any{ nil, // 64: master_pb.StorageBackend.PropertiesEntry (*SuperBlockExtra_ErasureCoding)(nil), // 65: master_pb.SuperBlockExtra.ErasureCoding (*LookupVolumeResponse_VolumeIdLocation)(nil), // 66: master_pb.LookupVolumeResponse.VolumeIdLocation - nil, // 67: master_pb.DataNodeInfo.DiskInfosEntry - nil, // 68: master_pb.RackInfo.DiskInfosEntry - nil, // 69: master_pb.DataCenterInfo.DiskInfosEntry - nil, // 70: master_pb.TopologyInfo.DiskInfosEntry - (*LookupEcVolumeResponse_EcShardIdLocation)(nil), // 71: master_pb.LookupEcVolumeResponse.EcShardIdLocation - (*ListClusterNodesResponse_ClusterNode)(nil), // 72: master_pb.ListClusterNodesResponse.ClusterNode - (*RaftListClusterServersResponse_ClusterServers)(nil), // 73: master_pb.RaftListClusterServersResponse.ClusterServers - (*volume_server_pb.VolumeServerState)(nil), // 74: volume_server_pb.VolumeServerState + nil, // 67: master_pb.DiskInfo.MaxVolumeCountByDiskEntry + nil, // 68: master_pb.DataNodeInfo.DiskInfosEntry + nil, // 69: master_pb.RackInfo.DiskInfosEntry + nil, // 70: master_pb.DataCenterInfo.DiskInfosEntry + nil, // 71: master_pb.TopologyInfo.DiskInfosEntry + (*LookupEcVolumeResponse_EcShardIdLocation)(nil), // 72: master_pb.LookupEcVolumeResponse.EcShardIdLocation + (*ListClusterNodesResponse_ClusterNode)(nil), // 73: master_pb.ListClusterNodesResponse.ClusterNode + (*RaftListClusterServersResponse_ClusterServers)(nil), // 74: master_pb.RaftListClusterServersResponse.ClusterServers + (*volume_server_pb.VolumeServerState)(nil), // 75: volume_server_pb.VolumeServerState } var file_master_proto_depIdxs = []int32{ 3, // 0: master_pb.Heartbeat.volumes:type_name -> master_pb.VolumeInformationMessage @@ -4886,7 +4924,7 @@ var file_master_proto_depIdxs = []int32{ 5, // 4: master_pb.Heartbeat.new_ec_shards:type_name -> master_pb.VolumeEcShardInformationMessage 5, // 5: master_pb.Heartbeat.deleted_ec_shards:type_name -> master_pb.VolumeEcShardInformationMessage 63, // 6: master_pb.Heartbeat.max_volume_counts:type_name -> master_pb.Heartbeat.MaxVolumeCountsEntry - 74, // 7: master_pb.Heartbeat.state:type_name -> volume_server_pb.VolumeServerState + 75, // 7: master_pb.Heartbeat.state:type_name -> volume_server_pb.VolumeServerState 0, // 8: master_pb.Heartbeat.disk_tags:type_name -> master_pb.DiskTag 6, // 9: master_pb.HeartbeatResponse.storage_backends:type_name -> master_pb.StorageBackend 64, // 10: master_pb.StorageBackend.properties:type_name -> master_pb.StorageBackend.PropertiesEntry @@ -4900,77 +4938,78 @@ var file_master_proto_depIdxs = []int32{ 22, // 18: master_pb.CollectionListResponse.collections:type_name -> master_pb.Collection 3, // 19: master_pb.DiskInfo.volume_infos:type_name -> master_pb.VolumeInformationMessage 5, // 20: master_pb.DiskInfo.ec_shard_infos:type_name -> master_pb.VolumeEcShardInformationMessage - 67, // 21: master_pb.DataNodeInfo.diskInfos:type_name -> master_pb.DataNodeInfo.DiskInfosEntry - 28, // 22: master_pb.RackInfo.data_node_infos:type_name -> master_pb.DataNodeInfo - 68, // 23: master_pb.RackInfo.diskInfos:type_name -> master_pb.RackInfo.DiskInfosEntry - 29, // 24: master_pb.DataCenterInfo.rack_infos:type_name -> master_pb.RackInfo - 69, // 25: master_pb.DataCenterInfo.diskInfos:type_name -> master_pb.DataCenterInfo.DiskInfosEntry - 30, // 26: master_pb.TopologyInfo.data_center_infos:type_name -> master_pb.DataCenterInfo - 70, // 27: master_pb.TopologyInfo.diskInfos:type_name -> master_pb.TopologyInfo.DiskInfosEntry - 31, // 28: master_pb.VolumeListResponse.topology_info:type_name -> master_pb.TopologyInfo - 71, // 29: master_pb.LookupEcVolumeResponse.shard_id_locations:type_name -> master_pb.LookupEcVolumeResponse.EcShardIdLocation - 6, // 30: master_pb.GetMasterConfigurationResponse.storage_backends:type_name -> master_pb.StorageBackend - 72, // 31: master_pb.ListClusterNodesResponse.cluster_nodes:type_name -> master_pb.ListClusterNodesResponse.ClusterNode - 73, // 32: master_pb.RaftListClusterServersResponse.cluster_servers:type_name -> master_pb.RaftListClusterServersResponse.ClusterServers - 16, // 33: master_pb.LookupVolumeResponse.VolumeIdLocation.locations:type_name -> master_pb.Location - 27, // 34: master_pb.DataNodeInfo.DiskInfosEntry.value:type_name -> master_pb.DiskInfo - 27, // 35: master_pb.RackInfo.DiskInfosEntry.value:type_name -> master_pb.DiskInfo - 27, // 36: master_pb.DataCenterInfo.DiskInfosEntry.value:type_name -> master_pb.DiskInfo - 27, // 37: master_pb.TopologyInfo.DiskInfosEntry.value:type_name -> master_pb.DiskInfo - 16, // 38: master_pb.LookupEcVolumeResponse.EcShardIdLocation.locations:type_name -> master_pb.Location - 1, // 39: master_pb.Seaweed.SendHeartbeat:input_type -> master_pb.Heartbeat - 9, // 40: master_pb.Seaweed.KeepConnected:input_type -> master_pb.KeepConnectedRequest - 14, // 41: master_pb.Seaweed.LookupVolume:input_type -> master_pb.LookupVolumeRequest - 17, // 42: master_pb.Seaweed.Assign:input_type -> master_pb.AssignRequest - 17, // 43: master_pb.Seaweed.StreamAssign:input_type -> master_pb.AssignRequest - 20, // 44: master_pb.Seaweed.Statistics:input_type -> master_pb.StatisticsRequest - 23, // 45: master_pb.Seaweed.CollectionList:input_type -> master_pb.CollectionListRequest - 25, // 46: master_pb.Seaweed.CollectionDelete:input_type -> master_pb.CollectionDeleteRequest - 32, // 47: master_pb.Seaweed.VolumeList:input_type -> master_pb.VolumeListRequest - 34, // 48: master_pb.Seaweed.LookupEcVolume:input_type -> master_pb.LookupEcVolumeRequest - 36, // 49: master_pb.Seaweed.VacuumVolume:input_type -> master_pb.VacuumVolumeRequest - 38, // 50: master_pb.Seaweed.DisableVacuum:input_type -> master_pb.DisableVacuumRequest - 40, // 51: master_pb.Seaweed.EnableVacuum:input_type -> master_pb.EnableVacuumRequest - 42, // 52: master_pb.Seaweed.VolumeMarkReadonly:input_type -> master_pb.VolumeMarkReadonlyRequest - 44, // 53: master_pb.Seaweed.GetMasterConfiguration:input_type -> master_pb.GetMasterConfigurationRequest - 46, // 54: master_pb.Seaweed.ListClusterNodes:input_type -> master_pb.ListClusterNodesRequest - 48, // 55: master_pb.Seaweed.LeaseAdminToken:input_type -> master_pb.LeaseAdminTokenRequest - 50, // 56: master_pb.Seaweed.ReleaseAdminToken:input_type -> master_pb.ReleaseAdminTokenRequest - 52, // 57: master_pb.Seaweed.Ping:input_type -> master_pb.PingRequest - 58, // 58: master_pb.Seaweed.RaftListClusterServers:input_type -> master_pb.RaftListClusterServersRequest - 54, // 59: master_pb.Seaweed.RaftAddServer:input_type -> master_pb.RaftAddServerRequest - 56, // 60: master_pb.Seaweed.RaftRemoveServer:input_type -> master_pb.RaftRemoveServerRequest - 60, // 61: master_pb.Seaweed.RaftLeadershipTransfer:input_type -> master_pb.RaftLeadershipTransferRequest - 18, // 62: master_pb.Seaweed.VolumeGrow:input_type -> master_pb.VolumeGrowRequest - 2, // 63: master_pb.Seaweed.SendHeartbeat:output_type -> master_pb.HeartbeatResponse - 12, // 64: master_pb.Seaweed.KeepConnected:output_type -> master_pb.KeepConnectedResponse - 15, // 65: master_pb.Seaweed.LookupVolume:output_type -> master_pb.LookupVolumeResponse - 19, // 66: master_pb.Seaweed.Assign:output_type -> master_pb.AssignResponse - 19, // 67: master_pb.Seaweed.StreamAssign:output_type -> master_pb.AssignResponse - 21, // 68: master_pb.Seaweed.Statistics:output_type -> master_pb.StatisticsResponse - 24, // 69: master_pb.Seaweed.CollectionList:output_type -> master_pb.CollectionListResponse - 26, // 70: master_pb.Seaweed.CollectionDelete:output_type -> master_pb.CollectionDeleteResponse - 33, // 71: master_pb.Seaweed.VolumeList:output_type -> master_pb.VolumeListResponse - 35, // 72: master_pb.Seaweed.LookupEcVolume:output_type -> master_pb.LookupEcVolumeResponse - 37, // 73: master_pb.Seaweed.VacuumVolume:output_type -> master_pb.VacuumVolumeResponse - 39, // 74: master_pb.Seaweed.DisableVacuum:output_type -> master_pb.DisableVacuumResponse - 41, // 75: master_pb.Seaweed.EnableVacuum:output_type -> master_pb.EnableVacuumResponse - 43, // 76: master_pb.Seaweed.VolumeMarkReadonly:output_type -> master_pb.VolumeMarkReadonlyResponse - 45, // 77: master_pb.Seaweed.GetMasterConfiguration:output_type -> master_pb.GetMasterConfigurationResponse - 47, // 78: master_pb.Seaweed.ListClusterNodes:output_type -> master_pb.ListClusterNodesResponse - 49, // 79: master_pb.Seaweed.LeaseAdminToken:output_type -> master_pb.LeaseAdminTokenResponse - 51, // 80: master_pb.Seaweed.ReleaseAdminToken:output_type -> master_pb.ReleaseAdminTokenResponse - 53, // 81: master_pb.Seaweed.Ping:output_type -> master_pb.PingResponse - 59, // 82: master_pb.Seaweed.RaftListClusterServers:output_type -> master_pb.RaftListClusterServersResponse - 55, // 83: master_pb.Seaweed.RaftAddServer:output_type -> master_pb.RaftAddServerResponse - 57, // 84: master_pb.Seaweed.RaftRemoveServer:output_type -> master_pb.RaftRemoveServerResponse - 61, // 85: master_pb.Seaweed.RaftLeadershipTransfer:output_type -> master_pb.RaftLeadershipTransferResponse - 62, // 86: master_pb.Seaweed.VolumeGrow:output_type -> master_pb.VolumeGrowResponse - 63, // [63:87] is the sub-list for method output_type - 39, // [39:63] is the sub-list for method input_type - 39, // [39:39] is the sub-list for extension type_name - 39, // [39:39] is the sub-list for extension extendee - 0, // [0:39] is the sub-list for field type_name + 67, // 21: master_pb.DiskInfo.max_volume_count_by_disk:type_name -> master_pb.DiskInfo.MaxVolumeCountByDiskEntry + 68, // 22: master_pb.DataNodeInfo.diskInfos:type_name -> master_pb.DataNodeInfo.DiskInfosEntry + 28, // 23: master_pb.RackInfo.data_node_infos:type_name -> master_pb.DataNodeInfo + 69, // 24: master_pb.RackInfo.diskInfos:type_name -> master_pb.RackInfo.DiskInfosEntry + 29, // 25: master_pb.DataCenterInfo.rack_infos:type_name -> master_pb.RackInfo + 70, // 26: master_pb.DataCenterInfo.diskInfos:type_name -> master_pb.DataCenterInfo.DiskInfosEntry + 30, // 27: master_pb.TopologyInfo.data_center_infos:type_name -> master_pb.DataCenterInfo + 71, // 28: master_pb.TopologyInfo.diskInfos:type_name -> master_pb.TopologyInfo.DiskInfosEntry + 31, // 29: master_pb.VolumeListResponse.topology_info:type_name -> master_pb.TopologyInfo + 72, // 30: master_pb.LookupEcVolumeResponse.shard_id_locations:type_name -> master_pb.LookupEcVolumeResponse.EcShardIdLocation + 6, // 31: master_pb.GetMasterConfigurationResponse.storage_backends:type_name -> master_pb.StorageBackend + 73, // 32: master_pb.ListClusterNodesResponse.cluster_nodes:type_name -> master_pb.ListClusterNodesResponse.ClusterNode + 74, // 33: master_pb.RaftListClusterServersResponse.cluster_servers:type_name -> master_pb.RaftListClusterServersResponse.ClusterServers + 16, // 34: master_pb.LookupVolumeResponse.VolumeIdLocation.locations:type_name -> master_pb.Location + 27, // 35: master_pb.DataNodeInfo.DiskInfosEntry.value:type_name -> master_pb.DiskInfo + 27, // 36: master_pb.RackInfo.DiskInfosEntry.value:type_name -> master_pb.DiskInfo + 27, // 37: master_pb.DataCenterInfo.DiskInfosEntry.value:type_name -> master_pb.DiskInfo + 27, // 38: master_pb.TopologyInfo.DiskInfosEntry.value:type_name -> master_pb.DiskInfo + 16, // 39: master_pb.LookupEcVolumeResponse.EcShardIdLocation.locations:type_name -> master_pb.Location + 1, // 40: master_pb.Seaweed.SendHeartbeat:input_type -> master_pb.Heartbeat + 9, // 41: master_pb.Seaweed.KeepConnected:input_type -> master_pb.KeepConnectedRequest + 14, // 42: master_pb.Seaweed.LookupVolume:input_type -> master_pb.LookupVolumeRequest + 17, // 43: master_pb.Seaweed.Assign:input_type -> master_pb.AssignRequest + 17, // 44: master_pb.Seaweed.StreamAssign:input_type -> master_pb.AssignRequest + 20, // 45: master_pb.Seaweed.Statistics:input_type -> master_pb.StatisticsRequest + 23, // 46: master_pb.Seaweed.CollectionList:input_type -> master_pb.CollectionListRequest + 25, // 47: master_pb.Seaweed.CollectionDelete:input_type -> master_pb.CollectionDeleteRequest + 32, // 48: master_pb.Seaweed.VolumeList:input_type -> master_pb.VolumeListRequest + 34, // 49: master_pb.Seaweed.LookupEcVolume:input_type -> master_pb.LookupEcVolumeRequest + 36, // 50: master_pb.Seaweed.VacuumVolume:input_type -> master_pb.VacuumVolumeRequest + 38, // 51: master_pb.Seaweed.DisableVacuum:input_type -> master_pb.DisableVacuumRequest + 40, // 52: master_pb.Seaweed.EnableVacuum:input_type -> master_pb.EnableVacuumRequest + 42, // 53: master_pb.Seaweed.VolumeMarkReadonly:input_type -> master_pb.VolumeMarkReadonlyRequest + 44, // 54: master_pb.Seaweed.GetMasterConfiguration:input_type -> master_pb.GetMasterConfigurationRequest + 46, // 55: master_pb.Seaweed.ListClusterNodes:input_type -> master_pb.ListClusterNodesRequest + 48, // 56: master_pb.Seaweed.LeaseAdminToken:input_type -> master_pb.LeaseAdminTokenRequest + 50, // 57: master_pb.Seaweed.ReleaseAdminToken:input_type -> master_pb.ReleaseAdminTokenRequest + 52, // 58: master_pb.Seaweed.Ping:input_type -> master_pb.PingRequest + 58, // 59: master_pb.Seaweed.RaftListClusterServers:input_type -> master_pb.RaftListClusterServersRequest + 54, // 60: master_pb.Seaweed.RaftAddServer:input_type -> master_pb.RaftAddServerRequest + 56, // 61: master_pb.Seaweed.RaftRemoveServer:input_type -> master_pb.RaftRemoveServerRequest + 60, // 62: master_pb.Seaweed.RaftLeadershipTransfer:input_type -> master_pb.RaftLeadershipTransferRequest + 18, // 63: master_pb.Seaweed.VolumeGrow:input_type -> master_pb.VolumeGrowRequest + 2, // 64: master_pb.Seaweed.SendHeartbeat:output_type -> master_pb.HeartbeatResponse + 12, // 65: master_pb.Seaweed.KeepConnected:output_type -> master_pb.KeepConnectedResponse + 15, // 66: master_pb.Seaweed.LookupVolume:output_type -> master_pb.LookupVolumeResponse + 19, // 67: master_pb.Seaweed.Assign:output_type -> master_pb.AssignResponse + 19, // 68: master_pb.Seaweed.StreamAssign:output_type -> master_pb.AssignResponse + 21, // 69: master_pb.Seaweed.Statistics:output_type -> master_pb.StatisticsResponse + 24, // 70: master_pb.Seaweed.CollectionList:output_type -> master_pb.CollectionListResponse + 26, // 71: master_pb.Seaweed.CollectionDelete:output_type -> master_pb.CollectionDeleteResponse + 33, // 72: master_pb.Seaweed.VolumeList:output_type -> master_pb.VolumeListResponse + 35, // 73: master_pb.Seaweed.LookupEcVolume:output_type -> master_pb.LookupEcVolumeResponse + 37, // 74: master_pb.Seaweed.VacuumVolume:output_type -> master_pb.VacuumVolumeResponse + 39, // 75: master_pb.Seaweed.DisableVacuum:output_type -> master_pb.DisableVacuumResponse + 41, // 76: master_pb.Seaweed.EnableVacuum:output_type -> master_pb.EnableVacuumResponse + 43, // 77: master_pb.Seaweed.VolumeMarkReadonly:output_type -> master_pb.VolumeMarkReadonlyResponse + 45, // 78: master_pb.Seaweed.GetMasterConfiguration:output_type -> master_pb.GetMasterConfigurationResponse + 47, // 79: master_pb.Seaweed.ListClusterNodes:output_type -> master_pb.ListClusterNodesResponse + 49, // 80: master_pb.Seaweed.LeaseAdminToken:output_type -> master_pb.LeaseAdminTokenResponse + 51, // 81: master_pb.Seaweed.ReleaseAdminToken:output_type -> master_pb.ReleaseAdminTokenResponse + 53, // 82: master_pb.Seaweed.Ping:output_type -> master_pb.PingResponse + 59, // 83: master_pb.Seaweed.RaftListClusterServers:output_type -> master_pb.RaftListClusterServersResponse + 55, // 84: master_pb.Seaweed.RaftAddServer:output_type -> master_pb.RaftAddServerResponse + 57, // 85: master_pb.Seaweed.RaftRemoveServer:output_type -> master_pb.RaftRemoveServerResponse + 61, // 86: master_pb.Seaweed.RaftLeadershipTransfer:output_type -> master_pb.RaftLeadershipTransferResponse + 62, // 87: master_pb.Seaweed.VolumeGrow:output_type -> master_pb.VolumeGrowResponse + 64, // [64:88] is the sub-list for method output_type + 40, // [40:64] is the sub-list for method input_type + 40, // [40:40] is the sub-list for extension type_name + 40, // [40:40] is the sub-list for extension extendee + 0, // [0:40] is the sub-list for field type_name } func init() { file_master_proto_init() } @@ -4984,7 +5023,7 @@ func file_master_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_master_proto_rawDesc), len(file_master_proto_rawDesc)), NumEnums: 0, - NumMessages: 74, + NumMessages: 75, NumExtensions: 0, NumServices: 1, }, diff --git a/weed/pb/master_pb/master_helper.go b/weed/pb/master_pb/master_helper.go index 75de27c64..0d9e92872 100644 --- a/weed/pb/master_pb/master_helper.go +++ b/weed/pb/master_pb/master_helper.go @@ -6,19 +6,12 @@ func (v *VolumeLocation) IsEmptyUrl() bool { return v.Url == "" || v.Url == ":0" } -// SplitByPhysicalDisk returns one DiskInfo per physical disk_id observed in -// VolumeInfos / EcShardInfos. The wire format keys DataNodeInfo.DiskInfos by -// disk type, so multiple same-type physical disks on one DataNode collapse -// into a single DiskInfo entry. Per-volume and per-shard records carry the -// real physical DiskId; this helper rebuilds a per-physical-disk view from -// those records so consumers (topology indexes, shell output) can target -// individual disks instead of treating each node as one big disk. -// -// ActiveVolumeCount and RemoteVolumeCount are computed exactly from each -// disk's VolumeInfos (read-only and remote-backed are known per-volume). -// MaxVolumeCount and FreeVolumeCount are not derivable from per-volume -// records, so they are split across reconstructed disks with the remainder -// distributed to the lowest disk ids — the sums are preserved exactly. +// SplitByPhysicalDisk returns one DiskInfo per physical disk. The wire format +// keys DataNodeInfo.DiskInfos by disk type, collapsing same-type disks into one +// entry, so this rebuilds the per-disk view. MaxVolumeCountByDisk, when present, +// is the authoritative set (empty disks included) and gives each disk its exact +// max; otherwise the set is the DiskIds seen on records with aggregate Max/Free +// split evenly. func (d *DiskInfo) SplitByPhysicalDisk() []*DiskInfo { if d == nil { return nil @@ -53,25 +46,6 @@ func (d *DiskInfo) SplitByPhysicalDisk() []*DiskInfo { return id } - diskIDs := make(map[uint32]struct{}) - for _, vi := range d.VolumeInfos { - diskIDs[normalize(vi.DiskId)] = struct{}{} - } - for _, eci := range d.EcShardInfos { - diskIDs[normalize(eci.DiskId)] = struct{}{} - } - if len(diskIDs) == 0 { - diskIDs[d.DiskId] = struct{}{} - } - - if len(diskIDs) == 1 { - for diskID := range diskIDs { - if diskID == d.DiskId { - return []*DiskInfo{d} - } - } - } - perDiskVolumes := make(map[uint32][]*VolumeInformationMessage) for _, vi := range d.VolumeInfos { id := normalize(vi.DiskId) @@ -83,6 +57,35 @@ func (d *DiskInfo) SplitByPhysicalDisk() []*DiskInfo { perDiskShards[id] = append(perDiskShards[id], eci) } + diskIDs := make(map[uint32]struct{}) + for id := range perDiskVolumes { + diskIDs[id] = struct{}{} + } + for id := range perDiskShards { + diskIDs[id] = struct{}{} + } + + // MaxVolumeCountByDisk (when present) is the authoritative set with each + // disk's exact max. + exactMax := len(d.MaxVolumeCountByDisk) > 0 + for diskID := range d.MaxVolumeCountByDisk { + diskIDs[diskID] = struct{}{} + } + + if len(diskIDs) == 0 { + diskIDs[d.DiskId] = struct{}{} + } + + // A lone disk equal to the aggregate needs no reconstruction; exactMax disks + // always reconstruct so they report their own max. + if !exactMax && len(diskIDs) == 1 { + for diskID := range diskIDs { + if diskID == d.DiskId { + return []*DiskInfo{d} + } + } + } + // Sort disk IDs so the remainder distribution is deterministic and the // reconstructed slice is in DiskId order, which is what downstream // renderers expect. @@ -116,11 +119,24 @@ func (d *DiskInfo) SplitByPhysicalDisk() []*DiskInfo { remoteCount++ } } + volumeCount := int64(len(perDiskVolumes[diskID])) + maxVolumeCount := share(d.MaxVolumeCount, i) + freeVolumeCount := share(d.FreeVolumeCount, i) + if exactMax { + maxVolumeCount = d.MaxVolumeCountByDisk[diskID] + // EC slots aren't subtracted (needs the configurable ratio, outside + // this package); planners subtract them from EcShardInfos. Clamp so an + // over-allocated disk can't report negative free. + freeVolumeCount = maxVolumeCount - (volumeCount - remoteCount) + if freeVolumeCount < 0 { + freeVolumeCount = 0 + } + } result = append(result, &DiskInfo{ Type: d.Type, - MaxVolumeCount: share(d.MaxVolumeCount, i), - VolumeCount: int64(len(perDiskVolumes[diskID])), - FreeVolumeCount: share(d.FreeVolumeCount, i), + MaxVolumeCount: maxVolumeCount, + VolumeCount: volumeCount, + FreeVolumeCount: freeVolumeCount, ActiveVolumeCount: activeCount, RemoteVolumeCount: remoteCount, VolumeInfos: perDiskVolumes[diskID], diff --git a/weed/pb/master_pb/master_helper_test.go b/weed/pb/master_pb/master_helper_test.go index ff202c50c..541143d54 100644 --- a/weed/pb/master_pb/master_helper_test.go +++ b/weed/pb/master_pb/master_helper_test.go @@ -5,6 +5,75 @@ import ( "testing" ) +// MaxVolumeCountByDisk surfaces empty disks and gives each its exact max. +func TestDiskInfoSplitByPhysicalDisk_includesEmptyDiskFromPhysicalDisks(t *testing.T) { + d := &DiskInfo{ + Type: "hdd", + MaxVolumeCount: 1000, // per-type aggregate + VolumeInfos: []*VolumeInformationMessage{ + {Id: 10, DiskId: 0}, + {Id: 11, DiskId: 1}, + }, + MaxVolumeCountByDisk: map[uint32]int64{ + 0: 350, + 1: 350, + 2: 300, // empty disk, no volumes/shards + }, + } + + got := d.SplitByPhysicalDisk() + if len(got) != 3 { + t.Fatalf("want 3 physical disks including the empty one, got %d", len(got)) + } + byID := map[uint32]*DiskInfo{} + for _, di := range got { + byID[di.DiskId] = di + } + + empty, ok := byID[2] + if !ok { + t.Fatalf("empty disk 2 not surfaced; got %v", byID) + } + if empty.VolumeCount != 0 { + t.Errorf("empty disk: want 0 volumes, got %d", empty.VolumeCount) + } + if empty.FreeVolumeCount != 300 { + t.Errorf("empty disk free: want its full max 300, got %d", empty.FreeVolumeCount) + } + + // Exact per-disk max, not the even split (which would give ~334/333/333). + if byID[0].MaxVolumeCount != 350 || byID[1].MaxVolumeCount != 350 || byID[2].MaxVolumeCount != 300 { + t.Errorf("want exact per-disk max 350/350/300, got %d/%d/%d", + byID[0].MaxVolumeCount, byID[1].MaxVolumeCount, byID[2].MaxVolumeCount) + } + if byID[0].FreeVolumeCount != 349 { + t.Errorf("disk 0 free: want 349 (350-1 volume), got %d", byID[0].FreeVolumeCount) + } +} + +// An over-allocated disk reports zero free, not negative. +func TestDiskInfoSplitByPhysicalDisk_clampsNegativeFreeOnOverAllocation(t *testing.T) { + d := &DiskInfo{ + Type: "hdd", + VolumeInfos: []*VolumeInformationMessage{ + {Id: 1, DiskId: 0}, + {Id: 2, DiskId: 0}, + {Id: 3, DiskId: 0}, + }, + MaxVolumeCountByDisk: map[uint32]int64{ + 0: 2, // 3 volumes on a max-2 disk + }, + } + + got := d.SplitByPhysicalDisk() + if len(got) != 1 { + t.Fatalf("want 1 disk, got %d", len(got)) + } + if got[0].FreeVolumeCount != 0 { + t.Errorf("over-allocated disk free: want clamped 0, got %d", got[0].FreeVolumeCount) + } +} + func TestDiskInfoSplitByPhysicalDisk_collapsesOnSingleDisk(t *testing.T) { d := &DiskInfo{ Type: "hdd", @@ -138,9 +207,9 @@ func TestDiskInfoSplitByPhysicalDisk_countsActiveAndRemoteExactly(t *testing.T) } cases := []struct { - id uint32 - wantActive int64 - wantRemote int64 + id uint32 + wantActive int64 + wantRemote int64 }{ {0, 1, 0}, {1, 1, 1}, diff --git a/weed/storage/store.go b/weed/storage/store.go index 70ca2962a..3b742fc2b 100644 --- a/weed/storage/store.go +++ b/weed/storage/store.go @@ -408,11 +408,13 @@ func (s *Store) GetRack() string { func (s *Store) CollectHeartbeat() *master_pb.Heartbeat { var volumeMessages []*master_pb.VolumeInformationMessage maxVolumeCounts := make(map[string]uint32) + // Per-disk effective max for DiskTag, captured alongside the per-type sum. + diskMaxByID := make(map[int]int32) var maxFileKey NeedleId collectionVolumeSize := make(map[string]int64) collectionVolumeDeletedBytes := make(map[string]int64) collectionVolumeReadOnlyCount := make(map[string]map[string]uint8) - for _, location := range s.Locations { + for diskID, location := range s.Locations { if location.isDiskUnavailable.Load() { continue } @@ -428,6 +430,7 @@ func (s *Store) CollectHeartbeat() *master_pb.Heartbeat { effectiveMaxCount = 0 } maxVolumeCounts[string(location.DiskType)] += uint32(effectiveMaxCount) + diskMaxByID[diskID] = effectiveMaxCount location.volumesLock.RLock() for _, v := range location.volumes { curMaxFileKey, volumeMessage := v.ToVolumeInformationMessage() @@ -540,8 +543,10 @@ func (s *Store) CollectHeartbeat() *master_pb.Heartbeat { var diskTags []*master_pb.DiskTag for diskID, loc := range s.Locations { diskTags = append(diskTags, &master_pb.DiskTag{ - DiskId: uint32(diskID), - Tags: append([]string(nil), loc.Tags...), + DiskId: uint32(diskID), + Tags: append([]string(nil), loc.Tags...), + Type: string(loc.DiskType), + MaxVolumeCount: int64(diskMaxByID[diskID]), }) } diff --git a/weed/topology/data_node.go b/weed/topology/data_node.go index 94c5e7c7f..b91a75114 100644 --- a/weed/topology/data_node.go +++ b/weed/topology/data_node.go @@ -25,7 +25,15 @@ type DataNode struct { IsTerminating bool MaintenanceMode bool - diskTags map[uint32][]string + // diskMetas holds each physical disk's tags, type, and capacity from the + // heartbeat DiskTags, including disks with no volumes or EC shards. + diskMetas map[uint32]diskMeta +} + +type diskMeta struct { + tags []string + diskType types.DiskType + maxVolumeCount int64 } func NewDataNode(id string) *DataNode { @@ -294,17 +302,35 @@ func (dn *DataNode) ToDataNodeInfo() *master_pb.DataNodeInfo { } dn.RLock() - diskTags := make(map[uint32][]string, len(dn.diskTags)) - for diskID, tags := range dn.diskTags { - diskTags[diskID] = append([]string(nil), tags...) + metas := make(map[uint32]diskMeta, len(dn.diskMetas)) + for diskID, meta := range dn.diskMetas { + metas[diskID] = meta } dn.RUnlock() for _, diskInfo := range m.DiskInfos { if diskInfo == nil { continue } - if tags, found := diskTags[diskInfo.DiskId]; found { - diskInfo.Tags = append([]string(nil), tags...) + if meta, found := metas[diskInfo.DiskId]; found { + diskInfo.Tags = append([]string(nil), meta.tags...) + } + // Max per physical disk of this type, empty and unavailable (max 0) ones + // included. Emit only when some disk reports capacity, so an older server + // sending all zeros leaves the map nil and falls back. + diskType := types.ToDiskType(diskInfo.Type) + maxByDisk := make(map[uint32]int64) + anyCapacity := false + for diskID, meta := range metas { + if meta.diskType != diskType { + continue + } + if meta.maxVolumeCount > 0 { + anyCapacity = true + } + maxByDisk[diskID] = meta.maxVolumeCount + } + if anyCapacity { + diskInfo.MaxVolumeCountByDisk = maxByDisk } } return m @@ -314,16 +340,21 @@ func (dn *DataNode) UpdateDiskTags(tags []*master_pb.DiskTag) { if len(tags) == 0 { return } - dn.Lock() - if dn.diskTags == nil { - dn.diskTags = make(map[uint32][]string, len(tags)) - } + // DiskTags is the full list on each full heartbeat; rebuild fresh to drop + // removed disks. + metas := make(map[uint32]diskMeta, len(tags)) for _, tagInfo := range tags { if tagInfo == nil { continue } - dn.diskTags[tagInfo.DiskId] = append([]string(nil), tagInfo.Tags...) + metas[tagInfo.DiskId] = diskMeta{ + tags: append([]string(nil), tagInfo.Tags...), + diskType: types.ToDiskType(tagInfo.Type), + maxVolumeCount: tagInfo.MaxVolumeCount, + } } + dn.Lock() + dn.diskMetas = metas dn.Unlock() } diff --git a/weed/topology/data_node_empty_disk_test.go b/weed/topology/data_node_empty_disk_test.go new file mode 100644 index 000000000..74c7e15e3 --- /dev/null +++ b/weed/topology/data_node_empty_disk_test.go @@ -0,0 +1,88 @@ +package topology + +import ( + "testing" + + "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" + "github.com/seaweedfs/seaweedfs/weed/sequence" +) + +// A disk with no volumes/shards surfaces via DiskTags. +func TestToDataNodeInfoReportsEmptyPhysicalDisks(t *testing.T) { + topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false) + dc := topo.GetOrCreateDataCenter("dc1") + rack := dc.GetOrCreateRack("rack1") + dn := rack.GetOrCreateDataNode("127.0.0.1", 34534, 0, "127.0.0.1", "", map[string]uint32{"": 1000}) + dn.AdjustMaxVolumeCounts(map[string]uint32{"": 1000}) + + dn.UpdateDiskTags([]*master_pb.DiskTag{ + {DiskId: 0, Type: "", MaxVolumeCount: 350}, + {DiskId: 1, Type: "", MaxVolumeCount: 350}, + {DiskId: 2, Type: "", MaxVolumeCount: 300}, // empty disk, never held a volume + }) + + info := dn.ToDataNodeInfo() + di, ok := info.DiskInfos[""] + if !ok { + t.Fatalf("missing HDD disk info") + } + if len(di.MaxVolumeCountByDisk) != 3 { + t.Fatalf("want 3 physical disks reported, got %d", len(di.MaxVolumeCountByDisk)) + } + + split := di.SplitByPhysicalDisk() + if len(split) != 3 { + t.Fatalf("SplitByPhysicalDisk: want 3 disks including the empty one, got %d", len(split)) + } + byID := map[uint32]*master_pb.DiskInfo{} + for _, d := range split { + byID[d.DiskId] = d + } + if byID[2] == nil { + t.Fatalf("empty disk 2 not surfaced") + } + if byID[2].MaxVolumeCount != 300 || byID[2].VolumeCount != 0 { + t.Errorf("empty disk 2: want max 300 / 0 volumes, got max %d / %d volumes", + byID[2].MaxVolumeCount, byID[2].VolumeCount) + } +} + +// A max-0 (unavailable) disk stays listed when the node reports capacity. +func TestToDataNodeInfoKeepsZeroCapacityDisk(t *testing.T) { + topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false) + dc := topo.GetOrCreateDataCenter("dc1") + rack := dc.GetOrCreateRack("rack1") + dn := rack.GetOrCreateDataNode("127.0.0.1", 34534, 0, "127.0.0.1", "", map[string]uint32{"": 700}) + dn.AdjustMaxVolumeCounts(map[string]uint32{"": 700}) + + dn.UpdateDiskTags([]*master_pb.DiskTag{ + {DiskId: 0, Type: "", MaxVolumeCount: 350}, + {DiskId: 1, Type: "", MaxVolumeCount: 350}, + {DiskId: 2, Type: "", MaxVolumeCount: 0}, // unavailable disk + }) + + di := dn.ToDataNodeInfo().DiskInfos[""] + if len(di.MaxVolumeCountByDisk) != 3 { + t.Fatalf("want 3 physical disks (incl the zero-capacity one), got %d", len(di.MaxVolumeCountByDisk)) + } +} + +// An older server sending no per-disk capacity leaves the map empty. +func TestToDataNodeInfoFallsBackWhenNoCapacityReported(t *testing.T) { + topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false) + dc := topo.GetOrCreateDataCenter("dc1") + rack := dc.GetOrCreateRack("rack1") + dn := rack.GetOrCreateDataNode("127.0.0.1", 34534, 0, "127.0.0.1", "", map[string]uint32{"": 700}) + dn.AdjustMaxVolumeCounts(map[string]uint32{"": 700}) + + // Older server: DiskTags carry disk_id (+tags) but no type/max. + dn.UpdateDiskTags([]*master_pb.DiskTag{ + {DiskId: 0}, + {DiskId: 1}, + }) + + di := dn.ToDataNodeInfo().DiskInfos[""] + if len(di.MaxVolumeCountByDisk) != 0 { + t.Fatalf("want no per-disk max (fallback), got %d", len(di.MaxVolumeCountByDisk)) + } +}