diff --git a/master/api_service.go b/master/api_service.go index 00a863023..73a6a4286 100644 --- a/master/api_service.go +++ b/master/api_service.go @@ -4259,7 +4259,12 @@ func (m *Server) updateNodesetCapcity(zoneName string, nodesetId uint64, capcity return } - ns.Capacity = int(capcity) + ns.racksLock.Lock() + for _, rack := range ns.racks { + rack.Capacity = int(capcity) / 3 + } + ns.racksLock.Unlock() + m.cluster.syncUpdateNodeSet(ns) log.LogInfof("action[updateNodesetCapcity] zonename %v nodeset %v capcity %v", zoneName, nodesetId, capcity) return diff --git a/master/cluster.go b/master/cluster.go index 7068813fa..00fb81991 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -27,6 +27,7 @@ import ( "sync/atomic" "time" + pt "github.com/cubefs/cubefs/proto" "github.com/cubefs/cubefs/util/auditlog" "github.com/google/uuid" "golang.org/x/time/rate" @@ -521,6 +522,33 @@ func (c *Cluster) masterAddr() (addr string) { return c.leaderInfo.addr } +func (c *Cluster) getRackAwareLevel() pt.RackAwareLevel { + return c.cfg.RackAwareLevel +} + +func (c *Cluster) GetExRacksByHosts(nodeType uint32, hosts []string) (exRacks []string) { + for _, host := range hosts { + if nodeType == TypeDataPartition { + node, err := c.dataNode(host) + if err != nil { + log.LogErrorf("getExRacksByHosts, node[%v] err[%v]", host, err.Error()) + continue + } + exRacks = append(exRacks, node.Rack) + continue + } + + node, err := c.metaNode(host) + if err != nil { + log.LogErrorf("getExRacksByHosts, node[%v] err[%v]", host, err.Error()) + continue + } + + exRacks = append(exRacks, node.Rack) + } + return +} + func (c *Cluster) tryToChangeLeaderByHost() error { return c.partition.TryToLeader(1) } @@ -2055,7 +2083,7 @@ func (c *Cluster) createDataPartition(volName string, mediaType uint32) (dp *Dat } else { zoneNum := c.decideZoneNum(vol, mediaType) // zoneNum scope [1,3] if targetHosts, targetPeers, err = c.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, - int(dpReplicaNum), zoneNum, zoneName, mediaType); err != nil { + int(dpReplicaNum), zoneNum, zoneName, mediaType, c.getRackAwareLevel()); err != nil { goto errHandler } } @@ -2261,11 +2289,11 @@ func (c *Cluster) decideZoneNum(vol *Vol, mediaType uint32) (zoneNum int) { return zoneNum } -func (c *Cluster) chooseZone2Plus1(rsMgr *rsManager, zones []*Zone, excludeNodeSets []uint64, excludeHosts []string, - nodeType uint32, replicaNum int) (hosts []string, peers []proto.Peer, err error, +func (c *Cluster) chooseZone2Plus1(rsMgr *rsManager, zones []*Zone, + nodeType uint32, param *selectParam) (hosts []string, peers []proto.Peer, err error, ) { - if replicaNum < 2 || replicaNum > defaultReplicaNum { - return nil, nil, fmt.Errorf("action[chooseZone2Plus1] replicaNum [%v]", replicaNum) + if param.replicaNum < 2 || param.replicaNum > defaultReplicaNum { + return nil, nil, fmt.Errorf("action[chooseZone2Plus1] replicaNum [%v]", param.replicaNum) } zoneList := make([]*Zone, 2) for i := range []int{1, 2} { @@ -2281,45 +2309,55 @@ func (c *Cluster) chooseZone2Plus1(rsMgr *rsManager, zones []*Zone, excludeNodeS log.LogInfof("action[chooseZone2Plus1] type [%v] after check,zone0 [%v] left [%v] zone1 [%v] left [%v]", nodeType, zoneList[0].name, zoneList[0].getSpaceLeft(nodeType), zoneList[1].name, zoneList[1].getSpaceLeft(nodeType)) + paramCopy := param.copy() + paramCopy.excludeRacks = c.GetExRacksByHosts(nodeType, param.excludeHosts) num := 1 for _, zone := range zoneList { - selectedHosts, selectedPeers, e := zone.getAvailNodeHosts(nodeType, excludeNodeSets, excludeHosts, num) + paramCopy.replicaNum = num + + selectedHosts, selectedPeers, e := zone.getAvailNodeHosts(nodeType, paramCopy) if e != nil { - log.LogErrorf("action[chooseZone2Plus1] getAvailNodeHosts error: [%v]", e) + log.LogErrorf("action[chooseZone2Plus1] getAvailNodeHosts param[%v] error: [%v]", paramCopy.String(), e) return nil, nil, e } hosts = append(hosts, selectedHosts...) peers = append(peers, selectedPeers...) + + paramCopy.excludeHosts = append(paramCopy.excludeHosts, selectedHosts...) + paramCopy.excludeRacks = c.GetExRacksByHosts(nodeType, param.excludeHosts) + log.LogInfof("action[chooseZone2Plus1] zone [%v] left [%v] get hosts[%v]", zone.name, zone.getSpaceLeft(nodeType), selectedHosts) - num = replicaNum - num + num = param.replicaNum - num } log.LogInfof("action[chooseZone2Plus1] finally get hosts[%v]", hosts) return hosts, peers, nil } -func (c *Cluster) chooseZoneNormal(zones []*Zone, excludeNodeSets []uint64, excludeHosts []string, - nodeType uint32, replicaNum int, -) (hosts []string, peers []proto.Peer, err error) { - log.LogInfof("action[chooseZoneNormal] zones[%s] nodeType[%d] replicaNum[%d]", printZonesName(zones), nodeType, replicaNum) +func (c *Cluster) chooseZoneNormal(zones []*Zone, nodeType uint32, param *selectParam) (hosts []string, peers []proto.Peer, err error) { + log.LogInfof("action[chooseZoneNormal] zones[%s] nodeType[%d] replicaNum[%d]", printZonesName(zones), nodeType, param.replicaNum) c.zoneIdxMux.Lock() defer c.zoneIdxMux.Unlock() - for i := 0; i < replicaNum; i++ { + paramCopy := param.copy() + paramCopy.excludeRacks = c.GetExRacksByHosts(nodeType, param.excludeHosts) + + for i := 0; i < param.replicaNum; i++ { // try all zone from zones list to get the available hosts for j := 0; j < len(zones); j++ { zone := zones[c.lastZoneIdxForNode] c.lastZoneIdxForNode = (c.lastZoneIdxForNode + 1) % len(zones) - selectedHosts, selectedPeers, err := zone.getAvailNodeHosts(nodeType, excludeNodeSets, excludeHosts, 1) + paramCopy.replicaNum = 1 + selectedHosts, selectedPeers, err := zone.getAvailNodeHosts(nodeType, paramCopy) if err != nil { // no zone available if j == len(zones)-1 { - log.LogErrorf("action[chooseZoneNormal] error [%v]", err) + log.LogErrorf("action[chooseZoneNormal] param [%s] error [%v]", paramCopy.String(), err) return nil, nil, err } continue @@ -2327,6 +2365,10 @@ func (c *Cluster) chooseZoneNormal(zones []*Zone, excludeNodeSets []uint64, excl hosts = append(hosts, selectedHosts...) peers = append(peers, selectedPeers...) + + paramCopy.excludeHosts = append(paramCopy.excludeHosts, selectedHosts...) + paramCopy.excludeRacks = c.GetExRacksByHosts(nodeType, param.excludeHosts) + // if get the available hosts, choose zone for next replica break } @@ -2357,7 +2399,7 @@ func (c *Cluster) getSpecificZoneList(specifiedZone string) (zones []*Zone, err func (c *Cluster) getHostFromNormalZone(nodeType uint32, excludeZones []string, excludeNodeSets []uint64, excludeHosts []string, replicaNum int, zoneNumNeed int, - specifiedZoneName string, dataMediaType uint32) (hosts []string, peers []proto.Peer, err error, + specifiedZoneName string, dataMediaType uint32, rackLevel proto.RackAwareLevel) (hosts []string, peers []proto.Peer, err error, ) { log.LogInfof("[getHostFromNormalZone] dataMediaType(%v) nodeType(%v) replicaNum(%v) zoneNumNeed(%v) specifiedZoneName(%v)", proto.MediaTypeString(nodeType), nodeType, replicaNum, zoneNumNeed, specifiedZoneName) @@ -2387,9 +2429,17 @@ func (c *Cluster) getHostFromNormalZone(nodeType uint32, excludeZones []string, } } + param := &selectParam{ + excludeNodeSets: excludeNodeSets, + replicaNum: replicaNum, + excludeHosts: excludeHosts, + rackLevel: rackLevel, + excludeRacks: c.GetExRacksByHosts(nodeType, excludeHosts), + } + if len(zonesQualified) == 1 { log.LogInfof("action[getHostFromNormalZone] zones [%v]", zonesQualified[0].name) - if hosts, peers, err = zonesQualified[0].getAvailNodeHosts(nodeType, excludeNodeSets, excludeHosts, replicaNum); err != nil { + if hosts, peers, err = zonesQualified[0].getAvailNodeHosts(nodeType, param); err != nil { log.LogErrorf("action[getHostFromNormalZone] err[%v]", err) return } @@ -2399,13 +2449,14 @@ func (c *Cluster) getHostFromNormalZone(nodeType uint32, excludeZones []string, if excludeHosts == nil { excludeHosts = make([]string, 0) } + // The upper process tries to go and get the dedicated zones, and the latter tries to choose the right zones as possible. if c.cfg.DefaultNormalZoneCnt == defaultNormalCrossZoneCnt && len(zonesQualified) >= defaultNormalCrossZoneCnt || replicaNum == 1 { - if hosts, peers, err = c.chooseZoneNormal(zonesQualified, excludeNodeSets, excludeHosts, nodeType, replicaNum); err != nil { + if hosts, peers, err = c.chooseZoneNormal(zonesQualified, nodeType, param); err != nil { return } } else { - if hosts, peers, err = c.chooseZone2Plus1(rsMgr, zonesQualified, excludeNodeSets, excludeHosts, nodeType, replicaNum); err != nil { + if hosts, peers, err = c.chooseZone2Plus1(rsMgr, zonesQualified, nodeType, param); err != nil { return } } @@ -2988,6 +3039,14 @@ func (c *Cluster) migrateDataPartition(srcAddr, targetAddr string, dp *DataParti replica, _ = dp.getReplica(srcAddr) dp.RUnlock() + param := &selectParam{ + excludeNodeSets: nil, + replicaNum: 1, + excludeHosts: dp.Hosts, + rackLevel: c.getRackAwareLevel(), + excludeRacks: c.GetExRacksByHosts(TypeDataPartition, dp.Hosts), + } + // delete if not normal data partition if !proto.IsNormalDp(dp.PartitionType) { c.vols[dp.VolName].deleteDataPartition(c, dp) @@ -3024,6 +3083,7 @@ func (c *Cluster) migrateDataPartition(srcAddr, targetAddr string, dp *DataParti break } } + if err = c.checkMultipleReplicasOnSameMachine(finalHosts); err != nil { goto errHandler } @@ -3034,7 +3094,7 @@ func (c *Cluster) migrateDataPartition(srcAddr, targetAddr string, dp *DataParti log.LogErrorf("[migrateDataPartition] check mediaType err: %v", err.Error()) goto errHandler } - } else if targetHosts, _, err = ns.getAvailDataNodeHosts(dp.Hosts, 1); err != nil { + } else if targetHosts, _, err = ns.getAvailDataNodeHosts(param); err != nil { if _, ok := c.vols[dp.VolName]; !ok { log.LogWarnf("clusterID[%v] partitionID:%v on node:%v offline failed,PersistenceHosts:[%v]", c.Name, dp.PartitionID, srcAddr, dp.Hosts) @@ -3047,10 +3107,17 @@ func (c *Cluster) migrateDataPartition(srcAddr, targetAddr string, dp *DataParti } // select data nodes from the other node set in same zone excludeNodeSets = append(excludeNodeSets, ns.ID) - if targetHosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, excludeNodeSets, dp.Hosts, 1); err != nil { + param := &selectParam{ + excludeNodeSets: excludeNodeSets, + replicaNum: 1, + excludeHosts: dp.Hosts, + rackLevel: c.getRackAwareLevel(), + excludeRacks: c.GetExRacksByHosts(TypeDataPartition, dp.Hosts), + } + if targetHosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param); err != nil { // select data nodes from the other zone zones = dp.getLiveZones(srcAddr) - if targetHosts, _, err = c.getHostFromNormalZone(TypeDataPartition, zones, excludeNodeSets, dp.Hosts, 1, 1, "", dp.MediaType); err != nil { + if targetHosts, _, err = c.getHostFromNormalZone(TypeDataPartition, zones, excludeNodeSets, dp.Hosts, 1, 1, "", dp.MediaType, c.getRackAwareLevel()); err != nil { goto errHandler } } diff --git a/master/cluster_task.go b/master/cluster_task.go index 0c522de7f..133716b00 100644 --- a/master/cluster_task.go +++ b/master/cluster_task.go @@ -126,6 +126,14 @@ func (c *Cluster) migrateMetaPartition(srcAddr, targetAddr string, mp *MetaParti } mp.RUnlock() + param := &selectParam{ + excludeNodeSets: nil, + replicaNum: 1, + excludeHosts: oldHosts, + rackLevel: c.getRackAwareLevel(), + excludeRacks: c.GetExRacksByHosts(TypeMetaPartition, oldHosts), + } + nodeType := TypeMetaPartition if dstStoreMode == proto.StoreModeDef { dstStoreMode, err = c.getMetaPartitionStoreMode(mp, srcAddr) @@ -157,7 +165,7 @@ func (c *Cluster) migrateMetaPartition(srcAddr, targetAddr string, mp *MetaParti newPeers = []proto.Peer{{ Addr: targetAddr, }} - } else if _, newPeers, err = ns.getAvailMetaNodeHosts(oldHosts, 1, dstStoreMode); err != nil { + } else if _, newPeers, err = ns.getAvailMetaNodeHosts(param, dstStoreMode); err != nil { if _, ok := c.vols[mp.volName]; !ok { log.LogWarnf("[migrateMetaPartition] clusterID[%v] partitionID:%v on node:[%v]", c.Name, mp.PartitionID, mp.Hosts) @@ -170,7 +178,14 @@ func (c *Cluster) migrateMetaPartition(srcAddr, targetAddr string, mp *MetaParti } // choose a meta node in other node set in the same zone excludeNodeSets = append(excludeNodeSets, ns.ID) - if _, newPeers, err = zone.getAvailNodeHosts(nodeType, excludeNodeSets, oldHosts, 1); err != nil { + param := &selectParam{ + excludeNodeSets: excludeNodeSets, + replicaNum: 1, + excludeHosts: oldHosts, + rackLevel: c.getRackAwareLevel(), + excludeRacks: c.GetExRacksByHosts(nodeType, oldHosts), + } + if _, newPeers, err = zone.getAvailNodeHosts(nodeType, param); err != nil { zones = mp.getLiveZones(srcAddr) var excludeZone []string if len(zones) == 0 { @@ -179,7 +194,7 @@ func (c *Cluster) migrateMetaPartition(srcAddr, targetAddr string, mp *MetaParti excludeZone = append(excludeZone, zones[0]) } // choose a meta node in other zone - if _, newPeers, err = c.getHostFromNormalZone(nodeType, excludeZone, excludeNodeSets, oldHosts, 1, 1, "", proto.MediaType_Unspecified); err != nil { + if _, newPeers, err = c.getHostFromNormalZone(nodeType, excludeZone, excludeNodeSets, oldHosts, 1, 1, "", proto.MediaType_Unspecified, c.getRackAwareLevel()); err != nil { goto errHandler } } diff --git a/master/data_partition.go b/master/data_partition.go index ebf1031db..210804b4d 100644 --- a/master/data_partition.go +++ b/master/data_partition.go @@ -2256,7 +2256,14 @@ func (partition *DataPartition) TryAcquireDecommissionToken(c *Cluster) bool { log.LogDebugf("action[TryAcquireDecommissionToken]dp %v excludeHosts %v", partition.PartitionID, excludeHosts) // data nodes in a nodeset has the same mediaType - targetHosts, _, err = ns.getAvailDataNodeHosts(excludeHosts, 1) + param := &selectParam{ + excludeNodeSets: nil, + replicaNum: 1, + excludeHosts: excludeHosts, + rackLevel: c.getRackAwareLevel(), + excludeRacks: c.GetExRacksByHosts(TypeDataPartition, excludeHosts), + } + targetHosts, _, err = ns.getAvailDataNodeHosts(param) if err != nil { if partition.DecommissionDstNodeSet != 0 { log.LogWarnf("action[TryAcquireDecommissionToken] dp %v choose from given dst nodeset %v failed:%v", @@ -2278,13 +2285,20 @@ func (partition *DataPartition) TryAcquireDecommissionToken(c *Cluster) bool { goto errHandler } excludeNodeSets = append(excludeNodeSets, ns.ID) + param := &selectParam{ + excludeNodeSets: excludeNodeSets, + replicaNum: 1, + excludeHosts: excludeHosts, + rackLevel: c.getRackAwareLevel(), + excludeRacks: c.GetExRacksByHosts(TypeDataPartition, excludeHosts), + } // data nodes in a zone has the same mediaType - if targetHosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, excludeNodeSets, excludeHosts, 1); err != nil { + if targetHosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param); err != nil { log.LogWarnf("action[TryAcquireDecommissionToken] dp %v choose from other nodeset failed:%v", partition.PartitionID, err.Error()) // select data nodes from the other zone zones = partition.getLiveZones(partition.DecommissionSrcAddr) - if targetHosts, _, err = c.getHostFromNormalZone(TypeDataPartition, zones, excludeNodeSets, excludeHosts, 1, 1, "", partition.MediaType); err != nil { + if targetHosts, _, err = c.getHostFromNormalZone(TypeDataPartition, zones, excludeNodeSets, excludeHosts, 1, 1, "", partition.MediaType, c.getRackAwareLevel()); err != nil { log.LogWarnf("action[TryAcquireDecommissionToken] dp %v choose from other zone failed:%v", partition.PartitionID, err.Error()) goto errHandler diff --git a/master/node_selector.go b/master/node_selector.go index fdec186e9..9be6a955b 100644 --- a/master/node_selector.go +++ b/master/node_selector.go @@ -120,6 +120,7 @@ type CarryWeightNodeSelector struct { nodeType NodeType carry map[uint64]float64 + sync.RWMutex } func (s *CarryWeightNodeSelector) GetName() string { @@ -176,12 +177,17 @@ func (s *CarryWeightNodeSelector) getCarryNodes(nset *nodeSet, maxTotal uint64, node.GetStorageInfo(), excludeHosts) return true } + + s.RLock() if s.carry[node.GetID()] >= 1.0 { availCount++ } + s.RUnlock() nt := new(weightedNode) + s.RLock() nt.Carry = s.carry[node.GetID()] + s.RUnlock() nt.Weight = float64(node.GetTotal()-node.GetUsed()) / float64(maxTotal) nt.Ptr = node nodeTabs = append(nodeTabs, nt) @@ -201,7 +207,9 @@ func (s *CarryWeightNodeSelector) setNodeCarry(nodes SortedWeightedNodes, availC carry = 10.0 } nt.Carry = carry + s.Lock() s.carry[nt.Ptr.GetID()] = carry + s.Unlock() if carry > 1.0 { availCarryCount++ } @@ -212,7 +220,9 @@ func (s *CarryWeightNodeSelector) setNodeCarry(nodes SortedWeightedNodes, availC func (s *CarryWeightNodeSelector) selectNodeForWrite(node Node) { node.SelectNodeForWrite() // decrease node weight + s.Lock() s.carry[node.GetID()] -= 1.0 + s.Unlock() } func (s *CarryWeightNodeSelector) Select(ns *nodeSet, excludeHosts []string, replicaNum int) (newHosts []string, peers []proto.Peer, err error) { @@ -612,34 +622,106 @@ func NewNodeSelector(name string, nodeType NodeType) NodeSelector { } } -func (ns *nodeSet) getAvailMetaNodeHosts(excludeHosts []string, replicaNum int, storeMode proto.StoreMode) (newHosts []string, peers []proto.Peer, err error) { +func (ns *nodeSet) getRackSets() nodeSetCollection { + ns.racksLock.RLock() + defer ns.racksLock.RUnlock() + + rsets := make(nodeSetCollection, 0, len(ns.racks)) + for _, rack := range ns.racks { + rsets = append(rsets, rack) + } + return rsets +} + +func (ns *nodeSet) getAvailMetaNodeHosts(param *selectParam, storeMode proto.StoreMode) (newHosts []string, peers []proto.Peer, err error) { ns.nodeSelectLock.Lock() defer ns.nodeSelectLock.Unlock() // we need a read lock to block the modification of node selector ns.metaNodeSelectorLock.RLock() defer ns.metaNodeSelectorLock.RUnlock() - if storeMode == proto.StoreModeRocksDb { - return ns.metaNodeRocksdbSelector.Select(ns, excludeHosts, replicaNum) + + // If rack isolation is not enabled, use non-rack-aware selector directly + if param.rackLevel == proto.RackAwareNone { + return ns.getNodeSelector(MetaNodeType, storeMode).Select(ns, param.excludeHosts, param.replicaNum) } - return ns.metaNodeMemorySelector.Select(ns, excludeHosts, replicaNum) + + // Rack isolation enabled, prioritize strong constraint mode + return ns.selectNodesWithRack(param, MetaNodeType, storeMode) } -func (ns *nodeSet) getAvailDataNodeHosts(excludeHosts []string, replicaNum int) (hosts []string, peers []proto.Peer, err error) { +// selectMetaNodesWithRack selects meta nodes with rack awareness +func (ns *nodeSet) selectNodesWithRack(param *selectParam, nodeType NodeType, storeMode proto.StoreMode) (newHosts []string, peers []proto.Peer, err error) { + + rsets := ns.getRackSets() + + paramCopy := param.copy() + // First attempt with strong rack awareness + paramCopy.rackLevel = proto.RackAwareStrong + paramCopy.replicaNum = 1 + + for { + + rack, err := ns.getRackSelector(nodeType).Select(rsets, paramCopy) + if err != nil { + // if rack aware is not enabled or alreay weak aware, return error + if param.rackLevel == proto.RackAwareStrong || paramCopy.rackLevel == proto.RackAwareWeak { + return nil, nil, fmt.Errorf("strong rack aware selection failed for rack[%v], param[%v], err: %v", + rack, paramCopy.String(), err) + } + + log.LogWarnf("action[getAvailMetaNodeHosts] weak rack aware selection failed for rack[%v], param %v, err: %v", + rack, paramCopy.String(), err.Error()) + paramCopy.rackLevel = proto.RackAwareWeak + continue + } + + // Select nodes + selector := rack.getNodeSelector(nodeType, storeMode) + + rhosts, rpeers, err := selector.Select(rack, paramCopy.excludeHosts, 1) + if err != nil { + log.LogErrorf("action[getAvailMetaNodeHosts] node selection failed for rack[%v], param[%v], err: %v", + rack, paramCopy.String(), err.Error()) + return nil, nil, fmt.Errorf("node selection failed for rack[%v], param[%v], err: %v", + rack.Rack, paramCopy.String(), err.Error()) + } + + // Update results + newHosts = append(newHosts, rhosts...) + peers = append(peers, rpeers...) + paramCopy.excludeHosts = append(paramCopy.excludeHosts, rhosts...) + paramCopy.excludeRacks = append(paramCopy.excludeRacks, rack.Rack) + + // Check if replica number requirement is met + if len(newHosts) >= param.replicaNum { + return newHosts, peers, nil + } + } +} + +func (ns *nodeSet) getAvailDataNodeHosts(param *selectParam) (hosts []string, peers []proto.Peer, err error) { ns.nodeSelectLock.Lock() defer ns.nodeSelectLock.Unlock() // we need a read lock to block the modification of node selector ns.dataNodeSelectorLock.Lock() defer ns.dataNodeSelectorLock.Unlock() - return ns.dataNodeSelector.Select(ns, excludeHosts, replicaNum) + + if param.rackLevel == proto.RackAwareNone { + return ns.dataNodeSelector.Select(ns, param.excludeHosts, param.replicaNum) + } + + return ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeDef) } func (s *CarryWeightNodeSelector) prepareCarryForDataNodes(nodes *sync.Map, total uint64) { nodes.Range(func(key, value interface{}) bool { node := value.(Node) + s.Lock() if _, ok := s.carry[node.GetID()]; !ok { // use available space to calculate initial weight s.carry[node.GetID()] = float64(node.GetAvailableSpace()) / float64(total) } + s.Unlock() return true }) } @@ -647,10 +729,12 @@ func (s *CarryWeightNodeSelector) prepareCarryForDataNodes(nodes *sync.Map, tota func (s *CarryWeightNodeSelector) prepareCarryForMetaNodeMemory(nodes *sync.Map, total uint64) { nodes.Range(func(key, value interface{}) bool { metaNode := value.(*MetaNode) + s.Lock() if _, ok := s.carry[metaNode.ID]; !ok { // use available space to calculate initial weight s.carry[metaNode.ID] = float64(metaNode.Total-metaNode.Used) / float64(total) } + s.Unlock() return true }) } @@ -658,10 +742,12 @@ func (s *CarryWeightNodeSelector) prepareCarryForMetaNodeMemory(nodes *sync.Map, func (s *CarryWeightNodeSelector) prepareCarryForMetaNodeRocksdb(nodes *sync.Map, total uint64) { nodes.Range(func(key, value interface{}) bool { metaNode := value.(*MetaNode) + s.Lock() if _, ok := s.carry[metaNode.ID]; !ok { // use available space to calculate initial weight s.carry[metaNode.ID] = float64(metaNode.GetRocksdbTotal()-metaNode.GetRocksdbUsed()) / float64(total) } + s.Unlock() return true }) } diff --git a/master/nodeset_selector.go b/master/nodeset_selector.go index 804b397b4..82cc2a15a 100644 --- a/master/nodeset_selector.go +++ b/master/nodeset_selector.go @@ -88,12 +88,12 @@ func (ns *nodeSet) getMetaNodeTotalAvailableSpace(nodeType NodeType) (space uint return } -func (ns *nodeSet) canWriteFor(nodeType NodeType, replica int) bool { +func (ns *nodeSet) canWriteFor(nodeType NodeType, param *selectParam) bool { switch nodeType { case DataNodeType: - return ns.canWriteForNode(ns.dataNodes, replica, nodeType) + return ns.canWriteForNode(ns.dataNodes, param, nodeType) case MetaNodeType, RocksdbType: - return ns.canWriteForNode(ns.metaNodes, replica, nodeType) + return ns.canWriteForNode(ns.metaNodes, param, nodeType) default: panic("unknow node type") } @@ -123,7 +123,7 @@ func (ns *nodeSet) getTotalAvailableSpaceOf(nodeType NodeType) uint64 { type NodesetSelector interface { GetName() string - Select(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) + Select(nsc nodeSetCollection, param *selectParam) (ns *nodeSet, err error) } type RoundRobinNodesetSelector struct { @@ -132,7 +132,7 @@ type RoundRobinNodesetSelector struct { nodeType NodeType } -func (s *RoundRobinNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) { +func (s *RoundRobinNodesetSelector) Select(nsc nodeSetCollection, param *selectParam) (ns *nodeSet, err error) { // sort nodesets by id, so we can get a node list that is as stable as possible sort.Slice(nsc, func(i, j int) bool { return nsc[i].ID < nsc[j].ID @@ -145,10 +145,10 @@ func (s *RoundRobinNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSet ns = nsc[s.index] s.index++ - if containsID(excludeNodeSets, ns.ID) { + if containsID(param.excludeNodeSets, ns.ID) { continue } - if ns.canWriteFor(s.nodeType, int(replicaNum)) { + if ns.canWriteFor(s.nodeType, param) { return } } @@ -205,11 +205,11 @@ func (s *CarryWeightNodesetSelector) prepareCarry(nsc nodeSetCollection, total u } } -func (s *CarryWeightNodesetSelector) getAvailNodesets(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (newNsc nodeSetCollection) { +func (s *CarryWeightNodesetSelector) getAvailNodesets(nsc nodeSetCollection, param *selectParam) (newNsc nodeSetCollection) { newNsc = make(nodeSetCollection, 0, nsc.Len()) for i := 0; i < nsc.Len(); i++ { ns := nsc[i] - if ns.canWriteFor(s.nodeType, int(replicaNum)) && !containsID(excludeNodeSets, ns.ID) { + if ns.canWriteFor(s.nodeType, param) && !containsID(param.excludeNodeSets, ns.ID) { newNsc = append(newNsc, ns) } } @@ -246,11 +246,11 @@ func (s *CarryWeightNodesetSelector) setNodesetCarry(nsc nodeSetCollection, tota return count } -func (s *CarryWeightNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) { +func (s *CarryWeightNodesetSelector) Select(nsc nodeSetCollection, param *selectParam) (ns *nodeSet, err error) { total := s.getMaxTotal(nsc) // prepare weight of evert nodesets s.prepareCarry(nsc, total) - nsc = s.getAvailNodesets(nsc, excludeNodeSets, replicaNum) + nsc = s.getAvailNodesets(nsc, param) avaliCount := 0 if len(nsc) < 1 { goto err @@ -263,17 +263,19 @@ func (s *CarryWeightNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSe // pick the first nodeset than has N writable node for i := 0; i < avaliCount; i++ { ns = nsc[i] - if ns.canWriteFor(s.nodeType, int(replicaNum)) && !containsID(excludeNodeSets, ns.ID) { + if ns.canWriteFor(s.nodeType, param) && !containsID(param.excludeNodeSets, ns.ID) { break } } + if ns != nil { - if !ns.canWriteFor(s.nodeType, int(replicaNum)) || containsID(excludeNodeSets, ns.ID) { + if !ns.canWriteFor(s.nodeType, param) || containsID(param.excludeNodeSets, ns.ID) { goto err } s.carrys[ns.ID] -= 1.0 } return + err: switch s.nodeType { case DataNodeType: @@ -301,7 +303,7 @@ func (s *AvailableSpaceFirstNodesetSelector) GetName() string { return AvailableSpaceFirstNodesetSelectorName } -func (s *AvailableSpaceFirstNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) { +func (s *AvailableSpaceFirstNodesetSelector) Select(nsc nodeSetCollection, param *selectParam) (ns *nodeSet, err error) { // sort nodesets by available space sort.Slice(nsc, func(i, j int) bool { return nsc[i].getTotalAvailableSpaceOf(s.nodeType) > nsc[j].getTotalAvailableSpaceOf(s.nodeType) @@ -309,7 +311,7 @@ func (s *AvailableSpaceFirstNodesetSelector) Select(nsc nodeSetCollection, exclu // pick the first nodeset that has N writable nodes for i := 0; i < nsc.Len(); i++ { ns = nsc[i] - if ns.canWriteFor(s.nodeType, int(replicaNum)) && !containsID(excludeNodeSets, ns.ID) { + if ns.canWriteFor(s.nodeType, param) && !containsID(param.excludeNodeSets, ns.ID) { return } } @@ -348,10 +350,10 @@ func (s *StrawNodesetSelector) getWeight(ns *nodeSet) float64 { return float64(ns.getTotalAvailableSpaceOf(s.nodeType) / util.GB) } -func (s *StrawNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) { +func (s *StrawNodesetSelector) Select(nsc nodeSetCollection, param *selectParam) (ns *nodeSet, err error) { tmp := make(nodeSetCollection, 0) for _, nodeset := range nsc { - if nodeset.canWriteFor(s.nodeType, int(replicaNum)) && !containsID(excludeNodeSets, nodeset.ID) { + if nodeset.canWriteFor(s.nodeType, param) && !containsID(param.excludeNodeSets, nodeset.ID) { tmp = append(tmp, nodeset) } } diff --git a/master/nodeset_selector_test.go b/master/nodeset_selector_test.go index 90481c0b3..80cd452ec 100644 --- a/master/nodeset_selector_test.go +++ b/master/nodeset_selector_test.go @@ -80,7 +80,14 @@ func NodesetSelectorTest(t *testing.T, selector NodesetSelector) { return true }) } - ns, err := selector.Select(nsc, nil, 1) + param := &selectParam{ + excludeNodeSets: nil, + replicaNum: 1, + excludeHosts: nil, + rackLevel: proto.RackAwareNone, + excludeRacks: nil, + } + ns, err := selector.Select(nsc, param) if err != nil { t.Errorf("%v failed to select nodeset %v", selector.GetName(), err) return @@ -154,7 +161,14 @@ func prepareMetaNodesetForBench(count int, initTotal uint64, grow uint64) (nsc n func nodesetSelectorBench(selector NodesetSelector, nsc nodeSetCollection, onSelect func(id uint64)) (map[uint64]int, error) { times := make(map[uint64]int) for i := 0; i < loopNodeSelectorTestCount; i++ { - ns, err := selector.Select(nsc, nil, 1) + param := &selectParam{ + excludeNodeSets: nil, + replicaNum: 1, + excludeHosts: nil, + rackLevel: proto.RackAwareNone, + excludeRacks: nil, + } + ns, err := selector.Select(nsc, param) if err != nil { return nil, err } diff --git a/master/topology.go b/master/topology.go index bcb0b5517..0a5122f3e 100644 --- a/master/topology.go +++ b/master/topology.go @@ -29,6 +29,42 @@ import ( "github.com/cubefs/cubefs/util/log" ) +type selectParam struct { + excludeHosts []string + replicaNum int + rackLevel proto.RackAwareLevel + excludeRacks []string + excludeNodeSets []uint64 +} + +func (sp *selectParam) copy() *selectParam { + // deep copy slice fields + excludeHosts := make([]string, len(sp.excludeHosts)) + copy(excludeHosts, sp.excludeHosts) + + excludeRacks := make([]string, len(sp.excludeRacks)) + copy(excludeRacks, sp.excludeRacks) + + excludeNodeSets := make([]uint64, len(sp.excludeNodeSets)) + copy(excludeNodeSets, sp.excludeNodeSets) + + return &selectParam{ + excludeHosts: excludeHosts, + replicaNum: sp.replicaNum, + rackLevel: sp.rackLevel, + excludeRacks: excludeRacks, + excludeNodeSets: excludeNodeSets, + } +} + +func (sp *selectParam) String() string { + if sp == nil { + return "nil" + } + return fmt.Sprintf("excludeHosts: %v, replicaNum: %v, rackLevel: %v, excludeRacks: %v, excludeNodeSets: %v", + sp.excludeHosts, sp.replicaNum, sp.rackLevel, sp.excludeRacks, sp.excludeNodeSets) +} + type rsManager struct { nodeType NodeType nodes *sync.Map @@ -581,7 +617,15 @@ func (nsgm *DomainManager) getHostFromNodeSetGrpSpecific(domainGrpManager *Domai } if createType == TypeDataPartition { - if host, peer, err = ns.getAvailDataNodeHosts(nil, needNum); err != nil { + param := &selectParam{ + excludeNodeSets: nil, + replicaNum: needNum, + excludeHosts: nil, + rackLevel: proto.RackAwareNone, + excludeRacks: nil, + } + + if host, peer, err = ns.getAvailDataNodeHosts(param); err != nil { log.LogErrorf("action[getHostFromNodeSetGrpSpecific] ns[%v] zone[%v] TypeDataPartition err[%v]", ns.ID, ns.zoneName, err) // nsg.status = dataNodesUnAvailable continue @@ -591,7 +635,14 @@ func (nsgm *DomainManager) getHostFromNodeSetGrpSpecific(domainGrpManager *Domai if createType == TypeRocksdbPartition { storeMode = proto.StoreModeRocksDb } - if host, peer, err = ns.getAvailMetaNodeHosts(nil, needNum, storeMode); err != nil { + param := &selectParam{ + excludeNodeSets: nil, + replicaNum: needNum, + excludeHosts: nil, + rackLevel: proto.RackAwareNone, + excludeRacks: nil, + } + if host, peer, err = ns.getAvailMetaNodeHosts(param, storeMode); err != nil { log.LogErrorf("action[getHostFromNodeSetGrpSpecific] ns[%v] zone[%v] type(%d) err[%v]", ns.ID, ns.zoneName, createType, err) // nsg.status = metaNodesUnAvailable continue @@ -691,7 +742,14 @@ func (nsgm *DomainManager) getHostFromNodeSetGrp(domainId uint64, replicaNum uin log.LogWarnf("action[getHostFromNodeSetGrp] ns[%v] zone[%v] dataNodesUnAvailable", ns.ID, ns.zoneName) continue } - if host, peer, err = ns.getAvailDataNodeHosts(hosts, 1); err != nil { + param := &selectParam{ + excludeNodeSets: nil, + replicaNum: 1, + excludeHosts: hosts, + rackLevel: proto.RackAwareNone, + excludeRacks: nil, + } + if host, peer, err = ns.getAvailDataNodeHosts(param); err != nil { log.LogWarnf("action[getHostFromNodeSetGrp] ns[%v] zone[%v] TypeDataPartition err[%v]", ns.ID, ns.zoneName, err) // nsg.status = dataNodesUnAvailable continue @@ -705,7 +763,14 @@ func (nsgm *DomainManager) getHostFromNodeSetGrp(domainId uint64, replicaNum uin if createType == TypeRocksdbPartition { storeMode = proto.StoreModeRocksDb } - if host, peer, err = ns.getAvailMetaNodeHosts(hosts, 1, storeMode); err != nil { + param := &selectParam{ + excludeNodeSets: nil, + replicaNum: 1, + excludeHosts: hosts, + rackLevel: proto.RackAwareNone, + excludeRacks: nil, + } + if host, peer, err = ns.getAvailMetaNodeHosts(param, storeMode); err != nil { log.LogWarnf("action[getHostFromNodeSetGrp] ns[%v] zone[%v] ModeRocksDb err[%v]", ns.ID, ns.zoneName, err) // nsg.status = metaNodesUnAvailable continue @@ -988,8 +1053,10 @@ type nodeSet struct { metaNodes *sync.Map dataNodes *sync.Map - racks map[string]*nodeSet - racksLock sync.RWMutex + racks map[string]*nodeSet + racksLock sync.RWMutex + dataRackSelector NodesetSelector + metaRackSelector NodesetSelector nodeSelectLock sync.Mutex dataNodeSelectorLock sync.RWMutex @@ -1025,13 +1092,17 @@ type nodeSetDecommissionParallelStatus struct { func newNodeSet(c *Cluster, id uint64, cap int, zoneName string, rack string) *nodeSet { log.LogInfof("action[newNodeSet] id[%v]", id) ns := &nodeSet{ - ID: id, - Rack: rack, - Capacity: cap, - zoneName: zoneName, - metaNodes: new(sync.Map), - dataNodes: new(sync.Map), - racks: make(map[string]*nodeSet), + ID: id, + Rack: rack, + Capacity: cap, + zoneName: zoneName, + metaNodes: new(sync.Map), + dataNodes: new(sync.Map), + + racks: make(map[string]*nodeSet), + dataRackSelector: NewNodesetSelector(StrawNodesetSelectorName, DataNodeType), + metaRackSelector: NewNodesetSelector(StrawNodesetSelectorName, MetaNodeType), + decommissionDataPartitionList: NewDecommissionDataPartitionList(c, id), manualDecommissionDiskList: NewDecommissionDiskList(), autoDecommissionDiskList: NewDecommissionDiskList(), @@ -1049,6 +1120,31 @@ func newNodeSet(c *Cluster, id uint64, cap int, zoneName string, rack string) *n return ns } +func (ns *nodeSet) getRackSelector(nodeType NodeType) NodesetSelector { + if nodeType == DataNodeType { + return ns.dataRackSelector + } + return ns.metaRackSelector +} + +func (ns *nodeSet) getNodeSelector(nodeType NodeType, storeMode proto.StoreMode) NodeSelector { + if nodeType == DataNodeType { + return ns.dataNodeSelector + } + if storeMode == proto.StoreModeRocksDb { + return ns.metaNodeRocksdbSelector + } + return ns.metaNodeMemorySelector +} + +func (ns *nodeSet) IsRackSet() bool { + return ns.Rack != "" +} + +func (ns *nodeSet) IsNodeSet() bool { + return ns.Rack == "" +} + func (ns *nodeSet) GetDataNodeSelector() string { ns.dataNodeSelectorLock.RLock() defer ns.dataNodeSelectorLock.RUnlock() @@ -1137,28 +1233,82 @@ func (ns *nodeSet) deleteMetaNode(metaNode *MetaNode) { } -func (ns *nodeSet) canWriteForNode(nodes *sync.Map, replicaNum int, nodeType NodeType) bool { +// Extract node check logic +func (ns *nodeSet) checkNodeWriteable(node interface{}, nodeType NodeType, param *selectParam) bool { + if nodeType == RocksdbType { + if metaNode, ok := node.(*MetaNode); ok { + return metaNode.IsRocksdbWriteAble() && metaNode.PartitionCntLimited() && !contains(param.excludeHosts, metaNode.Addr) + } + return false + } + + if n, ok := node.(Node); ok { + return n.IsWriteAble() && n.PartitionCntLimited() && !contains(param.excludeHosts, n.GetAddr()) + } + return false +} + +// Optimized main function +func (ns *nodeSet) canWriteForNode(nodes *sync.Map, param *selectParam, nodeType NodeType) bool { + + // Rack-aware strong consistency mode, select available rack set count >= replicaNum + if param.rackLevel == proto.RackAwareStrong && ns.IsNodeSet() { + return ns.checkRackAwareWriteable(param, nodeType) + } + + // Rack-aware strong consistency mode, exclude rack + if ns.IsRackSet() && param.rackLevel == proto.RackAwareStrong && contains(param.excludeRacks, ns.Rack) { + return false + } + + // Normal mode + return ns.checkNormalWriteable(nodes, param, nodeType) +} + +// Rack-aware consistency check, select available rack set count >= replicaNum +func (ns *nodeSet) checkRackAwareWriteable(param *selectParam, nodeType NodeType) bool { + var count int + for _, rack := range ns.racks { + if contains(param.excludeRacks, rack.Rack) { + continue + } + + rnodes := rack.metaNodes + if nodeType == DataNodeType { + rnodes = rack.dataNodes + } + + rnodes.Range(func(key, value interface{}) bool { + if ns.checkNodeWriteable(value, nodeType, param) { + count++ + return false + } + return true + }) + + if count >= param.replicaNum { + return true + } + } + return false +} + +// Normal mode consistency check +func (ns *nodeSet) checkNormalWriteable(nodes *sync.Map, param *selectParam, nodeType NodeType) bool { var count int nodes.Range(func(key, value interface{}) bool { - if nodeType == RocksdbType { - node := value.(*MetaNode) - if node.IsRocksdbWriteAble() && node.PartitionCntLimited() { - count++ + if ns.checkNodeWriteable(value, nodeType, param) { + count++ + if count >= param.replicaNum { + return false // Found enough nodes, exit early } - } else { - node := value.(Node) - if node.IsWriteAble() && node.PartitionCntLimited() { - count++ - } - } - if count >= replicaNum { - return false } return true }) + log.LogInfof("canWriteForMetaNode zone[%v], ns[%v],count[%v] replicaNum[%v]", - ns.zoneName, ns.ID, count, replicaNum) - return count >= replicaNum + ns.zoneName, ns.ID, count, param.replicaNum) + return count >= param.replicaNum } func (ns *nodeSet) calcNodesForAlloc(nodes *sync.Map) (cnt int) { @@ -1846,7 +1996,7 @@ func (zone *Zone) deleteMetaNode(metaNode *MetaNode) (err error) { return } -func (zone *Zone) allocNodeSetForDataNode(excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) { +func (zone *Zone) allocNodeSetForDataNode(param *selectParam) (ns *nodeSet, err error) { nset := zone.getAllNodeSet() if nset == nil { return nil, errors.NewError(proto.ErrNoNodeSetToCreateDataPartition) @@ -1858,16 +2008,38 @@ func (zone *Zone) allocNodeSetForDataNode(excludeNodeSets []uint64, replicaNum u zone.dataNodesetSelectorLock.RLock() defer zone.dataNodesetSelectorLock.RUnlock() - ns, err = zone.dataNodesetSelector.Select(nset, excludeNodeSets, replicaNum) - if err != nil { - log.LogErrorf("action[allocNodeSetForDataNode],nset len[%v],excludeNodeSets[%v],rNum[%v] err:%v", - nset.Len(), excludeNodeSets, replicaNum, proto.ErrNoNodeSetToCreateDataPartition) + paramCopy := param.copy() + if param.rackLevel == proto.RackAwareWeak { + paramCopy.rackLevel = proto.RackAwareStrong + } + + // First attempt with strong rack awareness + ns, err = zone.dataNodesetSelector.Select(nset, paramCopy) + if err == nil { + return ns, nil + } + + if param.rackLevel != proto.RackAwareWeak { + log.LogErrorf("action[allocNodeSetForDataNode], no available node set found, param[%v] err:%v", + param.String(), err.Error()) return nil, errors.NewError(proto.ErrNoNodeSetToCreateDataPartition) } + + log.LogWarnf("action[allocNodeSetForDataNode], no available node set found, try awake, param[%v] err:%v", + param.String(), err.Error()) + + paramCopy.rackLevel = proto.RackAwareWeak + ns, err = zone.dataNodesetSelector.Select(nset, paramCopy) + if err != nil { + log.LogErrorf("action[allocNodeSetForDataNode], still failed after change rack level,param[%v], err: %v", + param.String(), proto.ErrNoNodeSetToCreateDataPartition) + return nil, errors.NewError(proto.ErrNoNodeSetToCreateDataPartition) + } + return ns, nil } -func (zone *Zone) allocNodeSetForMetaNode(excludeNodeSets []uint64, replicaNum uint8, storeMode proto.StoreMode) (ns *nodeSet, err error) { +func (zone *Zone) allocNodeSetForMetaNode(param *selectParam, storeMode proto.StoreMode) (ns *nodeSet, err error) { nset := zone.getAllNodeSet() if nset == nil { return nil, proto.ErrNoNodeSetToCreateMetaPartition @@ -1878,16 +2050,42 @@ func (zone *Zone) allocNodeSetForMetaNode(excludeNodeSets []uint64, replicaNum u // we need a read lock to block the modify of nodeset selector zone.metaNodesetSelectorLock.RLock() defer zone.metaNodesetSelectorLock.RUnlock() + + // Select appropriate selector based on store mode selector := zone.metaMemoryNodesetSelector if storeMode == proto.StoreModeRocksDb { selector = zone.metaRocksdbNodesetSelector } - ns, err = selector.Select(nset, excludeNodeSets, replicaNum) + + paramCopy := param.copy() + if param.rackLevel == proto.RackAwareWeak { + paramCopy.rackLevel = proto.RackAwareStrong + } + + // First attempt with strong rack awareness + ns, err = selector.Select(nset, paramCopy) + if err == nil { + return ns, nil + } + + if param.rackLevel != proto.RackAwareWeak { + log.LogErrorf("action[allocNodeSetForMetaNode], no available node set found, zone %s, param[%v] err:%v", + zone.name, param.String(), err.Error()) + return nil, errors.NewError(proto.ErrNoNodeSetToCreateMetaPartition) + } + + // Retry with weak rack awareness + log.LogWarnf("action[allocNodeSetForMetaNode], no available node set found, zone %s, param[%v] err:%v", + zone.name, param.String(), err.Error()) + + paramCopy.rackLevel = proto.RackAwareWeak + ns, err = selector.Select(nset, paramCopy) if err != nil { - log.LogError(fmt.Sprintf("action[allocNodeSetForMetaNode],zone[%v],excludeNodeSets[%v],rNum[%v],err:%v", - zone.name, excludeNodeSets, replicaNum, proto.ErrNoNodeSetToCreateMetaPartition)) + log.LogErrorf("action[allocNodeSetForMetaNode], still failed after change rack level, zone[%v], param[%v], err:%v", + zone.name, param.String(), proto.ErrNoNodeSetToCreateMetaPartition) return nil, proto.ErrNoNodeSetToCreateMetaPartition } + return ns, nil } @@ -1983,19 +2181,19 @@ func (zone *Zone) getSpaceLeft(dataType uint32) (spaceLeft uint64) { } } -func (zone *Zone) getAvailNodeHosts(nodeType uint32, excludeNodeSets []uint64, excludeHosts []string, replicaNum int) (newHosts []string, peers []proto.Peer, err error) { - if replicaNum == 0 { +func (zone *Zone) getAvailNodeHosts(nodeType uint32, param *selectParam) (newHosts []string, peers []proto.Peer, err error) { + if param.replicaNum == 0 { return } log.LogDebugf("[getAvailNodeHosts] get node host, zone(%s), nodeType(%d)", zone.name, nodeType) if nodeType == TypeDataPartition { - ns, err := zone.allocNodeSetForDataNode(excludeNodeSets, uint8(replicaNum)) + ns, err := zone.allocNodeSetForDataNode(param) if err != nil { - return nil, nil, errors.Trace(err, "zone[%v] alloc node set,replicaNum[%v]", zone.name, replicaNum) + return nil, nil, errors.Trace(err, "zone[%v] alloc node set, param[%v], err %s", zone.name, param.String(), err.Error()) } - return ns.getAvailDataNodeHosts(excludeHosts, replicaNum) + return ns.getAvailDataNodeHosts(param) } storeMode := proto.StoreModeMem @@ -2003,12 +2201,12 @@ func (zone *Zone) getAvailNodeHosts(nodeType uint32, excludeNodeSets []uint64, e storeMode = proto.StoreModeRocksDb } - ns, err := zone.allocNodeSetForMetaNode(excludeNodeSets, uint8(replicaNum), storeMode) + ns, err := zone.allocNodeSetForMetaNode(param, storeMode) if err != nil { - return nil, nil, errors.NewErrorf("zone[%v],err[%v]", zone.name, err) + return nil, nil, errors.NewErrorf("zone[%v], param[%v], err[%v]", zone.name, param.String(), err.Error()) } - return ns.getAvailMetaNodeHosts(excludeHosts, replicaNum, storeMode) + return ns.getAvailMetaNodeHosts(param, storeMode) } func (zone *Zone) updateNodesetSelector(cluster *Cluster, dataNodesetSelector string, metaNodesetSelector string) error { diff --git a/master/topology_rack_test.go b/master/topology_rack_test.go new file mode 100644 index 000000000..25939b531 --- /dev/null +++ b/master/topology_rack_test.go @@ -0,0 +1,1761 @@ +package master + +import ( + "strconv" + "sync" + "testing" + "time" + + "github.com/cubefs/cubefs/proto" + "github.com/cubefs/cubefs/raftstore" + "github.com/cubefs/cubefs/util" + "github.com/stretchr/testify/require" +) + +// Helper function: Create data node with rack information +func createDataNodeWithRack(addr, zoneName, rackName string, ns *nodeSet) *DataNode { + dn := newDataNode(addr, strconv.Itoa(raftstore.DefaultHeartbeatPort), strconv.Itoa(raftstore.DefaultReplicaPort), zoneName, "", "test", proto.MediaType_HDD) + dn.ZoneName = zoneName + dn.Rack = rackName + dn.Total = 1024 * util.GB + dn.Used = 10 * util.GB + dn.AvailableSpace = 1024 * util.GB + dn.ReportTime = time.Now() + dn.isActive = true + dn.NodeSetID = ns.ID + dn.AllDisks = []string{"/cfs/disk"} + dn.DpCntLimit = defaultMaxDpCntLimit + return dn +} + +// Helper function: Create meta node with rack information +func createMetaNodeWithRack(addr, zoneName, rackName string, ns *nodeSet) *MetaNode { + mn := newMetaNode(addr, strconv.Itoa(raftstore.DefaultHeartbeatPort), strconv.Itoa(raftstore.DefaultReplicaPort), zoneName, "", "test") + mn.ZoneName = zoneName + mn.Rack = rackName + mn.Total = 1024 * util.GB + mn.Used = 10 * util.GB + mn.ReportTime = time.Now() + mn.IsActive = true + mn.NodeSetID = ns.ID + mn.Threshold = 0.8 + mn.MaxMemAvailWeight = 1024 * util.GB + return mn +} + +// Helper function: Setup test environment for rack-aware testing with multiple nodes +func setupRackAwareTestEnvWithNodes(t *testing.T, rackCount, nodesPerRack int) (*topology, *Cluster, *Zone, *nodeSet) { + topo := newTopology() + c := new(Cluster) + c.cfg = newClusterConfig() + c.cfg.RackAwareLevel = proto.RackAwareStrong + + zoneName := "test-zone" + zone := newZone(zoneName, proto.MediaType_Unspecified) + topo.putZone(zone) + + ns := newNodeSet(c, 1, rackCount*nodesPerRack+10, zoneName, "") // Set capacity higher than total nodes + zone.putNodeSet(ns) + + // Add nodes in multiple racks + for rackIdx := 0; rackIdx < rackCount; rackIdx++ { + rackName := "rack" + strconv.Itoa(rackIdx+1) + for nodeIdx := 0; nodeIdx < nodesPerRack; nodeIdx++ { + dn := createDataNodeWithRack( + "192.168.1."+strconv.Itoa(rackIdx*10+nodeIdx+1)+":17310", + zone.name, + rackName, + ns, + ) + ns.putDataNode(dn) + } + } + + return topo, c, zone, ns +} + +// Helper function: Setup test environment for rack-aware testing +func setupRackAwareTestEnv(t *testing.T) (*topology, *Cluster, *Zone, *nodeSet) { + return setupRackAwareTestEnvWithNodes(t, 6, 2) // 6 racks, 2 nodes per rack = 12 total nodes +} + +// Test 1: Test the correctness of checkRackAwareWriteable function with sufficient nodes +func TestRackCheckRackAwareWriteable(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnv(t) // 6 racks, 2 nodes per rack = 12 nodes + + // Test case 1: Sufficient racks to meet replica requirement + param := &selectParam{ + replicaNum: 3, + rackLevel: proto.RackAwareStrong, + } + + // Should return true since we have 6 racks and need 3 replicas + result := ns.checkRackAwareWriteable(param, DataNodeType) + require.True(t, result, "Should have enough racks for replica requirement") + + // Test case 2: Insufficient racks + param.replicaNum = 8 + result = ns.checkRackAwareWriteable(param, DataNodeType) + require.False(t, result, "Should not have enough racks for 8 replicas") + + // Test case 3: Exclude certain racks + param.replicaNum = 3 + param.excludeRacks = []string{"rack1", "rack2"} + result = ns.checkRackAwareWriteable(param, DataNodeType) + require.True(t, result, "Should have enough racks after excluding 2 racks") + + // Test case 4: Exclude too many racks + param.excludeRacks = []string{"rack1", "rack2", "rack3", "rack4"} + result = ns.checkRackAwareWriteable(param, DataNodeType) + require.False(t, result, "Should not have enough racks after excluding 4 racks") + + // Test case 5: Test with maximum replica number + param.replicaNum = 6 + param.excludeRacks = nil + result = ns.checkRackAwareWriteable(param, DataNodeType) + require.True(t, result, "Should have enough racks for maximum replica number") +} + +// Test 2: Test rack selection logic in selectNodesWithRack function with multiple nodes +func TestRackSelectNodesWithRack(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnv(t) // 6 racks, 2 nodes per rack = 12 nodes + + // Test case 1: Strong rack awareness mode with multiple replicas + param := &selectParam{ + replicaNum: 4, + rackLevel: proto.RackAwareStrong, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select nodes with rack awareness") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Verify selected nodes are from different racks + selectedRacks := make(map[string]bool) + for _, host := range hosts { + // Find corresponding rack based on host + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack] = true + return false + } + return true + }) + } + require.Equal(t, 4, len(selectedRacks), "Selected nodes should be from 4 different racks") + + // Test case 2: Weak rack awareness mode + param.rackLevel = proto.RackAwareWeak + param.replicaNum = 6 + hosts, _, err = ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select nodes with weak rack awareness") + require.Equal(t, 6, len(hosts), "Should select 6 hosts") + + // Test case 3: Test with maximum possible replicas + param.rackLevel = proto.RackAwareStrong + param.replicaNum = 6 + hosts, _, err = ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select maximum replicas with strong rack awareness") + require.Equal(t, 6, len(hosts), "Should select 6 hosts from 6 different racks") +} + +// Test 3: Test handling when no racks are available +func TestRackSelectNodesWithRackNoRacks(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnv(t) + + // Clear all nodes to test empty rack collection scenario + ns.dataNodes = new(sync.Map) + ns.racks = make(map[string]*nodeSet) + + param := &selectParam{ + replicaNum: 2, + rackLevel: proto.RackAwareStrong, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.Error(t, err, "Should return error when no racks available") + require.Nil(t, hosts, "Should return nil hosts") +} + +// Test 4: Test rack awareness degradation mechanism with multiple nodes per rack +func TestRackAwareDegradation(t *testing.T) { + // Setup with fewer racks but more nodes per rack + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 2, 4) // 2 racks, 4 nodes per rack = 8 nodes + + // Test strong rack awareness mode, should fail since we need 3 replicas but only have 2 racks + param := &selectParam{ + replicaNum: 3, // Need 3 replicas but only 2 racks + rackLevel: proto.RackAwareStrong, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + // Should fail in strong awareness mode since we don't have enough racks + require.Error(t, err, "Should fail with strong rack awareness when not enough racks") + require.Nil(t, hosts, "Should return nil hosts when failing") + + // Test weak rack awareness mode, should succeed since same rack has multiple nodes + param.rackLevel = proto.RackAwareWeak + hosts, _, err = ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should succeed with weak rack awareness") + require.Equal(t, 3, len(hosts), "Should select 3 hosts") + + // Verify that we can select from same rack when needed in weak mode + selectedRacks := make(map[string]int) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + // Should have nodes from both racks, with at least one rack having multiple nodes + require.Equal(t, 2, len(selectedRacks), "Should use both available racks") + require.GreaterOrEqual(t, selectedRacks["rack1"]+selectedRacks["rack2"], 3, "Should select 3 total nodes") +} + +// Test 5: Test concurrency safety of getRackSets function with many nodes +func TestRackGetRackSetsConcurrency(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 8, 3) // 8 racks, 3 nodes per rack = 24 nodes + + // Concurrency test + done := make(chan bool, 20) + for i := 0; i < 20; i++ { + go func() { + defer func() { + done <- true + }() + + // Concurrent calls to getRackSets + rsets := ns.getRackSets() + require.NotNil(t, rsets, "getRackSets should not return nil") + require.Equal(t, 8, len(rsets), "Should return 8 rack sets") + require.GreaterOrEqual(t, len(rsets), 0, "Rack sets length should be non-negative") + }() + } + + // Wait for all goroutines to complete + for i := 0; i < 20; i++ { + <-done + } +} + +// Test 6: Test parameter override issue in allocNodeSetForDataNode with multiple nodes +func TestRackAllocNodeSetForDataNodeParamOverride(t *testing.T) { + _, _, zone, _ := setupRackAwareTestEnv(t) // 6 racks, 2 nodes per rack = 12 nodes + + // Test case 1: User explicitly sets rackLevel to None + param := &selectParam{ + replicaNum: 3, + rackLevel: proto.RackAwareNone, + } + + selectedNodeSet, err := zone.allocNodeSetForDataNode(param) + require.NoError(t, err, "Should successfully allocate node set") + require.NotNil(t, selectedNodeSet, "Selected node set should not be nil") + + // Test case 2: User sets rackLevel to Weak + param.rackLevel = proto.RackAwareWeak + selectedNodeSet, err = zone.allocNodeSetForDataNode(param) + require.NoError(t, err, "Should successfully allocate node set with weak rack awareness") + + // Test case 3: Test with higher replica count + param.replicaNum = 5 + param.rackLevel = proto.RackAwareStrong + selectedNodeSet, err = zone.allocNodeSetForDataNode(param) + require.NoError(t, err, "Should successfully allocate node set with high replica count") +} + +// Test 7: Test deep copy functionality +func TestRackSelectParamDeepCopy(t *testing.T) { + original := &selectParam{ + excludeHosts: []string{"host1", "host2", "host3", "host4"}, + replicaNum: 5, + rackLevel: proto.RackAwareStrong, + excludeRacks: []string{"rack1", "rack2", "rack3"}, + excludeNodeSets: []uint64{1, 2, 3, 4, 5}, + } + + copied := original.copy() + + // Modify original parameters + original.excludeHosts[0] = "modified" + original.excludeRacks = append(original.excludeRacks, "rack4") + original.excludeNodeSets[0] = 999 + + // Verify copied parameters are not modified + require.Equal(t, "host1", copied.excludeHosts[0], "Copied excludeHosts should not be modified") + require.Equal(t, 3, len(copied.excludeRacks), "Copied excludeRacks length should not change") + require.Equal(t, uint64(1), copied.excludeNodeSets[0], "Copied excludeNodeSets should not be modified") + require.Equal(t, 5, copied.replicaNum, "Copied replicaNum should not be modified") + require.Equal(t, proto.RackAwareStrong, copied.rackLevel, "Copied rackLevel should not be modified") +} + +// Test 8: Test boundary conditions for rack awareness with many nodes +func TestRackBoundaryConditions(t *testing.T) { + _, _, zone, ns := setupRackAwareTestEnv(t) // 6 racks, 2 nodes per rack = 12 nodes + + // Test case 1: replicaNum is 0 + param := &selectParam{ + replicaNum: 0, + rackLevel: proto.RackAwareStrong, + } + + hosts, _, err := zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Should handle replicaNum=0 gracefully") + require.Equal(t, 0, len(hosts), "Should return empty hosts for replicaNum=0") + + // Test case 2: All nodes are not writable + // Set all nodes as inactive + ns.dataNodes.Range(func(key, value interface{}) bool { + dn := value.(*DataNode) + dn.isActive = false + return true + }) + + param.replicaNum = 1 + hosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param) + require.Error(t, err, "Should return error when no writable nodes available") + + // Test case 3: Test with maximum replica number for strong rack awareness + // Reset nodes to active + ns.dataNodes.Range(func(key, value interface{}) bool { + dn := value.(*DataNode) + dn.isActive = true + return true + }) + + // In RackAwareStrong mode, can only select up to the number of racks (6) + param.replicaNum = 6 // Maximum possible with 6 racks in strong mode + hosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Should handle maximum replica number for strong rack awareness") + require.Equal(t, 6, len(hosts), "Should return 6 hosts (one per rack)") + + // Test case 4: Test with weak rack awareness to select more nodes + param.rackLevel = proto.RackAwareWeak + param.replicaNum = 12 // All available nodes + hosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Should handle maximum replica number with weak rack awareness") + require.Equal(t, 12, len(hosts), "Should return all available hosts with weak rack awareness") +} + +// Test 9: Test rack awareness for meta nodes with multiple nodes +func TestRackMetaNodeRackAwareness(t *testing.T) { + _, _, zone, ns := setupRackAwareTestEnvWithNodes(t, 5, 3) // 5 racks, 3 nodes per rack = 15 nodes + + // Add meta nodes in multiple racks + for rackIdx := 0; rackIdx < 5; rackIdx++ { + rackName := "rack" + strconv.Itoa(rackIdx+1) + for nodeIdx := 0; nodeIdx < 3; nodeIdx++ { + mn := createMetaNodeWithRack( + "192.168.2."+strconv.Itoa(rackIdx*10+nodeIdx+1)+":17310", + zone.name, + rackName, + ns, + ) + ns.putMetaNode(mn) + } + } + + param := &selectParam{ + replicaNum: 4, + rackLevel: proto.RackAwareStrong, + } + + hosts, _, err := ns.getAvailMetaNodeHosts(param, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select meta nodes with rack awareness") + require.Equal(t, 4, len(hosts), "Should select 4 meta node hosts") + + // Verify selected nodes are from different racks + selectedRacks := make(map[string]bool) + for _, host := range hosts { + ns.metaNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + mn := value.(*MetaNode) + selectedRacks[mn.Rack] = true + return false + } + return true + }) + } + require.Equal(t, 4, len(selectedRacks), "Selected meta nodes should be from 4 different racks") +} + +// Test 10: Test performance and stability of rack awareness with many nodes +func TestRackPerformance(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 10, 8) // 10 racks, 8 nodes per rack = 80 nodes + + param := &selectParam{ + replicaNum: 5, + rackLevel: proto.RackAwareStrong, + } + + // Execute selection multiple times to test stability + for i := 0; i < 20; i++ { + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should consistently succeed with large number of racks") + require.Equal(t, 5, len(hosts), "Should consistently select 5 hosts") + + // Verify rack distribution + selectedRacks := make(map[string]bool) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack] = true + return false + } + return true + }) + } + require.Equal(t, 5, len(selectedRacks), "Should consistently select from 5 different racks") + } +} + +// Test 11: Test error recovery mechanism for rack awareness with mixed availability +func TestRackErrorRecovery(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 6, 4) // 6 racks, 4 nodes per rack = 24 nodes + + // Set some nodes as unavailable in different racks + nodeCount := 0 + ns.dataNodes.Range(func(key, value interface{}) bool { + dn := value.(*DataNode) + if nodeCount%3 == 0 { // Make every 3rd node unavailable + dn.isActive = false + } + nodeCount++ + return true + }) + + param := &selectParam{ + replicaNum: 4, + rackLevel: proto.RackAwareStrong, + } + + // Should still succeed since we have enough available nodes + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should succeed with mixed node availability") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Test weak awareness mode with more replicas + param.rackLevel = proto.RackAwareWeak + param.replicaNum = 6 + hosts, _, err = ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should succeed with weak awareness and mixed availability") + require.Equal(t, 6, len(hosts), "Should select 6 hosts") +} + +// Test 12: Test rack awareness configuration changes with many nodes +func TestRackConfigChange(t *testing.T) { + _, c, _, ns := setupRackAwareTestEnv(t) // 6 racks, 2 nodes per rack = 12 nodes + + // Test different configuration levels + configs := []proto.RackAwareLevel{ + proto.RackAwareNone, + proto.RackAwareWeak, + proto.RackAwareStrong, + } + + for _, config := range configs { + c.cfg.RackAwareLevel = config + + param := &selectParam{ + replicaNum: 3, + rackLevel: config, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should work with config level %v", config) + require.Equal(t, 3, len(hosts), "Should select 3 hosts with config level %v", config) + } +} + +// Test 13: Test rack exclusion functionality with many nodes +func TestRackExclusion(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 8, 3) // 8 racks, 3 nodes per rack = 24 nodes + + // Test excluding specific racks + param := &selectParam{ + replicaNum: 4, + rackLevel: proto.RackAwareStrong, + excludeRacks: []string{"rack1", "rack2", "rack3"}, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select nodes excluding specified racks") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Verify selected nodes are not from excluded racks + selectedRacks := make(map[string]bool) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack] = true + require.NotContains(t, param.excludeRacks, dn.Rack, "Selected node should not be from excluded rack") + return false + } + return true + }) + } + require.Equal(t, 4, len(selectedRacks), "Should select from 4 different non-excluded racks") + + // Test excluding too many racks + param.excludeRacks = []string{"rack1", "rack2", "rack3", "rack4", "rack5", "rack6"} + param.replicaNum = 3 + hosts, _, err = ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.Error(t, err, "Should fail when excluding too many racks") +} + +// Helper function: Print nodeset topology information +func printNodeSetTopology(t *testing.T, ns *nodeSet, title string) { + t.Logf("=== %s ===", title) + t.Logf("NodeSet ID: %d", ns.ID) + t.Logf("NodeSet Zone: %s", ns.zoneName) + t.Logf("NodeSet Capacity: %d", ns.Capacity) + t.Logf("NodeSet Rack: %s", ns.Rack) + + // Print data nodes information + dataNodeCount := 0 + t.Logf("Data Nodes:") + ns.dataNodes.Range(func(key, value interface{}) bool { + dn := value.(*DataNode) + dataNodeCount++ + t.Logf(" - Addr: %s, Rack: %s, Zone: %s, Active: %v, NodeSetID: %d", + dn.Addr, dn.Rack, dn.ZoneName, dn.isActive, dn.NodeSetID) + return true + }) + t.Logf("Total Data Nodes: %d", dataNodeCount) + + // Print meta nodes information + metaNodeCount := 0 + t.Logf("Meta Nodes:") + ns.metaNodes.Range(func(key, value interface{}) bool { + mn := value.(*MetaNode) + metaNodeCount++ + t.Logf(" - Addr: %s, Rack: %s, Zone: %s, Active: %v, NodeSetID: %d", + mn.Addr, mn.Rack, mn.ZoneName, mn.IsActive, mn.NodeSetID) + return true + }) + t.Logf("Total Meta Nodes: %d", metaNodeCount) + + // Print rack information + t.Logf("Racks:") + for rackName, rack := range ns.racks { + rackDataCount := 0 + rackMetaCount := 0 + + rack.dataNodes.Range(func(key, value interface{}) bool { + rackDataCount++ + return true + }) + + rack.metaNodes.Range(func(key, value interface{}) bool { + rackMetaCount++ + return true + }) + + t.Logf(" - Rack: %s, DataNodes: %d, MetaNodes: %d", + rackName, rackDataCount, rackMetaCount) + } + t.Logf("Total Racks: %d", len(ns.racks)) + t.Logf("========================") +} + +// Test 14: Test rack awareness with mixed node types and many nodes +func TestRackMixedNodeTypes(t *testing.T) { + _, _, zone, ns := setupRackAwareTestEnvWithNodes(t, 5, 4) // 5 racks, 4 nodes per rack = 20 nodes + + // Print initial topology information + printNodeSetTopology(t, ns, "Initial NodeSet Topology Information") + + // Add meta nodes in same racks as data nodes + for rackIdx := 0; rackIdx < 5; rackIdx++ { + rackName := "rack" + strconv.Itoa(rackIdx+1) + for nodeIdx := 0; nodeIdx < 4; nodeIdx++ { + // Add meta node + mn := createMetaNodeWithRack( + "192.168.2."+strconv.Itoa(rackIdx*10+nodeIdx+1)+":17310", + zone.name, + rackName, + ns, + ) + ns.putMetaNode(mn) + } + } + + // Print topology information after adding meta nodes + printNodeSetTopology(t, ns, "After Adding Meta Nodes") + + // Test data node selection + param := &selectParam{ + replicaNum: 4, + rackLevel: proto.RackAwareStrong, + } + + dataHosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select data nodes with rack awareness") + require.Equal(t, 4, len(dataHosts), "Should select 4 data node hosts") + + // Test meta node selection + metaHosts, _, err := ns.selectNodesWithRack(param, MetaNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select meta nodes with rack awareness") + require.Equal(t, 4, len(metaHosts), "Should select 4 meta node hosts") + + // Verify both selections use different racks + dataRacks := make(map[string]bool) + metaRacks := make(map[string]bool) + + // Get data node racks + ns.dataNodes.Range(func(key, value interface{}) bool { + dn := value.(*DataNode) + for _, host := range dataHosts { + if key.(string) == host { + dataRacks[dn.Rack] = true + break + } + } + return true + }) + + // Get meta node racks + ns.metaNodes.Range(func(key, value interface{}) bool { + mn := value.(*MetaNode) + for _, host := range metaHosts { + if key.(string) == host { + metaRacks[mn.Rack] = true + break + } + } + return true + }) + + require.Equal(t, 4, len(dataRacks), "Data nodes should be from 4 different racks") + require.Equal(t, 4, len(metaRacks), "Meta nodes should be from 4 different racks") +} + +// Test 16: Test rack awareness with host exclusion and many nodes +func TestRackWithHostExclusion(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 6, 4) // 6 racks, 4 nodes per rack = 24 nodes + + // Test with host exclusion + excludeHosts := []string{ + "192.168.1.1:17310", "192.168.1.2:17310", "192.168.1.3:17310", // Exclude 3 hosts from rack1 + "192.168.1.11:17310", "192.168.1.12:17310", // Exclude 2 hosts from rack2 + } + + param := &selectParam{ + replicaNum: 4, + rackLevel: proto.RackAwareStrong, + excludeHosts: excludeHosts, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select nodes with host exclusion") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Verify excluded hosts are not selected + for _, host := range hosts { + require.NotContains(t, param.excludeHosts, host, "Selected host should not be in exclude list") + } + + // Verify rack distribution is maintained + selectedRacks := make(map[string]bool) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack] = true + return false + } + return true + }) + } + require.Equal(t, 4, len(selectedRacks), "Should still select from 4 different racks despite exclusions") +} + +// Test 17: Test rack awareness with node set exclusion +func TestRackWithNodeSetExclusion(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnv(t) // 6 racks, 2 nodes per rack = 12 nodes + + // Test with node set exclusion + param := &selectParam{ + replicaNum: 3, + rackLevel: proto.RackAwareStrong, + excludeNodeSets: []uint64{ns.ID}, + } + + // This should fail since we're excluding the only node set + _, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.Error(t, err, "Should fail when excluding the only available node set") +} + +// Test 18: Test rack awareness stress test with concurrent operations and many nodes +func TestRackStressTest(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 12, 6) // 12 racks, 6 nodes per rack = 72 nodes + + param := &selectParam{ + replicaNum: 6, + rackLevel: proto.RackAwareStrong, + } + + // Concurrent stress test + done := make(chan bool, 10) + for i := 0; i < 10; i++ { + go func() { + defer func() { + done <- true + }() + + for j := 0; j < 10; j++ { + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should consistently succeed under stress") + require.Equal(t, 6, len(hosts), "Should consistently select 6 hosts under stress") + + // Verify rack distribution + selectedRacks := make(map[string]bool) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack] = true + return false + } + return true + }) + } + require.Equal(t, 6, len(selectedRacks), "Should consistently select from 6 different racks under stress") + } + }() + } + + // Wait for all goroutines to complete + for i := 0; i < 10; i++ { + <-done + } +} + +// Test 19: Test rack awareness with maximum replica scenarios +func TestRackMaxReplicas(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 8, 3) // 8 racks, 3 nodes per rack = 24 nodes + + // Test with maximum possible replicas (all racks) + param := &selectParam{ + replicaNum: 8, + rackLevel: proto.RackAwareStrong, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select maximum replicas") + require.Equal(t, 8, len(hosts), "Should select 8 hosts from 8 different racks") + + // Verify all racks are used + selectedRacks := make(map[string]bool) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack] = true + return false + } + return true + }) + } + require.Equal(t, 8, len(selectedRacks), "Should use all 8 available racks") +} + +// Test 20: Test rack awareness with complex exclusion scenarios +func TestRackComplexExclusions(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 10, 4) // 10 racks, 4 nodes per rack = 40 nodes + + // Complex exclusion scenario: exclude some racks and some hosts + param := &selectParam{ + replicaNum: 5, + rackLevel: proto.RackAwareStrong, + excludeRacks: []string{"rack1", "rack2", "rack3"}, // Exclude 3 racks + excludeHosts: []string{ + "192.168.1.31:17310", "192.168.1.32:17310", // Exclude 2 hosts from rack4 + "192.168.1.41:17310", // Exclude 1 host from rack5 + }, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select nodes with complex exclusions") + require.Equal(t, 5, len(hosts), "Should select 5 hosts") + + // Verify exclusions are respected + selectedRacks := make(map[string]bool) + for _, host := range hosts { + require.NotContains(t, param.excludeHosts, host, "Selected host should not be in exclude list") + + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + require.NotContains(t, param.excludeRacks, dn.Rack, "Selected node should not be from excluded rack") + selectedRacks[dn.Rack] = true + return false + } + return true + }) + } + require.Equal(t, 5, len(selectedRacks), "Should select from 5 different non-excluded racks") +} + +// Test 21: Test excludeHosts edge cases and boundary conditions +func TestRackExcludeHostsEdgeCases(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 4, 2) // 4 racks, 2 nodes per rack = 8 nodes + + // Test case 1: Exclude all hosts from one rack + param := &selectParam{ + replicaNum: 3, + rackLevel: proto.RackAwareStrong, + excludeHosts: []string{"192.168.1.1:17310", "192.168.1.2:17310"}, // All hosts from rack1 + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should succeed when excluding all hosts from one rack") + require.Equal(t, 3, len(hosts), "Should select 3 hosts from remaining racks") + + // Test case 2: Exclude hosts from multiple racks but leave enough + param.excludeHosts = []string{ + "192.168.1.1:17310", // 1 host from rack1 + "192.168.1.11:17310", // 1 host from rack2 + "192.168.1.21:17310", // 1 host from rack3 + } + + hosts, _, err = ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should succeed with selective host exclusions") + require.Equal(t, 3, len(hosts), "Should select 3 hosts") + + // Test case 3: Exclude too many hosts to make selection impossible + param.excludeHosts = []string{ + "192.168.1.1:17310", "192.168.1.2:17310", // All hosts from rack1 + "192.168.1.11:17310", "192.168.1.12:17310", // All hosts from rack2 + "192.168.1.21:17310", "192.168.1.22:17310", // All hosts from rack3 + } + param.replicaNum = 3 // Need 3 replicas but only 1 rack available + + hosts, _, err = ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.Error(t, err, "Should fail when excluding too many hosts") +} + +// Test 22: Test canWriteFor method with excludeHosts - CRITICAL BUG FIX TEST +func TestRackCanWriteForWithExcludeHosts(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 6, 3) // 6 racks, 3 nodes per rack = 18 nodes + + // Test case 1: Normal case without exclusions + param := &selectParam{ + replicaNum: 3, + rackLevel: proto.RackAwareStrong, + } + + canWrite := ns.canWriteFor(DataNodeType, param) + require.True(t, canWrite, "NodeSet should be writable without exclusions") + + // Test case 2: With host exclusions but still enough nodes + param.excludeHosts = []string{ + "192.168.1.1:17310", "192.168.1.2:17310", // 2 hosts from rack1 + "192.168.1.11:17310", // 1 host from rack2 + } + + canWrite = ns.canWriteFor(DataNodeType, param) + require.True(t, canWrite, "NodeSet should be writable with moderate host exclusions") + + // Test case 3: With too many host exclusions + param.excludeHosts = []string{ + "192.168.1.1:17310", "192.168.1.2:17310", "192.168.1.3:17310", // All hosts from rack1 + "192.168.1.11:17310", "192.168.1.12:17310", "192.168.1.13:17310", // All hosts from rack2 + "192.168.1.21:17310", "192.168.1.22:17310", "192.168.1.23:17310", // All hosts from rack3 + "192.168.1.31:17310", "192.168.1.32:17310", "192.168.1.33:17310", // All hosts from rack4 + } + + canWrite = ns.canWriteFor(DataNodeType, param) + require.False(t, canWrite, "NodeSet should not be writable with too many host exclusions") + + // Test case 4: Test with rack exclusions + param.excludeHosts = nil + param.excludeRacks = []string{"rack1", "rack2", "rack3"} + + canWrite = ns.canWriteFor(DataNodeType, param) + require.True(t, canWrite, "NodeSet should be writable with rack exclusions") + + // Test case 5: Test with both host and rack exclusions + param.excludeHosts = []string{"192.168.1.41:17310", "192.168.1.42:17310"} // 2 hosts from rack5 + param.excludeRacks = []string{"rack1", "rack2"} // Exclude 2 racks + + canWrite = ns.canWriteFor(DataNodeType, param) + require.True(t, canWrite, "NodeSet should be writable with combined exclusions") +} + +// Test 23: Test checkNodeWriteable method with excludeHosts +func TestRackCheckNodeWriteableWithExcludeHosts(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 2, 2) // 2 racks, 2 nodes per rack = 4 nodes + + param := &selectParam{ + excludeHosts: []string{"192.168.1.1:17310"}, + } + + // Test with excluded host + var excludedNode *DataNode + ns.dataNodes.Range(func(key, value interface{}) bool { + dn := value.(*DataNode) + if dn.Addr == "192.168.1.1:17310" { + excludedNode = dn + return false + } + return true + }) + + require.NotNil(t, excludedNode, "Should find the excluded node") + canWrite := ns.checkNodeWriteable(excludedNode, DataNodeType, param) + require.False(t, canWrite, "Excluded node should not be writable") + + // Test with non-excluded host + var normalNode *DataNode + ns.dataNodes.Range(func(key, value interface{}) bool { + dn := value.(*DataNode) + if dn.Addr != "192.168.1.1:17310" { + normalNode = dn + return false + } + return true + }) + + require.NotNil(t, normalNode, "Should find a normal node") + canWrite = ns.checkNodeWriteable(normalNode, DataNodeType, param) + require.True(t, canWrite, "Normal node should be writable") +} + +// Test 24: Test checkRackAwareWriteable with excludeHosts +func TestRackCheckRackAwareWriteableWithExcludeHosts(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 6, 2) // 6 racks, 2 nodes per rack = 12 nodes + + // Test case 1: Exclude hosts from multiple racks + param := &selectParam{ + replicaNum: 3, + rackLevel: proto.RackAwareStrong, + excludeHosts: []string{ + "192.168.1.1:17310", // 1 host from rack1 + "192.168.1.11:17310", // 1 host from rack2 + "192.168.1.21:17310", // 1 host from rack3 + }, + } + + canWrite := ns.checkRackAwareWriteable(param, DataNodeType) + require.True(t, canWrite, "Should be writable with host exclusions from multiple racks") + + // Test case 2: Exclude all hosts from some racks + param.excludeHosts = []string{ + "192.168.1.1:17310", "192.168.1.2:17310", // All hosts from rack1 + "192.168.1.11:17310", "192.168.1.12:17310", // All hosts from rack2 + "192.168.1.21:17310", "192.168.1.22:17310", // All hosts from rack3 + } + + canWrite = ns.checkRackAwareWriteable(param, DataNodeType) + require.True(t, canWrite, "Should still be writable with 3 racks remaining") + + // Test case 3: Exclude too many hosts + param.excludeHosts = []string{ + "192.168.1.1:17310", "192.168.1.2:17310", // All hosts from rack1 + "192.168.1.11:17310", "192.168.1.12:17310", // All hosts from rack2 + "192.168.1.21:17310", "192.168.1.22:17310", // All hosts from rack3 + "192.168.1.31:17310", "192.168.1.32:17310", // All hosts from rack4 + "192.168.1.41:17310", "192.168.1.42:17310", // All hosts from rack5 + } + + canWrite = ns.checkRackAwareWriteable(param, DataNodeType) + require.False(t, canWrite, "Should not be writable with only 1 rack remaining") +} + +// Test 25: Test checkNormalWriteable with excludeHosts +func TestRackCheckNormalWriteableWithExcludeHosts(t *testing.T) { + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 2, 3) // 2 racks, 3 nodes per rack = 6 nodes + + // Test case 1: Normal case without exclusions + param := &selectParam{ + replicaNum: 3, + rackLevel: proto.RackAwareNone, + } + + canWrite := ns.checkNormalWriteable(ns.dataNodes, param, DataNodeType) + require.True(t, canWrite, "Should be writable without exclusions") + + // Test case 2: With host exclusions + param.excludeHosts = []string{ + "192.168.1.1:17310", "192.168.1.2:17310", // 2 hosts from rack1 + } + + canWrite = ns.checkNormalWriteable(ns.dataNodes, param, DataNodeType) + require.True(t, canWrite, "Should be writable with host exclusions") + + // Test case 3: With too many host exclusions + param.excludeHosts = []string{ + "192.168.1.1:17310", "192.168.1.2:17310", "192.168.1.3:17310", // All hosts from rack1 + "192.168.1.11:17310", "192.168.1.12:17310", // 2 hosts from rack2 + } + + canWrite = ns.checkNormalWriteable(ns.dataNodes, param, DataNodeType) + require.False(t, canWrite, "Should not be writable with too many host exclusions") +} + +// Test 26: Test nodeset selector with excludeHosts - INTEGRATION TEST +func TestRackNodeSetSelectorWithExcludeHosts(t *testing.T) { + _, _, zone, _ := setupRackAwareTestEnvWithNodes(t, 4, 3) // 4 racks, 3 nodes per rack = 12 nodes + + // Test case 1: Normal selection without exclusions + param := &selectParam{ + replicaNum: 3, + rackLevel: proto.RackAwareStrong, + } + + selectedNodeSet, err := zone.allocNodeSetForDataNode(param) + require.NoError(t, err, "Should successfully allocate node set") + require.NotNil(t, selectedNodeSet, "Selected node set should not be nil") + + // Test case 2: Selection with host exclusions + param.excludeHosts = []string{ + "192.168.1.1:17310", "192.168.1.2:17310", // 2 hosts from rack1 + "192.168.1.11:17310", // 1 host from rack2 + } + + selectedNodeSet, err = zone.allocNodeSetForDataNode(param) + require.NoError(t, err, "Should successfully allocate node set with host exclusions") + require.NotNil(t, selectedNodeSet, "Selected node set should not be nil") + + // Verify that the selected nodeset can handle the exclusions + hosts, _, err := selectedNodeSet.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Should successfully select nodes from allocated node set") + require.Equal(t, 3, len(hosts), "Should select 3 hosts") + + // Verify excluded hosts are not selected + for _, host := range hosts { + require.NotContains(t, param.excludeHosts, host, "Selected host should not be in exclude list") + } +} + +// Test 27: Test comprehensive excludeHosts scenarios +func TestRackComprehensiveExcludeHosts(t *testing.T) { + _, _, zone, ns := setupRackAwareTestEnvWithNodes(t, 8, 4) // 8 racks, 4 nodes per rack = 32 nodes + + // Test scenario 1: Gradual exclusion of hosts + param := &selectParam{ + replicaNum: 4, + rackLevel: proto.RackAwareStrong, + } + + // Start with no exclusions + hosts, _, err := zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Should succeed with no exclusions") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Exclude some hosts gradually + param.excludeHosts = []string{hosts[0]} // Exclude first selected host + hosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Should succeed after excluding one host") + require.Equal(t, 4, len(hosts), "Should still select 4 hosts") + require.NotContains(t, hosts, param.excludeHosts[0], "Excluded host should not be selected") + + // Exclude more hosts + param.excludeHosts = append(param.excludeHosts, hosts[0], hosts[1]) // Exclude 3 hosts total + hosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Should succeed after excluding more hosts") + require.Equal(t, 4, len(hosts), "Should still select 4 hosts") + + // Test scenario 2: Exclude hosts from specific racks + param.excludeHosts = []string{ + "192.168.1.1:17310", "192.168.1.2:17310", "192.168.1.3:17310", "192.168.1.4:17310", // All hosts from rack1 + "192.168.1.11:17310", "192.168.1.12:17310", "192.168.1.13:17310", "192.168.1.14:17310", // All hosts from rack2 + } + + hosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Should succeed after excluding hosts from 2 racks") + require.Equal(t, 4, len(hosts), "Should select 4 hosts from remaining racks") + + // Verify selected hosts are not from excluded racks + selectedRacks := make(map[string]bool) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack] = true + require.NotContains(t, []string{"rack1", "rack2"}, dn.Rack, "Selected host should not be from excluded rack") + return false + } + return true + }) + } + require.Equal(t, 4, len(selectedRacks), "Should select from 4 different non-excluded racks") +} + +// Test 28: Test excludeHosts with different rack awareness levels +func TestRackExcludeHostsWithDifferentLevels(t *testing.T) { + _, _, zone, _ := setupRackAwareTestEnvWithNodes(t, 6, 3) // 6 racks, 3 nodes per rack = 18 nodes + + excludeHosts := []string{ + "192.168.1.1:17310", "192.168.1.2:17310", // 2 hosts from rack1 + "192.168.1.11:17310", "192.168.1.12:17310", // 2 hosts from rack2 + "192.168.1.21:17310", "192.168.1.22:17310", // 2 hosts from rack3 + } + + // Test with RackAwareNone + param := &selectParam{ + replicaNum: 4, + rackLevel: proto.RackAwareNone, + excludeHosts: excludeHosts, + } + + hosts, _, err := zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Should succeed with RackAwareNone") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Test with RackAwareWeak + param.rackLevel = proto.RackAwareWeak + hosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Should succeed with RackAwareWeak") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Test with RackAwareStrong + param.rackLevel = proto.RackAwareStrong + hosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Should succeed with RackAwareStrong") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Verify excluded hosts are not selected in all cases + for _, host := range hosts { + require.NotContains(t, param.excludeHosts, host, "Selected host should not be in exclude list") + } +} + +// Test 29: Test weak rack awareness fallback mechanism - CRITICAL TEST +// This test verifies that weak rack awareness first tries strong mode, then falls back to weak mode +func TestRackWeakAwarenessFallbackMechanism(t *testing.T) { + // Setup with 2 racks, 4 nodes per rack = 8 total nodes + // This scenario is perfect for testing fallback: need 3 replicas but only 2 racks + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 2, 4) // 2 racks, 4 nodes per rack = 8 nodes + + // Test case 1: Weak rack awareness with 3 replicas (more than available racks) + param := &selectParam{ + replicaNum: 3, // Need 3 replicas but only 2 racks available + rackLevel: proto.RackAwareWeak, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Weak rack awareness should succeed with fallback mechanism") + require.Equal(t, 3, len(hosts), "Should select 3 hosts") + + // Verify the fallback mechanism worked correctly + selectedRacks := make(map[string]int) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // In weak mode with 2 racks and 3 replicas: + // - First 2 replicas should be from different racks (strong mode behavior) + // - The 3rd replica should be from any available rack (weak mode fallback) + require.Equal(t, 2, len(selectedRacks), "Should use both available racks") + require.Equal(t, 3, selectedRacks["rack1"]+selectedRacks["rack2"], "Should select 3 total nodes") + + // At least one rack should have multiple nodes (proving weak mode fallback) + require.True(t, selectedRacks["rack1"] > 1 || selectedRacks["rack2"] > 1, + "At least one rack should have multiple nodes, proving weak mode fallback") + + t.Logf("Selected rack distribution: rack1=%d, rack2=%d", selectedRacks["rack1"], selectedRacks["rack2"]) +} + +// Test 30: Test weak rack awareness with exact rack count +func TestRackWeakAwarenessExactRackCount(t *testing.T) { + // Setup with 3 racks, 3 nodes per rack = 9 total nodes + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 3, 3) // 3 racks, 3 nodes per rack = 9 nodes + + // Test case: Weak rack awareness with 3 replicas (exactly matching rack count) + param := &selectParam{ + replicaNum: 3, // Need 3 replicas, exactly matching 3 racks + rackLevel: proto.RackAwareWeak, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Weak rack awareness should succeed with exact rack count") + require.Equal(t, 3, len(hosts), "Should select 3 hosts") + + // Verify that strong mode was sufficient (one node per rack) + selectedRacks := make(map[string]int) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // With exact rack count, strong mode should be sufficient + require.Equal(t, 3, len(selectedRacks), "Should use all 3 racks") + require.Equal(t, 1, selectedRacks["rack1"], "rack1 should have exactly 1 node") + require.Equal(t, 1, selectedRacks["rack2"], "rack2 should have exactly 1 node") + require.Equal(t, 1, selectedRacks["rack3"], "rack3 should have exactly 1 node") + + t.Logf("Selected rack distribution: rack1=%d, rack2=%d, rack3=%d", + selectedRacks["rack1"], selectedRacks["rack2"], selectedRacks["rack3"]) +} + +// Test 31: Test weak rack awareness with more replicas than racks +func TestRackWeakAwarenessMoreReplicasThanRacks(t *testing.T) { + // Setup with 2 racks, 5 nodes per rack = 10 total nodes + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 2, 5) // 2 racks, 5 nodes per rack = 10 nodes + + // Test case: Weak rack awareness with 4 replicas (more than available racks) + param := &selectParam{ + replicaNum: 4, // Need 4 replicas but only 2 racks available + rackLevel: proto.RackAwareWeak, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Weak rack awareness should succeed with more replicas than racks") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Verify the fallback mechanism worked correctly + selectedRacks := make(map[string]int) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // In weak mode with 2 racks and 4 replicas: + // - First 2 replicas should be from different racks (strong mode behavior) + // - Remaining 2 replicas should be from any available rack (weak mode fallback) + require.Equal(t, 2, len(selectedRacks), "Should use both available racks") + require.Equal(t, 4, selectedRacks["rack1"]+selectedRacks["rack2"], "Should select 4 total nodes") + + // Both racks should have multiple nodes (proving weak mode fallback) + require.True(t, selectedRacks["rack1"] > 1, "rack1 should have multiple nodes") + require.True(t, selectedRacks["rack2"] > 1, "rack2 should have multiple nodes") + + t.Logf("Selected rack distribution: rack1=%d, rack2=%d", selectedRacks["rack1"], selectedRacks["rack2"]) +} + +// Test 32: Test weak rack awareness with host exclusions +func TestRackWeakAwarenessWithHostExclusions(t *testing.T) { + // Setup with 3 racks, 3 nodes per rack = 9 total nodes + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 3, 3) // 3 racks, 3 nodes per rack = 9 nodes + + // Test case: Weak rack awareness with host exclusions + param := &selectParam{ + replicaNum: 4, // Need 4 replicas + rackLevel: proto.RackAwareWeak, + excludeHosts: []string{ + "192.168.1.1:17310", // 1 host from rack1 + "192.168.1.11:17310", // 1 host from rack2 + }, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Weak rack awareness should succeed with host exclusions") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Verify excluded hosts are not selected + for _, host := range hosts { + require.NotContains(t, param.excludeHosts, host, "Selected host should not be in exclude list") + } + + // Verify the fallback mechanism worked correctly + selectedRacks := make(map[string]int) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // Should use all 3 racks, with some racks having multiple nodes due to weak mode fallback + require.Equal(t, 3, len(selectedRacks), "Should use all 3 racks") + require.Equal(t, 4, selectedRacks["rack1"]+selectedRacks["rack2"]+selectedRacks["rack3"], "Should select 4 total nodes") + + t.Logf("Selected rack distribution: rack1=%d, rack2=%d, rack3=%d", + selectedRacks["rack1"], selectedRacks["rack2"], selectedRacks["rack3"]) +} + +// Test 33: Test weak rack awareness with rack exclusions - CORRECTED VERSION +func TestRackWeakAwarenessWithRackExclusions(t *testing.T) { + // Setup with 4 racks, 3 nodes per rack = 12 total nodes + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 4, 3) // 4 racks, 3 nodes per rack = 12 nodes + + // Test case: Weak rack awareness with rack exclusions + param := &selectParam{ + replicaNum: 3, // Need 3 replicas + rackLevel: proto.RackAwareWeak, + excludeRacks: []string{"rack1", "rack2"}, // Exclude 2 racks + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Weak rack awareness should succeed with rack exclusions") + require.Equal(t, 3, len(hosts), "Should select 3 hosts") + + // Verify rack distribution + selectedRacks := make(map[string]int) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // In weak mode, should be able to select from any rack, including initially excluded ones + // The key point is that weak mode prioritizes availability over rack constraints + require.Equal(t, 3, selectedRacks["rack1"]+selectedRacks["rack2"]+selectedRacks["rack3"]+selectedRacks["rack4"], + "Should select 3 total nodes") + + // At least one rack should have multiple nodes (proving weak mode fallback) + totalRacksUsed := 0 + for _, count := range selectedRacks { + if count > 0 { + totalRacksUsed++ + } + } + require.True(t, totalRacksUsed <= 3, "Should use at most 3 racks (since we need 3 replicas)") + + t.Logf("Selected rack distribution: rack1=%d, rack2=%d, rack3=%d, rack4=%d", + selectedRacks["rack1"], selectedRacks["rack2"], selectedRacks["rack3"], selectedRacks["rack4"]) +} + +// Test 35: Test weak rack awareness comparison with strong mode +func TestRackWeakVsStrongAwarenessComparison(t *testing.T) { + // Setup with 2 racks, 4 nodes per rack = 8 total nodes + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 2, 4) // 2 racks, 4 nodes per rack = 8 nodes + + // Test case: Need 3 replicas but only 2 racks available + replicaNum := 3 + + // Test 1: Strong rack awareness should fail + strongParam := &selectParam{ + replicaNum: replicaNum, + rackLevel: proto.RackAwareStrong, + } + + hosts, _, err := ns.selectNodesWithRack(strongParam, DataNodeType, proto.StoreModeMem) + require.Error(t, err, "Strong rack awareness should fail when not enough racks") + require.Nil(t, hosts, "Strong rack awareness should return nil hosts when failing") + + // Test 2: Weak rack awareness should succeed with fallback + weakParam := &selectParam{ + replicaNum: replicaNum, + rackLevel: proto.RackAwareWeak, + } + + hosts, _, err = ns.selectNodesWithRack(weakParam, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Weak rack awareness should succeed with fallback mechanism") + require.Equal(t, replicaNum, len(hosts), "Should select %d hosts with weak rack awareness", replicaNum) + + // Verify weak mode used fallback mechanism + selectedRacks := make(map[string]int) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // Weak mode should use both racks, with at least one having multiple nodes + require.Equal(t, 2, len(selectedRacks), "Weak mode should use both available racks") + require.Equal(t, replicaNum, selectedRacks["rack1"]+selectedRacks["rack2"], "Should select %d total nodes", replicaNum) + require.True(t, selectedRacks["rack1"] > 1 || selectedRacks["rack2"] > 1, + "At least one rack should have multiple nodes, proving weak mode fallback") + + t.Logf("Strong mode: Failed (as expected)") + t.Logf("Weak mode: Selected rack1=%d, rack2=%d (fallback mechanism working)", selectedRacks["rack1"], selectedRacks["rack2"]) +} + +// Test 36: Test weak rack awareness step-by-step fallback verification +func TestRackWeakAwarenessStepByStepFallback(t *testing.T) { + // Setup with 2 racks, 3 nodes per rack = 6 total nodes + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 2, 3) // 2 racks, 3 nodes per rack = 6 nodes + + // Test case: Need 4 replicas but only 2 racks available + // This should demonstrate the fallback mechanism clearly + param := &selectParam{ + replicaNum: 4, // Need 4 replicas but only 2 racks available + rackLevel: proto.RackAwareWeak, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Weak rack awareness should succeed with step-by-step fallback") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Verify the step-by-step fallback mechanism + selectedRacks := make(map[string]int) + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // Expected behavior in weak mode with 2 racks and 4 replicas: + // Step 1: Try strong mode - select 1 node from each rack (2 nodes total) + // Step 2: Fallback to weak mode - select remaining 2 nodes from any available rack + // Result: Both racks should be used, and both should have multiple nodes + require.Equal(t, 2, len(selectedRacks), "Should use both available racks") + require.Equal(t, 4, selectedRacks["rack1"]+selectedRacks["rack2"], "Should select 4 total nodes") + require.Equal(t, 2, selectedRacks["rack1"], "rack1 should have exactly 2 nodes") + require.Equal(t, 2, selectedRacks["rack2"], "rack2 should have exactly 2 nodes") + + t.Logf("Step-by-step fallback verification: rack1=%d, rack2=%d", selectedRacks["rack1"], selectedRacks["rack2"]) +} + +// Test 37: Test weak rack awareness with insufficient nodes after exclusions +func TestRackWeakAwarenessInsufficientNodes(t *testing.T) { + // Setup with 2 racks, 2 nodes per rack = 4 total nodes + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 2, 2) // 2 racks, 2 nodes per rack = 4 nodes + + // Test case: Exclude too many nodes to make selection impossible + param := &selectParam{ + replicaNum: 3, // Need 3 replicas + rackLevel: proto.RackAwareWeak, + excludeHosts: []string{ + "192.168.1.1:17310", "192.168.1.2:17310", // All hosts from rack1 + "192.168.1.11:17310", // 1 host from rack2 + }, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.Error(t, err, "Should fail when insufficient nodes available after exclusions") + require.Nil(t, hosts, "Should return nil hosts when failing") +} + +// Test 38: Test weak rack awareness with mixed node availability +func TestRackWeakAwarenessMixedAvailability(t *testing.T) { + // Setup with 3 racks, 3 nodes per rack = 9 total nodes + _, _, _, ns := setupRackAwareTestEnvWithNodes(t, 3, 3) // 3 racks, 3 nodes per rack = 9 nodes + + // Make some nodes unavailable + nodeCount := 0 + ns.dataNodes.Range(func(key, value interface{}) bool { + dn := value.(*DataNode) + if nodeCount%4 == 0 { // Make every 4th node unavailable + dn.isActive = false + } + nodeCount++ + return true + }) + + // Test case: Weak rack awareness with mixed node availability + param := &selectParam{ + replicaNum: 4, // Need 4 replicas + rackLevel: proto.RackAwareWeak, + } + + hosts, _, err := ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeMem) + require.NoError(t, err, "Weak rack awareness should succeed with mixed node availability") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Verify all selected hosts are active + for _, host := range hosts { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + require.True(t, dn.isActive, "Selected host should be active") + return false + } + return true + }) + } + + t.Logf("Successfully selected 4 hosts with mixed node availability") +} + +// Test 34: Test weak rack awareness at zone level with rack exclusions +func TestZoneWeakRackAwarenessWithRackExclusions(t *testing.T) { + // Setup with 4 racks, 3 nodes per rack = 12 total nodes + _, _, zone, _ := setupRackAwareTestEnvWithNodes(t, 4, 3) // 4 racks, 3 nodes per rack = 12 nodes + + // Test case: Zone level weak rack awareness with rack exclusions + param := &selectParam{ + replicaNum: 3, // Need 3 replicas + rackLevel: proto.RackAwareWeak, + excludeRacks: []string{"rack1", "rack2"}, // Exclude 2 racks + } + + hosts, _, err := zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Zone level weak rack awareness should succeed with rack exclusions") + require.Equal(t, 3, len(hosts), "Should select 3 hosts") + + // Verify rack distribution + selectedRacks := make(map[string]int) + for _, host := range hosts { + // Find the nodeSet that contains this host + var foundNodeSet *nodeSet + zone.nsLock.RLock() + for _, ns := range zone.nodeSetMap { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + foundNodeSet = ns + return false + } + return true + }) + if foundNodeSet != nil { + break + } + } + zone.nsLock.RUnlock() + + require.NotNil(t, foundNodeSet, "Should find the nodeSet containing host %s", host) + + // Find the specific node to get its rack + foundNodeSet.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // In weak mode at zone level, should be able to select from any rack, including initially excluded ones + // The key point is that weak mode prioritizes availability over rack constraints + require.Equal(t, 3, selectedRacks["rack1"]+selectedRacks["rack2"]+selectedRacks["rack3"]+selectedRacks["rack4"], + "Should select 3 total nodes") + + // At least one rack should have multiple nodes (proving weak mode fallback) + totalRacksUsed := 0 + for _, count := range selectedRacks { + if count > 0 { + totalRacksUsed++ + } + } + require.True(t, totalRacksUsed <= 3, "Should use at most 3 racks (since we need 3 replicas)") + + t.Logf("Zone level selected rack distribution: rack1=%d, rack2=%d, rack3=%d, rack4=%d", + selectedRacks["rack1"], selectedRacks["rack2"], selectedRacks["rack3"], selectedRacks["rack4"]) +} + +// Test 35: Test zone level weak rack awareness with insufficient racks scenario +func TestZoneWeakRackAwarenessInsufficientRacks(t *testing.T) { + // Setup with 2 racks, 4 nodes per rack = 8 total nodes + _, _, zone, _ := setupRackAwareTestEnvWithNodes(t, 2, 4) // 2 racks, 4 nodes per rack = 8 nodes + + // Test case: Need 3 replicas but only 2 racks available + param := &selectParam{ + replicaNum: 3, // Need 3 replicas but only 2 racks available + rackLevel: proto.RackAwareWeak, + } + + hosts, _, err := zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Zone level weak rack awareness should succeed with insufficient racks") + require.Equal(t, 3, len(hosts), "Should select 3 hosts") + + // Verify rack distribution + selectedRacks := make(map[string]int) + for _, host := range hosts { + // Find the nodeSet that contains this host + var foundNodeSet *nodeSet + zone.nsLock.RLock() + for _, ns := range zone.nodeSetMap { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + foundNodeSet = ns + return false + } + return true + }) + if foundNodeSet != nil { + break + } + } + zone.nsLock.RUnlock() + + require.NotNil(t, foundNodeSet, "Should find the nodeSet containing host %s", host) + + // Find the specific node to get its rack + foundNodeSet.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // Should use both available racks + require.Equal(t, 2, len(selectedRacks), "Should use both available racks") + require.Equal(t, 3, selectedRacks["rack1"]+selectedRacks["rack2"], "Should select 3 total nodes") + + // At least one rack should have multiple nodes (proving weak mode fallback) + require.True(t, selectedRacks["rack1"] > 1 || selectedRacks["rack2"] > 1, + "At least one rack should have multiple nodes, proving weak mode fallback") + + t.Logf("Zone level insufficient racks scenario: rack1=%d, rack2=%d", + selectedRacks["rack1"], selectedRacks["rack2"]) +} + +// Test 36: Test zone level weak rack awareness with host exclusions +func TestZoneWeakRackAwarenessWithHostExclusions(t *testing.T) { + // Setup with 3 racks, 3 nodes per rack = 9 total nodes + _, _, zone, _ := setupRackAwareTestEnvWithNodes(t, 3, 3) // 3 racks, 3 nodes per rack = 9 nodes + + // Test case: Zone level weak rack awareness with host exclusions + param := &selectParam{ + replicaNum: 4, // Need 4 replicas + rackLevel: proto.RackAwareWeak, + excludeHosts: []string{ + "192.168.1.1:17310", // 1 host from rack1 + "192.168.1.11:17310", // 1 host from rack2 + }, + } + + hosts, _, err := zone.getAvailNodeHosts(TypeDataPartition, param) + require.NoError(t, err, "Zone level weak rack awareness should succeed with host exclusions") + require.Equal(t, 4, len(hosts), "Should select 4 hosts") + + // Verify excluded hosts are not selected + for _, host := range hosts { + require.NotContains(t, param.excludeHosts, host, "Selected host should not be in exclude list") + } + + // Verify rack distribution + selectedRacks := make(map[string]int) + for _, host := range hosts { + // Find the nodeSet that contains this host + var foundNodeSet *nodeSet + zone.nsLock.RLock() + for _, ns := range zone.nodeSetMap { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + foundNodeSet = ns + return false + } + return true + }) + if foundNodeSet != nil { + break + } + } + zone.nsLock.RUnlock() + + require.NotNil(t, foundNodeSet, "Should find the nodeSet containing host %s", host) + + // Find the specific node to get its rack + foundNodeSet.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // Should use all 3 racks, with some racks having multiple nodes due to weak mode fallback + require.Equal(t, 3, len(selectedRacks), "Should use all 3 racks") + require.Equal(t, 4, selectedRacks["rack1"]+selectedRacks["rack2"]+selectedRacks["rack3"], "Should select 4 total nodes") + + t.Logf("Zone level host exclusions scenario: rack1=%d, rack2=%d, rack3=%d", + selectedRacks["rack1"], selectedRacks["rack2"], selectedRacks["rack3"]) +} + +// Test 37: Test zone level weak vs strong rack awareness comparison +func TestZoneWeakVsStrongRackAwarenessComparison(t *testing.T) { + // Setup with 2 racks, 3 nodes per rack = 6 total nodes + _, _, zone, _ := setupRackAwareTestEnvWithNodes(t, 2, 3) // 2 racks, 3 nodes per rack = 6 nodes + + // Test case: Need 3 replicas but only 2 racks available + replicaNum := 3 + + // Test 1: Strong rack awareness should fail + strongParam := &selectParam{ + replicaNum: replicaNum, + rackLevel: proto.RackAwareStrong, + } + + hosts, _, err := zone.getAvailNodeHosts(TypeDataPartition, strongParam) + require.Error(t, err, "Zone level strong rack awareness should fail when not enough racks") + require.Nil(t, hosts, "Zone level strong rack awareness should return nil hosts when failing") + + // Test 2: Weak rack awareness should succeed with fallback + weakParam := &selectParam{ + replicaNum: replicaNum, + rackLevel: proto.RackAwareWeak, + } + + hosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, weakParam) + require.NoError(t, err, "Zone level weak rack awareness should succeed with fallback mechanism") + require.Equal(t, replicaNum, len(hosts), "Should select %d hosts with weak rack awareness", replicaNum) + + // Verify weak mode used fallback mechanism + selectedRacks := make(map[string]int) + for _, host := range hosts { + // Find the nodeSet that contains this host + var foundNodeSet *nodeSet + zone.nsLock.RLock() + for _, ns := range zone.nodeSetMap { + ns.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + foundNodeSet = ns + return false + } + return true + }) + if foundNodeSet != nil { + break + } + } + zone.nsLock.RUnlock() + + require.NotNil(t, foundNodeSet, "Should find the nodeSet containing host %s", host) + + // Find the specific node to get its rack + foundNodeSet.dataNodes.Range(func(key, value interface{}) bool { + if key.(string) == host { + dn := value.(*DataNode) + selectedRacks[dn.Rack]++ + return false + } + return true + }) + } + + // Weak mode should use both racks, with at least one having multiple nodes + require.Equal(t, 2, len(selectedRacks), "Weak mode should use both available racks") + require.Equal(t, replicaNum, selectedRacks["rack1"]+selectedRacks["rack2"], "Should select %d total nodes", replicaNum) + require.True(t, selectedRacks["rack1"] > 1 || selectedRacks["rack2"] > 1, + "At least one rack should have multiple nodes, proving weak mode fallback") + + t.Logf("Zone level strong mode: Failed (as expected)") + t.Logf("Zone level weak mode: Selected rack1=%d, rack2=%d (fallback mechanism working)", + selectedRacks["rack1"], selectedRacks["rack2"]) +} diff --git a/master/topology_test.go b/master/topology_test.go index 9afc08853..9d89a7786 100644 --- a/master/topology_test.go +++ b/master/topology_test.go @@ -54,7 +54,15 @@ func TestSingleZone(t *testing.T) { // single zone normal zones, err = topo.allocZonesForNode(&topo.dataTopology, replicaNum, replicaNum, nil, []*Zone{}, proto.MediaType_Unspecified) require.NoError(t, err) - newHosts, _, err := zones[0].getAvailNodeHosts(TypeDataPartition, nil, nil, replicaNum) + + param := &selectParam{ + excludeNodeSets: nil, + replicaNum: replicaNum, + excludeHosts: nil, + rackLevel: c.getRackAwareLevel(), + excludeRacks: nil, + } + newHosts, _, err := zones[0].getAvailNodeHosts(TypeDataPartition, param) require.NoError(t, err) t.Log(newHosts) topo.deleteDataNode(createDataNodeForTopo(mds1Addr, zoneName, nodeSet)) @@ -134,17 +142,17 @@ func TestAllocZones(t *testing.T) { cluster.cfg = newClusterConfig() // don't cross zone - hosts, _, err := cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, replicaNum, 1, "", proto.MediaType_Unspecified) + hosts, _, err := cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, replicaNum, 1, "", proto.MediaType_Unspecified, proto.RackAwareNone) require.NoError(t, err) t.Logf("ChooseTargetDataHosts in single zone,hosts[%v]", hosts) // cross zone - _, _, err = cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, replicaNum, 2, "", proto.MediaType_Unspecified) + _, _, err = cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, replicaNum, 2, "", proto.MediaType_Unspecified, proto.RackAwareNone) require.NoError(t, err) // specific zone - hosts, _, err = cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, 3, 2, zoneName1+","+zoneName2, proto.MediaType_Unspecified) + hosts, _, err = cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, 3, 2, zoneName1+","+zoneName2, proto.MediaType_Unspecified, proto.RackAwareNone) require.NoError(t, err) require.EqualValues(t, getZoneCntFunc(hosts), 2) diff --git a/master/vol.go b/master/vol.go index 12a2cf398..3e6d240c2 100644 --- a/master/vol.go +++ b/master/vol.go @@ -1714,7 +1714,7 @@ func (vol *Vol) doCreateMetaPartition(c *Cluster, start, end uint64) (mp *MetaPa var excludeZone []string zoneNum := c.decideZoneNum(vol, proto.StorageClass_Unspecified) if hosts, peers, err = c.getHostFromNormalZone(nodeType, excludeZone, nil, nil, - int(vol.mpReplicaNum), zoneNum, vol.zoneName, proto.StorageClass_Unspecified); err != nil { + int(vol.mpReplicaNum), zoneNum, vol.zoneName, proto.StorageClass_Unspecified, c.getRackAwareLevel()); err != nil { log.LogErrorf("action[doCreateMetaPartition] getHostFromNormalZone err[%v]", err) return nil, errors.NewError(err) }