diff --git a/weed/storage/erasure_coding/ecbalancer/place.go b/weed/storage/erasure_coding/ecbalancer/place.go index 84a43c632..d42d96f24 100644 --- a/weed/storage/erasure_coding/ecbalancer/place.go +++ b/weed/storage/erasure_coding/ecbalancer/place.go @@ -81,15 +81,17 @@ type PlaceResult struct { type PlacementMode int const ( - // PlaceStrict: caps and ReplicaPlacement are hard. Place fails rather than - // violate them, so the caller can defer (leave the volume as-is and retry). + // PlaceStrict: caps, the total-shards-per-rack cap and ReplicaPlacement are + // hard. Place fails rather than violate them, so the caller can defer (leave + // the volume as-is and retry). PlaceStrict PlacementMode = iota // PlaceDurabilityFirst (used by both encode and repair): relax per-type caps -> - // data/parity anti-affinity -> ReplicaPlacement, in that order, until each shard - // lands, reporting what was relaxed in PlaceResult.Relaxed. The per-disk - // durability cap (<= parityShards per disk) is never relaxed. Fails only if no - // disk has free capacity. Encode places best-effort this way and rebalancing - // tightens the spread afterward. + // data/parity anti-affinity -> ReplicaPlacement -> the total-shards-per-rack + // cap, in that order, until each shard lands, reporting what was relaxed in + // PlaceResult.Relaxed. The per-disk durability cap (<= parityShards per disk) + // is never relaxed. Fails only if no eligible disk has room for another shard + // of the volume. Encode places best-effort this way and rebalancing tightens + // the spread afterward. PlaceDurabilityFirst ) @@ -101,6 +103,7 @@ type relaxation struct { caps bool antiAffinity bool rp bool + rackTotal bool } func (r relaxation) relaxedNames() []string { @@ -114,16 +117,23 @@ func (r relaxation) relaxedNames() []string { if !r.rp { n = append(n, "replica-placement") } + if !r.rackTotal { + n = append(n, "rack-total-cap") + } return n } -var strictAttempts = []relaxation{{caps: true, antiAffinity: true, rp: true}} +var strictAttempts = []relaxation{{caps: true, antiAffinity: true, rp: true, rackTotal: true}} +// The total-shards-per-rack cap is sized from real capacity (see rackTotalCap), +// so it is relaxed last and only as a safety net; with it dropped, the per-disk +// cap is the only constraint left. var durabilityAttempts = []relaxation{ - {caps: true, antiAffinity: true, rp: true}, - {caps: false, antiAffinity: true, rp: true}, - {caps: false, antiAffinity: false, rp: true}, - {caps: false, antiAffinity: false, rp: false}, + {caps: true, antiAffinity: true, rp: true, rackTotal: true}, + {caps: false, antiAffinity: true, rp: true, rackTotal: true}, + {caps: false, antiAffinity: false, rp: true, rackTotal: true}, + {caps: false, antiAffinity: false, rp: false, rackTotal: true}, + {caps: false, antiAffinity: false, rp: false, rackTotal: false}, } type placedEntry struct { @@ -236,12 +246,17 @@ func (t *Topology) tryPlace(vk volKey, need []int, dataShards, parityShards int, } } - // Even per-rack caps divide by racks that actually have an eligible free disk, - // not all racks (the snapshot keeps every disk type/tag), so a valid tiered - // cluster — e.g. SSDs in only 2 of 4 racks — is not capped impossibly low. + // Even per-rack caps divide by racks that can actually take another shard of + // the volume, not all racks (the snapshot keeps every disk type/tag), so a + // valid tiered cluster — e.g. SSDs in only 2 of 4 racks — is not capped + // impossibly low. + rackRoom := make(map[string]int, len(rackKeys)) + rackRoomRP := make(map[string]int, len(rackKeys)) numEligibleRacks := 0 for _, rk := range rackKeys { - if rackHasFreeDisk(racks[rk], eligible) { + rackRoom[rk] = rackShardRoom(racks[rk], vk, eligible, parityShards) + rackRoomRP[rk] = rackShardRoomUnderRP(racks[rk], vk, eligible, parityShards, rp) + if rackRoom[rk] > 0 { numEligibleRacks++ } } @@ -249,6 +264,24 @@ func (t *Topology) tryPlace(vk volKey, need []int, dataShards, parityShards int, numEligibleRacks = 1 } + // The per-type caps spread data and parity independently, so together they + // can still stack e.g. 2 data + 1 parity on one rack. Cap the TOTAL per rack + // as well, at the lowest value the racks' free capacity allows. Attempts that + // enforce ReplicaPlacement size it from the room left under SameRackCount, so + // a rack of few nodes does not count disk room its nodes may not use. + totalShards := len(need) + for _, n := range rackShardCount { + totalShards += n + } + maxTotalPerRack := rackTotalCap(rackKeys, rackShardCount, rackRoom, totalShards) + maxTotalPerRackRP := rackTotalCap(rackKeys, rackShardCount, rackRoomRP, totalShards) + rackCap := func(rl relaxation) int { + if rl.rp { + return maxTotalPerRackRP + } + return maxTotalPerRack + } + attempts := strictAttempts if mode == PlaceDurabilityFirst { attempts = durabilityAttempts @@ -264,7 +297,7 @@ func (t *Topology) tryPlace(vk volKey, need []int, dataShards, parityShards int, typeTotal = parityShards } for _, rl := range attempts { - node, diskID, spilled, ok := chooseShardDest(vk, sid, isData, dataShards, typeTotal, numEligibleRacks, parityShards, racks, rackKeys, rp, eligible, prefer, shardsPerRack[isData], rackShardCount, bearing, rl) + node, diskID, spilled, ok := chooseShardDest(vk, sid, isData, dataShards, typeTotal, numEligibleRacks, parityShards, rackCap(rl), racks, rackKeys, rp, eligible, prefer, shardsPerRack[isData], rackShardCount, bearing, rl) if !ok { continue } @@ -303,7 +336,7 @@ func (t *Topology) tryPlace(vk volKey, need []int, dataShards, parityShards int, e.node.freeSlots++ racks[e.rackKey].freeSlots++ } - return nil, fmt.Errorf("cannot place EC shard %d of volume %d (collection %q)", sid, vk.vid, vk.collection) + return nil, fmt.Errorf("cannot place EC shard %d of volume %d (collection %q) (rack total cap %d, per-disk cap %d)", sid, vk.vid, vk.collection, rackCap(attempts[len(attempts)-1]), parityShards) } } @@ -316,11 +349,13 @@ func (t *Topology) tryPlace(vk volKey, need []int, dataShards, parityShards int, } // chooseShardDest selects a (node, disk) for one shard at the given relaxation -// level: pick a rack (even per-type cap + ReplicaPlacement caps + two-pass -// anti-affinity to the opposite type), then the least-loaded eligible node, then -// the best eligible disk. The third return reports whether the disk spilled off -// the soft-preferred type. ok=false when no rack/node/disk fits. -func chooseShardDest(vk volKey, sid int, isData bool, dataShards, typeTotal, numEligibleRacks, maxPerDisk int, racks map[string]*rack, rackKeys []string, rp *super_block.ReplicaPlacement, eligible func(*disk) bool, prefer func(*disk) bool, shardsPerRackType map[string][]int, rackShardCount map[string]int, bearing map[bool]map[string]bool, rl relaxation) (*Node, uint32, bool, bool) { +// level: pick a rack (total-shards-per-rack cap + even per-type cap + +// ReplicaPlacement caps + two-pass anti-affinity to the opposite type), then the +// least-loaded eligible node, then the best eligible disk. If no node in the +// chosen rack fits (e.g. every node is at SameRackCount), the next-best rack is +// tried. The third return reports whether the disk spilled off the +// soft-preferred type. ok=false when no rack/node/disk fits. +func chooseShardDest(vk volKey, sid int, isData bool, dataShards, typeTotal, numEligibleRacks, maxPerDisk, maxTotalPerRack int, racks map[string]*rack, rackKeys []string, rp *super_block.ReplicaPlacement, eligible func(*disk) bool, prefer func(*disk) bool, shardsPerRackType map[string][]int, rackShardCount map[string]int, bearing map[bool]map[string]bool, rl relaxation) (*Node, uint32, bool, bool) { maxPerRack := numEligibleRacks*typeTotal + 1 // effectively unlimited when caps are relaxed if rl.caps { if maxPerRack = ceilDivide(typeTotal, numEligibleRacks); maxPerRack < 1 { @@ -336,74 +371,127 @@ func chooseShardDest(vk volKey, sid int, isData bool, dataShards, typeTotal, num if !rl.rp { rp = nil } - // A rack is eligible only if it is under the per-rack shard cap (DiffRackCount), - // enforced only when set (and relaxed with rp). + // A rack is eligible only if it is under the total-shards-per-rack cap and, + // when set, the per-rack shard cap (DiffRackCount). The total cap does not + // depend on ReplicaPlacement, so it also holds when rp is nil. withinLimit := func(r string) bool { - if rp == nil { - return true + if rl.rackTotal && rackShardCount[r] >= maxTotalPerRack { + return false } - if rp.DiffRackCount > 0 && rackShardCount[r] >= rp.DiffRackCount { + if rp != nil && rp.DiffRackCount > 0 && rackShardCount[r] >= rp.DiffRackCount { return false } return true } - destRack, ok := pickTarget(rackKeys, shardsPerRackType, maxPerRack, anti, - func(r string) bool { return racks[r].freeSlots > 0 && rackHasFreeDisk(racks[r], eligible) }, - withinLimit) - if !ok { - return nil, 0, false, false + tried := map[string]bool{} + hasRoom := func(r string) bool { + return !tried[r] && racks[r].freeSlots > 0 && rackShardRoom(racks[r], vk, eligible, maxPerDisk) > 0 } - node := pickNodeInRackEligible(racks[destRack], vk, rp, eligible) - if node == nil { - return nil, 0, false, false + for { + destRack, ok := pickTarget(rackKeys, shardsPerRackType, maxPerRack, anti, hasRoom, withinLimit) + if !ok { + return nil, 0, false, false + } + if node := pickNodeInRackEligible(racks[destRack], vk, rp, eligible, maxPerDisk); node != nil { + if diskID, ok, spilled := pickBestDiskEligible(node, vk, eligible, prefer, sid, dataShards, maxPerDisk); ok { + return node, diskID, spilled, true + } + } + tried[destRack] = true } - diskID, ok, spilled := pickBestDiskEligible(node, vk, eligible, prefer, sid, dataShards, maxPerDisk) - if !ok { - return nil, 0, false, false - } - return node, diskID, spilled, true } -// nodeHasFreeDisk reports whether the node has a free disk satisfying eligible. -func nodeHasFreeDisk(n *Node, eligible func(*disk) bool) bool { +// rackTotalCap returns the smallest per-rack total c such that the racks can +// hold totalShards shards of the volume with none above c, given each rack's +// shards already placed (held) and its room for more. On a uniform cluster this +// is ceil(totalShards/racks); a nearly full rack raises it just enough for the +// other racks to absorb its share, so the cap alone never makes a feasible +// placement fail. +// +// It is the most even spread the free capacity allows, not a durability +// guarantee: losing k racks loses up to k*c shards, which the volume survives +// only while k*c <= parityShards. With few racks no placement can achieve that. +func rackTotalCap(rackKeys []string, held, room map[string]int, totalShards int) int { + for c := 1; c < totalShards; c++ { + fits := 0 + for _, rk := range rackKeys { + fits += max(held[rk], min(c, held[rk]+room[rk])) + } + if fits >= totalShards { + return c + } + } + return totalShards +} + +// diskShardRoom returns how many more shards of the volume the disk can take: +// its free slots, bounded by the per-disk durability cap (maxPerDisk shards of +// one volume per disk; <= 0 disables it). +func diskShardRoom(n *Node, d *disk, vk volKey, maxPerDisk int) int { + room := d.freeSlots + if maxPerDisk > 0 { + held := 0 + if info := n.shards[vk]; info != nil { + held = info.diskShardBits[d.diskID].Count() + } + room = min(room, maxPerDisk-held) + } + return max(room, 0) +} + +// nodeShardRoom returns how many more shards of the volume the node's eligible +// disks can take, bounded by the node's free slots. +func nodeShardRoom(n *Node, vk volKey, eligible func(*disk) bool, maxPerDisk int) int { + room := 0 for _, d := range n.disks { - if d.freeSlots > 0 && eligible(d) { - return true + if eligible(d) { + room += diskShardRoom(n, d, vk, maxPerDisk) } } - return false + return max(min(room, n.freeSlots), 0) } -// rackHasFreeDisk reports whether any node in the rack has a free eligible disk. -func rackHasFreeDisk(r *rack, eligible func(*disk) bool) bool { +// rackShardRoom returns how many more shards of the volume the rack can take. +func rackShardRoom(r *rack, vk volKey, eligible func(*disk) bool, maxPerDisk int) int { + room := 0 for _, n := range r.nodes { - if n.freeSlots > 0 && nodeHasFreeDisk(n, eligible) { - return true - } + room += nodeShardRoom(n, vk, eligible, maxPerDisk) } - return false + return room } -// pickNodeInRackEligible is pickNodeInRack restricted to nodes that have a free -// eligible disk. FromActiveTopology keeps all disk types/tags in the snapshot, so -// without this a node with free volume slots but no eligible disk could be chosen. +// rackShardRoomUnderRP is rackShardRoom with each node further bounded by the +// shards it may still take under ReplicaPlacement's SameRackCount (max shards of +// the volume per node), which pickNodeInRackEligible enforces. +func rackShardRoomUnderRP(r *rack, vk volKey, eligible func(*disk) bool, maxPerDisk int, rp *super_block.ReplicaPlacement) int { + if rp == nil || rp.SameRackCount <= 0 { + return rackShardRoom(r, vk, eligible, maxPerDisk) + } + room := 0 + for _, n := range r.nodes { + room += min(nodeShardRoom(n, vk, eligible, maxPerDisk), max(rp.SameRackCount-volumeShardCount(n, vk), 0)) + } + return room +} + +// pickNodeInRackEligible is pickNodeInRack restricted to nodes with an eligible +// disk that can take another shard of the volume (free slot, under maxPerDisk). +// FromActiveTopology keeps all disk types/tags in the snapshot, so without this a +// node with free volume slots but no eligible disk could be chosen. // // Among eligible nodes it ranks by fewest shards of the volume per machine, then per // node, with free capacity breaking ties. The free-capacity tie-break (not sorted id) // keeps the lowest-id machine from winning every volume's first shard against the // shared encode snapshot and piling up load. -func pickNodeInRackEligible(r *rack, vk volKey, rp *super_block.ReplicaPlacement, eligible func(*disk) bool) *Node { +func pickNodeInRackEligible(r *rack, vk volKey, rp *super_block.ReplicaPlacement, eligible func(*disk) bool, maxPerDisk int) *Node { machineShards := countShardsByHost(vk, r.nodes) machineFree := freeSlotsByHost(r.nodes) var best *Node var bestMCount, bestMFree, bestNCount, bestNFree int for _, id := range sortedNodeKeys(r.nodes) { node := r.nodes[id] - if node.freeSlots <= 0 { - continue - } - if !nodeHasFreeDisk(node, eligible) { + if nodeShardRoom(node, vk, eligible, maxPerDisk) <= 0 { continue } count := volumeShardCount(node, vk) diff --git a/weed/storage/erasure_coding/ecbalancer/place_rack_cap_test.go b/weed/storage/erasure_coding/ecbalancer/place_rack_cap_test.go new file mode 100644 index 000000000..f3ca0d780 --- /dev/null +++ b/weed/storage/erasure_coding/ecbalancer/place_rack_cap_test.go @@ -0,0 +1,301 @@ +package ecbalancer + +import ( + "fmt" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding" + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" +) + +// Total-shards-per-rack cap for Place / PlaceDurabilityFirst. +// +// The per-type caps (ceil(data/racks) data, ceil(parity/racks) parity) spread +// each type evenly on their own, but together they allow e.g. 2 data + 1 parity +// on one rack. Place also caps the TOTAL per rack at the lowest value the racks' +// free capacity allows (rackTotalCap), under both modes and with a nil +// ReplicaPlacement. + +// placedTotalsPerRack checks every shard was placed and returns the number of +// shards per rack. +func placedTotalsPerRack(t *testing.T, res *PlaceResult) map[string]int { + t.Helper() + if len(res.Destinations) != erasure_coding.TotalShardsCount { + t.Fatalf("placed %d shards, want %d", len(res.Destinations), erasure_coding.TotalShardsCount) + } + totals := map[string]int{} + for _, d := range res.Destinations { + totals[d.Rack]++ + } + t.Logf("total shards per rack: %v", totals) + return totals +} + +func assertMaxPerRack(t *testing.T, totals map[string]int, maxAllowed int) { + t.Helper() + for rk, n := range totals { + if n > maxAllowed { + t.Errorf("rack %s holds %d shards, want at most %d", rk, n, maxAllowed) + } + } +} + +func assertNotRelaxed(t *testing.T, res *PlaceResult, name string) { + t.Helper() + for _, r := range res.Relaxed { + if r == name { + t.Errorf("unexpected %q relaxation, got %v", name, res.Relaxed) + } + } +} + +// setRackFreeSlots leaves the rack with `slots` free slots, all on the first +// node's disk 0 (buildPlaceTopo gives each node one disk). Node and disk free +// slots are kept consistent, since rack capacity is derived from them. +func setRackFreeSlots(topo *Topology, rackKey string, slots int) { + first := true + for _, id := range sortedNodeKeys(topo.nodes) { + n := topo.nodes[id] + if n.rack != rackKey { + continue + } + free := 0 + if first { + free, first = slots, false + } + for _, d := range n.disks { + d.freeSlots = free + } + n.freeSlots = free + } +} + +// TestPlaceTotalShardsPerRackCap: 10+4 over 8 racks of two nodes each places at +// most ceil(14/8) = 2 shards per rack, in both modes. Before the cap this was +// 3,3,2,2,1,1,1,1: two racks with 3 each. +func TestPlaceTotalShardsPerRackCap(t *testing.T) { + modes := []struct { + name string + mode PlacementMode + }{{"strict", PlaceStrict}, {"durability-first", PlaceDurabilityFirst}} + for _, m := range modes { + t.Run(m.name, func(t *testing.T) { + topo := buildPlaceTopo(8, 2, 50) + res, err := topo.Place(1, "c1", allShards(), Constraints{}, m.mode) + if err != nil { + t.Fatalf("Place: %v", err) + } + assertMaxPerRack(t, placedTotalsPerRack(t, res), 2) + }) + } +} + +// TestPlaceTotalCapWithStarvedRack: a rack with a single free slot still takes +// one shard and the other seven racks keep to 2 each. +func TestPlaceTotalCapWithStarvedRack(t *testing.T) { + topo := buildPlaceTopo(8, 2, 50) + setRackFreeSlots(topo, "dc1:rack0", 1) + + res, err := topo.Place(1, "c1", allShards(), Constraints{}, PlaceDurabilityFirst) + if err != nil { + t.Fatalf("Place: %v", err) + } + totals := placedTotalsPerRack(t, res) + assertMaxPerRack(t, totals, 2) + assertNotRelaxed(t, res, "rack-total-cap") +} + +// TestPlaceTotalCapRisesForNearlyFullRacks: when nearly full racks cannot take +// an even share, the cap rises so the roomy racks absorb it instead of failing. +// A plain ceil(14/eligibleRacks) cap failed both cases. +func TestPlaceTotalCapRisesForNearlyFullRacks(t *testing.T) { + cases := []struct { + name string + racks int + starved map[string]int // rack -> free slots + wantCap int + wantHeld map[string]int // exact shard count expected on starved racks + }{ + { + // ceil(14/3) = 5 would fit only 5+5+1 = 11. + name: "3 racks, one with 1 slot", + racks: 3, + starved: map[string]int{"dc1:rack2": 1}, + wantCap: 7, + wantHeld: map[string]int{"rack2": 1}, + }, + { + // ceil(14/8) = 2 would fit only 2+7 = 9. + name: "8 racks, seven with 1 slot", + racks: 8, + starved: map[string]int{ + "dc1:rack1": 1, "dc1:rack2": 1, "dc1:rack3": 1, "dc1:rack4": 1, + "dc1:rack5": 1, "dc1:rack6": 1, "dc1:rack7": 1, + }, + wantCap: 7, + wantHeld: map[string]int{ + "rack1": 1, "rack2": 1, "rack3": 1, "rack4": 1, + "rack5": 1, "rack6": 1, "rack7": 1, + }, + }, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + topo := buildPlaceTopo(tc.racks, 2, 50) + for rk, slots := range tc.starved { + setRackFreeSlots(topo, rk, slots) + } + res, err := topo.Place(1, "c1", allShards(), Constraints{}, PlaceDurabilityFirst) + if err != nil { + t.Fatalf("Place: %v", err) + } + totals := placedTotalsPerRack(t, res) + assertMaxPerRack(t, totals, tc.wantCap) + for rk, want := range tc.wantHeld { + if totals[rk] != want { + t.Errorf("rack %s holds %d shards, want %d", rk, totals[rk], want) + } + } + assertNotRelaxed(t, res, "rack-total-cap") + }) + } +} + +// TestPlaceTotalCapKeepsPreferredTagTier: a preferred-tag tier whose racks have +// uneven room must still place the whole volume on the tagged disks rather than +// spill to untagged ones. +func TestPlaceTotalCapKeepsPreferredTagTier(t *testing.T) { + topo := NewTopology() + for r := 0; r < 3; r++ { + for n := 0; n < 2; n++ { + taggedFree := 50 + if r == 2 { + taggedFree = 0 + if n == 0 { + taggedFree = 2 + } + } + node := topo.AddNode(fmt.Sprintf("10.0.%d.%d:8080", r, n), "dc1", fmt.Sprintf("dc1:rack%d", r), taggedFree+50) + node.AddDisk(0, "", taggedFree, 0) + node.AddDiskTags(0, []string{"ec-pool"}) + node.AddDisk(1, "", 50, 0) + } + } + + res, err := topo.Place(1, "c1", allShards(), Constraints{PreferredTags: []string{"ec-pool"}}, PlaceDurabilityFirst) + if err != nil { + t.Fatalf("Place: %v", err) + } + if res.SpilledOutsidePreferredTags { + t.Fatal("placement spilled outside the preferred tags although the tagged disks can hold every shard") + } + for sid, d := range res.Destinations { + if d.DiskID != 0 { + t.Errorf("shard %d landed on untagged disk %d of %s", sid, d.DiskID, d.Node) + } + } + placedTotalsPerRack(t, res) +} + +// TestPlaceSkipsRackWithFullDisks: once a rack's disks hit the per-disk cap +// (parityShards shards of the volume), Place moves on to a rack that still has +// room instead of failing the shard. +func TestPlaceSkipsRackWithFullDisks(t *testing.T) { + topo := buildPlaceTopo(1, 3, 50) // rack0: 3 nodes + small := topo.AddNode("10.0.1.0:8080", "dc1", "dc1:rack1", 50) + small.AddDisk(0, "", 50, 0) // rack1: a single disk, so at most 4 shards + + res, err := topo.Place(1, "c1", allShards(), Constraints{}, PlaceDurabilityFirst) + if err != nil { + t.Fatalf("Place: %v", err) + } + totals := placedTotalsPerRack(t, res) + if totals["rack1"] > erasure_coding.ParityShardsCount { + t.Errorf("rack1 holds %d shards on one disk, want at most %d", totals["rack1"], erasure_coding.ParityShardsCount) + } +} + +// TestPlaceTotalCapCountsSameRackCount: with SameRackCount=1, a one-node rack +// takes a single shard however much disk room it has, so the cap must be sized +// from the room left under that limit. Four one-node racks and four three-node +// racks fit 10+4 as 1 per small rack and up to 3 per large rack; a cap sized from +// disk room alone is 2, which fits only 4 + 4*2 = 12, so strict placement failed +// and durability-first relaxed replica placement. +func TestPlaceTotalCapCountsSameRackCount(t *testing.T) { + rp := &super_block.ReplicaPlacement{SameRackCount: 1} // <=1 shard per node + modes := []struct { + name string + mode PlacementMode + }{{"strict", PlaceStrict}, {"durability-first", PlaceDurabilityFirst}} + for _, m := range modes { + t.Run(m.name, func(t *testing.T) { + topo := NewTopology() + for r := 0; r < 8; r++ { + nodes := 1 + if r >= 4 { + nodes = 3 + } + for n := 0; n < nodes; n++ { + node := topo.AddNode(fmt.Sprintf("10.0.%d.%d:8080", r, n), "dc1", fmt.Sprintf("dc1:rack%d", r), 50) + node.AddDisk(0, "", 50, 0) + } + } + res, err := topo.Place(1, "c1", allShards(), Constraints{ReplicaPlacement: rp}, m.mode) + if err != nil { + t.Fatalf("Place: %v", err) + } + assertMaxPerRack(t, placedTotalsPerRack(t, res), 3) + assertNotRelaxed(t, res, "replica-placement") + assertNotRelaxed(t, res, "rack-total-cap") + perNode := map[string]int{} + for _, d := range res.Destinations { + perNode[d.Node]++ + } + for node, n := range perNode { + if n > rp.SameRackCount { + t.Errorf("node %s holds %d shards, want at most %d", node, n, rp.SameRackCount) + } + } + }) + } +} + +func TestRackTotalCap(t *testing.T) { + keys := func(n int) []string { + out := make([]string, n) + for i := range out { + out[i] = fmt.Sprintf("r%d", i) + } + return out + } + uniform := func(n, room int) map[string]int { + out := map[string]int{} + for _, k := range keys(n) { + out[k] = room + } + return out + } + cases := []struct { + name string + racks int + held map[string]int + room map[string]int + total int + want int + }{ + {name: "even over 8", racks: 8, room: uniform(8, 50), total: 14, want: 2}, + {name: "even over 5", racks: 5, room: uniform(5, 50), total: 14, want: 3}, + {name: "one rack with 1 slot", racks: 3, room: map[string]int{"r0": 50, "r1": 50, "r2": 1}, total: 14, want: 7}, + {name: "one rack full", racks: 3, room: map[string]int{"r0": 50, "r1": 50}, total: 14, want: 7}, + // Repair: r0 already holds 5 survivors; the other racks share the rest. + {name: "held above even share", racks: 3, held: map[string]int{"r0": 5}, room: uniform(3, 50), total: 14, want: 5}, + {name: "not enough room", racks: 2, room: uniform(2, 3), total: 14, want: 14}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + if got := rackTotalCap(keys(tc.racks), tc.held, tc.room, tc.total); got != tc.want { + t.Errorf("rackTotalCap = %d, want %d", got, tc.want) + } + }) + } +} diff --git a/weed/worker/tasks/erasure_coding/detection.go b/weed/worker/tasks/erasure_coding/detection.go index ee73969be..60234e7d5 100644 --- a/weed/worker/tasks/erasure_coding/detection.go +++ b/weed/worker/tasks/erasure_coding/detection.go @@ -485,9 +485,10 @@ func buildNodeAddressMap(at *topology.ActiveTopology) map[string]string { // in the same detection cycle see the reduced capacity. Rebuilding it per volume // is O(volumes × topology) and times out on large clusters. // -// Encode is lenient (PlaceDurabilityFirst): it relaxes caps/anti-affinity/RP as -// needed rather than fail, and prefers the source disk type but spills if that -// type can't hold every shard. rp is the resolved replica placement (may be nil). +// Encode is lenient (PlaceDurabilityFirst): it relaxes caps/anti-affinity/RP and, +// last, the total-shards-per-rack cap as needed, failing only when no eligible +// disk has room. It prefers the source disk type but spills if that type can't +// hold every shard. rp is the resolved replica placement (may be nil). func planECDestinations(snap *ecbalancer.Topology, nodeAddresses map[string]string, metric *types.VolumeHealthMetrics, ecConfig *Config, rp *super_block.ReplicaPlacement, dataShards, parityShards int) (*topology.MultiDestinationPlan, [][]uint32, error) { if snap == nil { return nil, nil, fmt.Errorf("EC placement snapshot not available") @@ -529,8 +530,8 @@ func planECDestinations(snap *ecbalancer.Topology, nodeAddresses map[string]stri if len(res.Relaxed) > 0 { // Encode is best-effort (PlaceDurabilityFirst): it relaxes these constraints // rather than defer when the cluster can't satisfy them. Surface it so a tight - // replica placement isn't silently weakened; rebalancing tightens the spread. - glog.Warningf("EC volume %d: placed with relaxed constraints %v; replica placement not fully satisfied (rebalancing will adjust)", metric.VolumeID, res.Relaxed) + // replica placement or rack spread isn't silently weakened. + glog.Warningf("EC volume %d: placed with relaxed placement constraints %v", metric.VolumeID, res.Relaxed) } // Group the per-shard destinations into one plan per (node,disk), iterating