diff --git a/weed/admin/topology/active_topology_test.go b/weed/admin/topology/active_topology_test.go index a8b3f8085..79aef9465 100644 --- a/weed/admin/topology/active_topology_test.go +++ b/weed/admin/topology/active_topology_test.go @@ -408,6 +408,57 @@ func TestECPlanningNotBlockedByUnrelatedBalance(t *testing.T) { "EC must still see all 4 disks even with an unrelated in-flight balance") } +// #9369: same-type physical disks collapse into one DiskInfo at the master; +// the active topology must still expose one entry per physical disk_id. +func TestECPlannerSeesEachPhysicalDisk(t *testing.T) { + topology := NewActiveTopology(10) + + const numServers = 7 + const disksPerServer = 2 + nodes := make([]*master_pb.DataNodeInfo, 0, numServers) + for i := 1; i <= numServers; i++ { + volumeInfos := make([]*master_pb.VolumeInformationMessage, 0, disksPerServer) + for d := uint32(0); d < disksPerServer; d++ { + volumeInfos = append(volumeInfos, &master_pb.VolumeInformationMessage{ + Id: uint32(i*10 + int(d)), + DiskId: d, + DiskType: "hdd", + }) + } + nodes = append(nodes, &master_pb.DataNodeInfo{ + Id: fmt.Sprintf("127.0.0.1:%d", 8080+i), + DiskInfos: map[string]*master_pb.DiskInfo{ + "hdd": { + DiskId: 0, + VolumeCount: int64(disksPerServer), + MaxVolumeCount: 200, + VolumeInfos: volumeInfos, + }, + }, + }) + } + + require.NoError(t, topology.UpdateTopology(&master_pb.TopologyInfo{ + DataCenterInfos: []*master_pb.DataCenterInfo{{ + Id: "dc1", + RackInfos: []*master_pb.RackInfo{{ + Id: "rack1", + DataNodeInfos: nodes, + }}, + }}, + })) + + candidates := topology.GetDisksWithEffectiveCapacity(TaskTypeErasureCoding, "", 0) + assert.Equal(t, numServers*disksPerServer, len(candidates)) + + seen := make(map[string]bool, len(candidates)) + for _, c := range candidates { + key := fmt.Sprintf("%s:%d", c.NodeID, c.DiskID) + assert.False(t, seen[key], "duplicate placement target %s", key) + seen[key] = true + } +} + // TestPublicInterfaces tests the public interface methods func TestPublicInterfaces(t *testing.T) { topology := NewActiveTopology(10) diff --git a/weed/admin/topology/topology_management.go b/weed/admin/topology/topology_management.go index 8073ae3da..6b47695ab 100644 --- a/weed/admin/topology/topology_management.go +++ b/weed/admin/topology/topology_management.go @@ -8,6 +8,74 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" ) +// splitDiskInfoByPhysicalDisk returns one master_pb.DiskInfo per physical +// disk_id observed in VolumeInfos / EcShardInfos. Multiple same-type physical +// disks collapse to one DiskInfo at the master; per-volume/per-shard records +// keep the original disk_id and are the authoritative signal here. Capacity +// is split evenly — the wire format doesn't carry per-disk capacity yet. +func splitDiskInfoByPhysicalDisk(diskInfo *master_pb.DiskInfo) []*master_pb.DiskInfo { + if diskInfo == nil { + return nil + } + + // Records with DiskId=0 and a non-zero outer DiskId belong to the outer + // disk — handles older payloads / fixtures that omit the per-record id. + normalize := func(id uint32) uint32 { + if id == 0 && diskInfo.DiskId != 0 { + return diskInfo.DiskId + } + return id + } + + diskIDs := make(map[uint32]struct{}) + for _, vi := range diskInfo.VolumeInfos { + diskIDs[normalize(vi.DiskId)] = struct{}{} + } + for _, eci := range diskInfo.EcShardInfos { + diskIDs[normalize(eci.DiskId)] = struct{}{} + } + if len(diskIDs) == 0 { + diskIDs[diskInfo.DiskId] = struct{}{} + } + + if len(diskIDs) == 1 { + for diskID := range diskIDs { + if diskID == diskInfo.DiskId { + return []*master_pb.DiskInfo{diskInfo} + } + } + } + + perDiskVolumes := make(map[uint32][]*master_pb.VolumeInformationMessage) + for _, vi := range diskInfo.VolumeInfos { + perDiskVolumes[normalize(vi.DiskId)] = append(perDiskVolumes[normalize(vi.DiskId)], vi) + } + perDiskShards := make(map[uint32][]*master_pb.VolumeEcShardInformationMessage) + for _, eci := range diskInfo.EcShardInfos { + perDiskShards[normalize(eci.DiskId)] = append(perDiskShards[normalize(eci.DiskId)], eci) + } + + count := int64(len(diskIDs)) + share := func(total int64) int64 { return total / count } + + result := make([]*master_pb.DiskInfo, 0, len(diskIDs)) + for diskID := range diskIDs { + result = append(result, &master_pb.DiskInfo{ + Type: diskInfo.Type, + MaxVolumeCount: share(diskInfo.MaxVolumeCount), + VolumeCount: int64(len(perDiskVolumes[diskID])), + FreeVolumeCount: share(diskInfo.FreeVolumeCount), + ActiveVolumeCount: share(diskInfo.ActiveVolumeCount), + RemoteVolumeCount: share(diskInfo.RemoteVolumeCount), + VolumeInfos: perDiskVolumes[diskID], + EcShardInfos: perDiskShards[diskID], + DiskId: diskID, + Tags: append([]string(nil), diskInfo.Tags...), + }) + } + return result +} + // CountTopologyResources counts datacenters, nodes, and disks in topology info func CountTopologyResources(topologyInfo *master_pb.TopologyInfo) (dcCount, nodeCount, diskCount int) { if topologyInfo == nil { @@ -73,24 +141,28 @@ func (at *ActiveTopology) UpdateTopology(topologyInfo *master_pb.TopologyInfo) e disks: make(map[uint32]*activeDisk), } - // Add disks for this node + // One activeDisk per physical disk_id (#9369): the master keys + // DiskInfos by disk type, so same-type disks must be split out. for diskType, diskInfo := range nodeInfo.DiskInfos { - disk := &activeDisk{ - DiskInfo: &DiskInfo{ - NodeID: nodeInfo.Id, - DiskID: diskInfo.DiskId, - DiskType: diskType, - DataCenter: dc.Id, - Rack: rack.Id, - DiskInfo: diskInfo, - }, - } + perDiskInfos := splitDiskInfoByPhysicalDisk(diskInfo) + for _, perDisk := range perDiskInfos { + disk := &activeDisk{ + DiskInfo: &DiskInfo{ + NodeID: nodeInfo.Id, + DiskID: perDisk.DiskId, + DiskType: diskType, + DataCenter: dc.Id, + Rack: rack.Id, + DiskInfo: perDisk, + }, + } - diskKey := fmt.Sprintf("%s:%d", nodeInfo.Id, diskInfo.DiskId) - glog.V(3).Infof("UpdateTopology: adding disk key=%q nodeId=%q diskId=%d diskType=%q address=%q grpcPort=%d volumes=%d maxVolumes=%d", - diskKey, nodeInfo.Id, diskInfo.DiskId, diskType, nodeInfo.Address, nodeInfo.GrpcPort, diskInfo.VolumeCount, diskInfo.MaxVolumeCount) - node.disks[diskInfo.DiskId] = disk - at.disks[diskKey] = disk + diskKey := fmt.Sprintf("%s:%d", nodeInfo.Id, perDisk.DiskId) + glog.V(3).Infof("UpdateTopology: adding disk key=%q nodeId=%q diskId=%d diskType=%q address=%q grpcPort=%d volumes=%d maxVolumes=%d", + diskKey, nodeInfo.Id, perDisk.DiskId, diskType, nodeInfo.Address, nodeInfo.GrpcPort, perDisk.VolumeCount, perDisk.MaxVolumeCount) + node.disks[perDisk.DiskId] = disk + at.disks[diskKey] = disk + } } at.nodes[nodeInfo.Id] = node @@ -205,23 +277,27 @@ func (at *ActiveTopology) rebuildIndexes() { at.volumeIndex = make(map[uint32][]string) at.ecShardIndex = make(map[uint32][]string) - // Rebuild indexes from current topology + // Index by the per-record DiskId (not the outer DiskInfo.DiskId) so the + // keys match the per-physical-disk activeDisk entries — see #9369. for _, dc := range at.topologyInfo.DataCenterInfos { for _, rack := range dc.RackInfos { for _, nodeInfo := range rack.DataNodeInfos { for _, diskInfo := range nodeInfo.DiskInfos { - diskKey := fmt.Sprintf("%s:%d", nodeInfo.Id, diskInfo.DiskId) - - // Index volumes for _, volumeInfo := range diskInfo.VolumeInfos { - volumeID := volumeInfo.Id - at.volumeIndex[volumeID] = append(at.volumeIndex[volumeID], diskKey) + diskID := volumeInfo.DiskId + if diskID == 0 && diskInfo.DiskId != 0 { + diskID = diskInfo.DiskId + } + diskKey := fmt.Sprintf("%s:%d", nodeInfo.Id, diskID) + at.volumeIndex[volumeInfo.Id] = append(at.volumeIndex[volumeInfo.Id], diskKey) } - - // Index EC shards for _, ecShardInfo := range diskInfo.EcShardInfos { - volumeID := ecShardInfo.Id - at.ecShardIndex[volumeID] = append(at.ecShardIndex[volumeID], diskKey) + diskID := ecShardInfo.DiskId + if diskID == 0 && diskInfo.DiskId != 0 { + diskID = diskInfo.DiskId + } + diskKey := fmt.Sprintf("%s:%d", nodeInfo.Id, diskID) + at.ecShardIndex[ecShardInfo.Id] = append(at.ecShardIndex[ecShardInfo.Id], diskKey) } } } diff --git a/weed/worker/tasks/erasure_coding/detection.go b/weed/worker/tasks/erasure_coding/detection.go index 453bb8fb6..6d63738f4 100644 --- a/weed/worker/tasks/erasure_coding/detection.go +++ b/weed/worker/tasks/erasure_coding/detection.go @@ -151,7 +151,9 @@ func Detection(ctx context.Context, metrics []*types.VolumeHealthMetrics, cluste if planner == nil { planner = newECPlacementPlanner(clusterInfo.ActiveTopology, ecConfig.PreferredTags) } - multiPlan, err := planECDestinations(planner, metric, ecConfig) + dataShards := erasure_coding.DataShardsCount + parityShards := erasure_coding.ParityShardsCount + multiPlan, err := planECDestinations(planner, metric, ecConfig, dataShards, parityShards) if err != nil { glog.Warningf("Failed to plan EC destinations for volume %d: %v", metric.VolumeID, err) consecutivePlanningFailures++ @@ -168,7 +170,7 @@ func Detection(ctx context.Context, metrics []*types.VolumeHealthMetrics, cluste // Calculate expected shard size for EC operation // Each data shard will be approximately volumeSize / dataShards - expectedShardSize := uint64(metric.Size) / uint64(erasure_coding.DataShardsCount) + expectedShardSize := uint64(metric.Size) / uint64(dataShards) // Add pending EC shard task to ActiveTopology for capacity management @@ -281,10 +283,10 @@ func Detection(ctx context.Context, metrics []*types.VolumeHealthMetrics, cluste Sources: sourcesProto, // Unified targets - all EC shard destinations - Targets: createECTargets(multiPlan), + Targets: createECTargets(multiPlan, dataShards, parityShards), TaskParams: &worker_pb.TaskParams_ErasureCodingParams{ - ErasureCodingParams: createECTaskParams(multiPlan), + ErasureCodingParams: createECTaskParams(dataShards, parityShards), }, } @@ -591,14 +593,20 @@ func (p *ecPlacementPlanner) buildCandidateSets(shardsNeeded int) [][]*placement return candidateSets } -// planECDestinations plans the destinations for erasure coding operation -// This function implements EC destination planning logic directly in the detection phase -func planECDestinations(planner *ecPlacementPlanner, metric *types.VolumeHealthMetrics, ecConfig *Config) (*topology.MultiDestinationPlan, error) { +// planECDestinations plans the destinations for erasure coding operation. +// dataShards/parityShards are parameters so callers can drive non-10+4 ratios. +func planECDestinations(planner *ecPlacementPlanner, metric *types.VolumeHealthMetrics, ecConfig *Config, dataShards, parityShards int) (*topology.MultiDestinationPlan, error) { if planner == nil || planner.activeTopology == nil { return nil, fmt.Errorf("active topology not available for EC placement") } - // Calculate expected shard size for EC operation - expectedShardSize := uint64(metric.Size) / uint64(erasure_coding.DataShardsCount) + if dataShards <= 0 || parityShards <= 0 { + return nil, fmt.Errorf("invalid EC ratio: dataShards=%d parityShards=%d", dataShards, parityShards) + } + totalShards := dataShards + parityShards + // Survive losing one disk: each disk holds at most parityShards shards, + // so we need at least ceil(totalShards / parityShards) disks. + minTotalDisks := (totalShards + parityShards - 1) / parityShards + expectedShardSize := uint64(metric.Size) / uint64(dataShards) // Get source node information from topology var sourceRack, sourceDC string @@ -626,12 +634,18 @@ func planECDestinations(planner *ecPlacementPlanner, metric *types.VolumeHealthM } // Select best disks for EC placement with rack/DC diversity using the cached planner - selectedDisks, err := planner.selectDestinations(sourceRack, sourceDC, erasure_coding.TotalShardsCount) + selectedDisks, err := planner.selectDestinations(sourceRack, sourceDC, totalShards) if err != nil { return nil, err } - if len(selectedDisks) < erasure_coding.MinTotalDisks { - return nil, fmt.Errorf("found %d disks, but could not find %d suitable destinations for EC placement", len(selectedDisks), erasure_coding.MinTotalDisks) + if len(selectedDisks) < minTotalDisks { + return nil, fmt.Errorf("found %d disks, but could not find %d suitable destinations for EC placement", len(selectedDisks), minTotalDisks) + } + // One shard per (server, disk_id): #9185's disk_id-aware ReceiveFile rejects + // a second shard on the same disk. + if len(selectedDisks) < totalShards { + return nil, fmt.Errorf("found %d disks, but EC %d+%d needs %d distinct (server, disk_id) targets", + len(selectedDisks), dataShards, parityShards, totalShards) } var plans []*topology.DestinationPlan @@ -688,54 +702,48 @@ func planECDestinations(planner *ecPlacementPlanner, metric *types.VolumeHealthM }, nil } -// createECTargets creates unified TaskTarget structures from the multi-destination plan -// with proper shard ID assignment during planning phase -func createECTargets(multiPlan *topology.MultiDestinationPlan) []*worker_pb.TaskTarget { +// createECTargets builds TaskTargets with one shard per plan entry. +// planECDestinations ensures numTargets == totalShards. +func createECTargets(multiPlan *topology.MultiDestinationPlan, dataShards, parityShards int) []*worker_pb.TaskTarget { var targets []*worker_pb.TaskTarget numTargets := len(multiPlan.Plans) + totalShards := dataShards + parityShards - // Create shard assignment arrays for each target (round-robin distribution) targetShards := make([][]uint32, numTargets) for i := range targetShards { targetShards[i] = make([]uint32, 0) } - - // Distribute shards in round-robin fashion to spread both data and parity shards - // This ensures each target gets a mix of data shards (0-9) and parity shards (10-13) - for shardId := uint32(0); shardId < uint32(erasure_coding.TotalShardsCount); shardId++ { - targetIndex := int(shardId) % numTargets - targetShards[targetIndex] = append(targetShards[targetIndex], shardId) + for shardId := 0; shardId < totalShards; shardId++ { + targetIndex := shardId % numTargets + targetShards[targetIndex] = append(targetShards[targetIndex], uint32(shardId)) } - // Create targets with assigned shard IDs for i, plan := range multiPlan.Plans { target := &worker_pb.TaskTarget{ Node: plan.TargetAddress, DiskId: plan.TargetDisk, Rack: plan.TargetRack, DataCenter: plan.TargetDC, - ShardIds: targetShards[i], // Round-robin assigned shards + ShardIds: targetShards[i], EstimatedSize: plan.ExpectedSize, } targets = append(targets, target) - // Log shard assignment with data/parity classification - dataShards := make([]uint32, 0) - parityShards := make([]uint32, 0) + assignedData := make([]uint32, 0) + assignedParity := make([]uint32, 0) for _, shardId := range targetShards[i] { - if shardId < uint32(erasure_coding.DataShardsCount) { - dataShards = append(dataShards, shardId) + if int(shardId) < dataShards { + assignedData = append(assignedData, shardId) } else { - parityShards = append(parityShards, shardId) + assignedParity = append(assignedParity, shardId) } } glog.V(2).Infof("EC planning: target %s assigned shards %v (data: %v, parity: %v)", - plan.TargetNode, targetShards[i], dataShards, parityShards) + plan.TargetNode, targetShards[i], assignedData, assignedParity) } glog.V(1).Infof("EC planning: distributed %d shards across %d targets using round-robin (data shards 0-%d, parity shards %d-%d)", - erasure_coding.TotalShardsCount, numTargets, - erasure_coding.DataShardsCount-1, erasure_coding.DataShardsCount, erasure_coding.TotalShardsCount-1) + totalShards, numTargets, dataShards-1, dataShards, totalShards-1) return targets } @@ -779,10 +787,10 @@ func convertTaskSourcesToProtobuf(sources []topology.TaskSourceSpec, volumeID ui } // createECTaskParams creates clean EC task parameters (destinations now in unified targets) -func createECTaskParams(multiPlan *topology.MultiDestinationPlan) *worker_pb.ErasureCodingTaskParams { +func createECTaskParams(dataShards, parityShards int) *worker_pb.ErasureCodingTaskParams { return &worker_pb.ErasureCodingTaskParams{ - DataShards: erasure_coding.DataShardsCount, // Standard data shards - ParityShards: erasure_coding.ParityShardsCount, // Standard parity shards + DataShards: int32(dataShards), + ParityShards: int32(parityShards), } } diff --git a/weed/worker/tasks/erasure_coding/detection_test.go b/weed/worker/tasks/erasure_coding/detection_test.go index f47355a08..1f51b8801 100644 --- a/weed/worker/tasks/erasure_coding/detection_test.go +++ b/weed/worker/tasks/erasure_coding/detection_test.go @@ -57,7 +57,7 @@ func TestPlanECDestinationsUsesPlanner(t *testing.T) { Collection: "", } - plan, err := planECDestinations(planner, metric, NewDefaultConfig()) + plan, err := planECDestinations(planner, metric, NewDefaultConfig(), erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount) require.NoError(t, err) require.NotNil(t, plan) assert.Equal(t, erasure_coding.TotalShardsCount, len(plan.Plans)) @@ -140,7 +140,8 @@ func TestDetectionContextCancellation(t *testing.T) { } func TestDetectionMaxResultsHonorsLimit(t *testing.T) { - activeTopology := buildActiveTopology(t, 4, []string{"hdd"}, 20, 0) + // One node per shard so each shard gets its own disk (#9369). + activeTopology := buildActiveTopology(t, erasure_coding.TotalShardsCount, []string{"hdd"}, 20, 0) clusterInfo := &types.ClusterInfo{ActiveTopology: activeTopology} metrics := buildVolumeMetricsForIDs(3) @@ -150,6 +151,68 @@ func TestDetectionMaxResultsHonorsLimit(t *testing.T) { assert.True(t, hasMore) } +// #9369: 7 servers × 2 physical HDDs must yield 14 distinct (server, disk_id) +// destinations, not 7 destinations doubled up on the same disk. +func TestPlanECDestinationsSpreadsAcrossPhysicalDisks(t *testing.T) { + const numServers = 7 + const disksPerServer = 2 + + activeTopology := topology.NewActiveTopology(10) + nodes := make([]*master_pb.DataNodeInfo, 0, numServers) + for i := 1; i <= numServers; i++ { + volumeInfos := make([]*master_pb.VolumeInformationMessage, 0, disksPerServer) + for d := uint32(0); d < disksPerServer; d++ { + volumeInfos = append(volumeInfos, &master_pb.VolumeInformationMessage{ + Id: uint32(i*100 + int(d)), + DiskId: d, + DiskType: "hdd", + }) + } + nodes = append(nodes, &master_pb.DataNodeInfo{ + Id: fmt.Sprintf("127.0.0.1:%d", 8080+i), + DiskInfos: map[string]*master_pb.DiskInfo{ + "hdd": { + DiskId: 0, + VolumeCount: int64(disksPerServer), + MaxVolumeCount: 200, + VolumeInfos: volumeInfos, + }, + }, + }) + } + require.NoError(t, activeTopology.UpdateTopology(&master_pb.TopologyInfo{ + DataCenterInfos: []*master_pb.DataCenterInfo{{ + Id: "dc1", + RackInfos: []*master_pb.RackInfo{{ + Id: "rack1", + DataNodeInfos: nodes, + }}, + }}, + })) + + planner := newECPlacementPlanner(activeTopology, nil) + require.NotNil(t, planner) + + metric := &types.VolumeHealthMetrics{ + VolumeID: 42, + Server: "127.0.0.1:8081", + Size: 100 * 1024 * 1024, + Collection: "", + } + + plan, err := planECDestinations(planner, metric, NewDefaultConfig(), erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount) + require.NoError(t, err) + require.NotNil(t, plan) + require.Equal(t, erasure_coding.TotalShardsCount, len(plan.Plans)) + + seen := make(map[string]bool, len(plan.Plans)) + for _, p := range plan.Plans { + key := fmt.Sprintf("%s:%d", p.TargetNode, p.TargetDisk) + assert.False(t, seen[key], "duplicate (server,disk_id) target %s", key) + seen[key] = true + } +} + func TestPlanECDestinationsFailsWithInsufficientCapacity(t *testing.T) { activeTopology := buildActiveTopology(t, 1, []string{"hdd"}, 1, 1) planner := newECPlacementPlanner(activeTopology, nil) @@ -162,7 +225,7 @@ func TestPlanECDestinationsFailsWithInsufficientCapacity(t *testing.T) { Collection: "", } - _, err := planECDestinations(planner, metric, NewDefaultConfig()) + _, err := planECDestinations(planner, metric, NewDefaultConfig(), erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount) require.Error(t, err) }