Line data Source code
1 : /**
2 : * Copyright (c) 2025 Huawei Technologies Co., Ltd.
3 : * This program is free software, you can redistribute it and/or modify it under the terms and conditions of
4 : * CANN Open Software License Agreement Version 2.0 (the "License").
5 : * Please refer to the License for details. You may not use this file except in compliance with the License.
6 : * THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
7 : * INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE.
8 : * See LICENSE in the root of the software repository for the full text of the License.
9 : */
10 :
11 : #include "topoinfo_detect.h"
12 : #include <string>
13 : #include "adapter_rts_common.h"
14 : #include "hccl_whitelist.h"
15 : #include "hccl_socket.h"
16 : #include "sal_pub.h"
17 : #include "device_capacity.h"
18 : #include "preempt_port_manager.h"
19 :
20 : using namespace std;
21 : namespace hccl {
22 : const u32 TOPO_EXCHANGE_SERVER_STATUS_IDLE = 0;
23 : const u32 TOPO_EXCHANGE_SERVER_STATUS_RUNING = 1;
24 : const u32 TOPO_EXCHANGE_SERVER_STATUS_ERROR = 2;
25 : UniversalConcurrentMap<u32, volatile u32> TopoInfoDetect::g_topoExchangeServerStatus_;
26 :
27 20 : TopoInfoDetect::TopoInfoDetect() : deviceLogicID_(INVALID_INT), localRankInfo_(),
28 20 : clusterTopoInfo_(), isInterSuperPodRetryEnable_(GetExternalInputInterSuperPodRetryEnable())
29 : {
30 20 : }
31 :
32 20 : TopoInfoDetect::~TopoInfoDetect()
33 : {
34 20 : if (exchangeServerThreadPtr_ && exchangeServerThreadPtr_->joinable()) {
35 13 : exchangeServerThreadPtr_->join();
36 : }
37 20 : exchangeServerThreadPtr_ = nullptr;
38 20 : pTopoExchangeServer_ = nullptr;
39 20 : (void)Teardown();
40 20 : return;
41 20 : }
42 :
43 0 : HcclResult TopoInfoDetect::GetServerConnections(std::map<u32, std::shared_ptr<HcclSocket>> &connectSockets)
44 : {
45 0 : if (pTopoExchangeServer_) {
46 0 : return pTopoExchangeServer_->GetConnections(connectSockets);
47 : } else {
48 0 : return HCCL_SUCCESS;
49 : }
50 : }
51 :
52 0 : HcclResult TopoInfoDetect::GetAgentListenSocket(HcclSocketPortConfig &commPortConfig)
53 : {
54 : // 将抢占的端口传入comm connection参数中
55 0 : commPortConfig = commPortConfig_;
56 0 : return HCCL_SUCCESS;
57 : }
58 :
59 0 : HcclResult TopoInfoDetect::GetAgentConnection(std::shared_ptr<HcclSocket> &connectSocket)
60 : {
61 0 : CHK_SMART_PTR_NULL(pTopoExchangeAgent_);
62 0 : return pTopoExchangeAgent_->GetConnection(connectSocket);
63 : }
64 :
65 0 : HcclResult TopoInfoDetect::SendGroupLeaderPort(std::shared_ptr<HcclSocket> &connectSocket, HcclRankHandle &rankHandle)
66 : {
67 0 : CHK_SMART_PTR_NULL(pTopoExchangeAgent_);
68 0 : HcclResult ret = pTopoExchangeAgent_->SendGroupLeaderPortInfo(connectSocket, rankHandle);
69 0 : CHK_PRT_RET(ret != HCCL_SUCCESS,
70 : HCCL_ERROR("[Setup][SendGroupLeaderPortInfo]SendGroupLeaderPortInfo to root failed, ret[%d]",
71 : ret), ret);
72 0 : return HCCL_SUCCESS;
73 : }
74 :
75 0 : void TopoInfoDetect::SetupTopoGroupLeader(s32 devicePhysicID, s32 deviceLogicID, HcclIpAddress hostIP, u32 hostPort,
76 : vector<HcclIpAddress> whitelist, HcclNetDevCtx netDevCtx, std::shared_ptr<HcclSocket> listenSocket,
77 : std::shared_ptr<HcclSocket> grpLeaderToRoot, bool isMasterInfo)
78 : {
79 : //给当前线程添加名字
80 0 : SetThreadName("Hccl_TopoDetect_GroupLeader");
81 :
82 0 : HcclResult ret = hrtSetDevice(deviceLogicID);
83 0 : if (ret != HCCL_SUCCESS) {
84 0 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
85 0 : status = TOPO_EXCHANGE_SERVER_STATUS_ERROR;
86 0 : });
87 0 : HCCL_ERROR("[Setup][TopoExchangeServer]set device[%d] failed, ret[%u]", deviceLogicID, ret);
88 0 : return;
89 : }
90 :
91 0 : pTopoExchangeServer_.reset(new (nothrow) TopoInfoExchangeServer(hostIP, hostPort, whitelist, netDevCtx,
92 0 : listenSocket, grpLeaderToRoot, rootInfo_.identifier));
93 0 : if (!pTopoExchangeServer_) {
94 0 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
95 0 : status = TOPO_EXCHANGE_SERVER_STATUS_ERROR;
96 0 : });
97 0 : HCCL_ERROR("[Setup][TopoExchangeServer]build topoExchangeServer failed. ");
98 : } else {
99 0 : ret = isMasterInfo ? pTopoExchangeServer_->SetupByMasterInfo() : pTopoExchangeServer_->SetupGroupLeader();
100 0 : if (ret != HCCL_SUCCESS) {
101 0 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
102 0 : status = TOPO_EXCHANGE_SERVER_STATUS_ERROR;
103 0 : });
104 0 : HCCL_ERROR("[Setup][TopoExchangeServer]setup topoExchangeServer failed, ret[%u]", ret);
105 : }
106 : }
107 :
108 0 : ret = hrtResetDevice(deviceLogicID);
109 0 : if (ret != HCCL_SUCCESS) {
110 0 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
111 0 : status = TOPO_EXCHANGE_SERVER_STATUS_ERROR;
112 0 : });
113 0 : HCCL_ERROR("[Setup][TopoExchangeServer]reset device[%d] failed, ret[%u]", deviceLogicID, ret);
114 0 : return;
115 : }
116 0 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
117 0 : status = TOPO_EXCHANGE_SERVER_STATUS_IDLE;
118 0 : });
119 : }
120 :
121 13 : void TopoInfoDetect::SetupTopoExchangeServer(s32 devicePhysicID, s32 deviceLogicID, HcclIpAddress hostIP, u32 hostPort,
122 : vector<HcclIpAddress> whitelist, HcclNetDevCtx netDevCtx,
123 : std::shared_ptr<HcclSocket> listenSocket, bool isMasterInfo)
124 : {
125 : //给当前线程添加名字
126 13 : SetThreadName("Hccl_TopoDetect");
127 :
128 13 : HcclResult ret = hrtSetDevice(deviceLogicID);
129 13 : if (ret != HCCL_SUCCESS) {
130 0 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
131 0 : status = TOPO_EXCHANGE_SERVER_STATUS_ERROR;
132 0 : });
133 0 : HCCL_ERROR("[Setup][TopoExchangeServer]set device[%d] failed, ret[%u]", deviceLogicID, ret);
134 0 : return;
135 : }
136 :
137 13 : pTopoExchangeServer_.reset(new (nothrow) TopoInfoExchangeServer(hostIP, hostPort, whitelist, netDevCtx,
138 39 : listenSocket, rootInfo_.identifier));
139 13 : if (!pTopoExchangeServer_) {
140 0 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
141 0 : status = TOPO_EXCHANGE_SERVER_STATUS_ERROR;
142 0 : });
143 0 : HCCL_ERROR("[Setup][TopoExchangeServer]build topoExchangeServer failed. ");
144 : } else {
145 13 : ret = isMasterInfo ? pTopoExchangeServer_->SetupByMasterInfo() : pTopoExchangeServer_->Setup();
146 13 : if (ret != HCCL_SUCCESS) {
147 13 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
148 13 : status = TOPO_EXCHANGE_SERVER_STATUS_ERROR;
149 13 : });
150 13 : HCCL_ERROR("[Setup][TopoExchangeServer]setup topoExchangeServer failed, ret[%u]", ret);
151 : }
152 : }
153 :
154 13 : ret = hrtResetDevice(deviceLogicID);
155 13 : if (ret != HCCL_SUCCESS) {
156 0 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
157 0 : status = TOPO_EXCHANGE_SERVER_STATUS_ERROR;
158 0 : });
159 0 : HCCL_ERROR("[Setup][TopoExchangeServer]reset device[%d] failed, ret[%u]", deviceLogicID, ret);
160 0 : return;
161 : }
162 13 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
163 13 : status = TOPO_EXCHANGE_SERVER_STATUS_IDLE;
164 13 : });
165 : }
166 0 : HcclResult TopoInfoDetect::SetupServerByMasterInfo(const HcclIpAddress& masterIP, u32 masterPort, const HcclRootHandle &rootInfo)
167 : {
168 0 : CHK_RET(hrtGetDevice(&deviceLogicID_));
169 0 : CHK_RET(hrtGetDevicePhyIdByIndex(deviceLogicID_, devicePhysicID_));
170 0 : vector<HcclIpAddress> whitelist;
171 0 : if (GetExternalInputHcclEnableWhitelist() == HCCL_WHITELIST_ON) {
172 0 : CHK_RET(ReadHostSocketWhitelist(whitelist));
173 : }
174 0 : rootInfo_ = rootInfo;
175 0 : CHK_RET(HcclNetInit(NICDeployment::NIC_DEPLOYMENT_HOST, devicePhysicID_, deviceLogicID_, true));
176 0 : HcclResult ret = StartRootNetwork(masterIP, masterPort, GetExternalInputHostSocketPortRange());
177 0 : CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_ERROR("[%s][%s]%s failed, masterIP[%s] and masterPort[%u] ret[%u]",
178 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), __func__,
179 : masterIP.GetReadableAddress(), masterPort, ret), ret);
180 :
181 0 : if (GetExternalInputHcclEnableWhitelist() == HCCL_WHITELIST_ON) {
182 0 : CHK_RET(AddSocketWhiteList(masterPort, whitelist));
183 : }
184 :
185 0 : u32 hostPort = GetExternalInputMasterInfo().port;
186 0 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
187 0 : status = TOPO_EXCHANGE_SERVER_STATUS_RUNING;
188 0 : });
189 :
190 0 : thread threadHandle(&TopoInfoDetect::SetupTopoExchangeServer, this, devicePhysicID_, deviceLogicID_,
191 0 : masterIP, GetExternalInputMasterInfo().port, whitelist, serverPortCtx_, listenSocket_, true);
192 0 : threadHandle.detach();
193 :
194 0 : return HCCL_SUCCESS;
195 0 : }
196 :
197 15 : HcclResult TopoInfoDetect::SetupServer(HcclRootHandle &rootInfo)
198 : {
199 15 : CHK_RET(hrtGetDevice(&deviceLogicID_));
200 :
201 15 : vector<HcclIpAddress> whitelist;
202 15 : if (GetExternalInputHcclEnableWhitelist() == HCCL_WHITELIST_ON) {
203 3 : CHK_RET(ReadHostSocketWhitelist(whitelist));
204 : }
205 14 : HcclIpAddress hostIP = GetBootstrapHostIP();
206 14 : CHK_RET(hrtGetDevicePhyIdByIndex(deviceLogicID_, devicePhysicID_, true));
207 14 : HCCL_INFO("[Setup][hcclIfBasePort]deviceLogicID_[%u], devicePhysicID_[%u]", deviceLogicID_, devicePhysicID_);
208 :
209 : // true代表感知白名单disable配置
210 14 : CHK_RET(HcclNetInit(NICDeployment::NIC_DEPLOYMENT_HOST, devicePhysicID_, deviceLogicID_, true));
211 :
212 14 : CHK_RET(GetRootHostIP(whitelist, hostIP, devicePhysicID_));
213 13 : SetBootstrapHostIP(hostIP);
214 :
215 13 : u32 deviceNum = 0;
216 13 : CHK_RET(hrtGetDeviceCount(&deviceNum));
217 13 : CHK_PRT_RET((static_cast<u32>(deviceLogicID_) >= deviceNum),
218 : HCCL_ERROR("[Setup][Server]deviceLogicID[%d] is invalid,deviceNum[%d].", deviceLogicID_, deviceNum),
219 : HCCL_E_PARA);
220 :
221 13 : u32 hostPort = HCCL_INVALID_PORT;
222 13 : std::vector<HcclSocketPortRange> portRanges;
223 13 : if (!GetExternalInputHostPortSwitch()) {
224 13 : if (GetExternalInputHcclIfBasePort() == HCCL_INVALID_PORT) {
225 : // 若没有设置HCCL_HOST_SOCKET_PORT_RANGE和HCCL_IF_BASE_PORT 使用自动调整监听端口range[60000,60031]
226 13 : HCCL_RUN_INFO("[Setup][Server] user not set base port and port range, use default port range[%u, %u]", HOST_CONTROL_BASE_PORT, HOST_CONTROL_BASE_PORT + 31);
227 13 : portRanges.push_back({HOST_CONTROL_BASE_PORT, HOST_CONTROL_BASE_PORT + 31});
228 : } else {
229 0 : hostPort = devicePhysicID_ + GetExternalInputHcclIfBasePort();
230 : }
231 : } else {
232 0 : portRanges = GetExternalInputHostSocketPortRange();
233 : }
234 13 : HcclResult ret = StartRootNetwork(hostIP, hostPort, portRanges);
235 13 : CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_ERROR("[%s][%s]%s failed, hostIP[%s] and hostPort[%u] ret[%u]",
236 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), __func__,
237 : hostIP.GetReadableAddress(), hostPort, ret), ret);
238 13 : CHK_RET(GenerateRootInfo(hostIP, hostPort, devicePhysicID_, rootInfo_));
239 13 : if (GetExternalInputHcclEnableWhitelist() == HCCL_WHITELIST_ON) {
240 1 : CHK_RET(AddSocketWhiteList(hostPort, whitelist));
241 : }
242 :
243 13 : g_topoExchangeServerStatus_.EmplaceAndUpdate(hostPort, [] (volatile u32 &status) {
244 13 : status = TOPO_EXCHANGE_SERVER_STATUS_RUNING;
245 13 : });
246 13 : exchangeServerThreadPtr_.reset(new (nothrow) thread(&TopoInfoDetect::SetupTopoExchangeServer, this,
247 13 : devicePhysicID_, deviceLogicID_, hostIP, hostPort, whitelist, serverPortCtx_, listenSocket_, false));
248 13 : CHK_SMART_PTR_NULL(exchangeServerThreadPtr_);
249 :
250 13 : rootInfo = rootInfo_;
251 13 : HCCL_INFO("setup topo exchange server complete, identifier[%s]", rootInfo.identifier);
252 13 : return HCCL_SUCCESS;
253 15 : }
254 :
255 3 : HcclResult TopoInfoDetect::GroupLeaderListen(HcclRankHandle &rankHandle, vector<HcclIpAddress> &whitelist)
256 : {
257 3 : CHK_RET(hrtGetDevice(&deviceLogicID_));
258 :
259 3 : if (GetExternalInputHcclEnableWhitelist() == HCCL_WHITELIST_ON) {
260 0 : CHK_RET(ReadHostSocketWhitelist(whitelist));
261 : }
262 3 : HcclIpAddress hostIP = GetBootstrapHostIP();
263 3 : CHK_RET(hrtGetDevicePhyIdByIndex(deviceLogicID_, devicePhysicID_, true));
264 3 : HCCL_INFO("[Setup][hcclIfBasePort]deviceLogicID_[%u], devicePhysicID_[%u]", deviceLogicID_, devicePhysicID_);
265 :
266 : // true代表感知白名单disable配置
267 3 : CHK_RET(HcclNetInit(NICDeployment::NIC_DEPLOYMENT_HOST, devicePhysicID_, deviceLogicID_, true));
268 :
269 3 : CHK_RET(GetRootHostIP(whitelist, hostIP, devicePhysicID_));
270 3 : SetBootstrapHostIP(hostIP);
271 :
272 3 : u32 deviceNum = INVALID_INT;
273 3 : CHK_RET(hrtGetDeviceCount(&deviceNum));
274 3 : CHK_PRT_RET((static_cast<u32>(deviceLogicID_) >= deviceNum),
275 : HCCL_ERROR("[Setup][GroupLeader]deviceLogicID[%d] is invalid,deviceNum[%d].", deviceLogicID_, deviceNum),
276 : HCCL_E_PARA);
277 :
278 3 : u32 hostPort = HCCL_INVALID_PORT ;
279 3 : std::vector<HcclSocketPortRange> portRanges;
280 3 : if (!GetExternalInputHostPortSwitch()) {
281 : // 不开启host侧端口范围配置, 则使用默认端口
282 2 : if (GetExternalInputHcclIfBasePort() == HCCL_INVALID_PORT) {
283 : // 若没有设置HCCL_HOST_SOCKET_PORT_RANGE和HCCL_IF_BASE_PORT 使用自动调整监听端口range[60000,60031]
284 2 : HCCL_RUN_INFO("[Setup][GroupLeader] user not set base port and port range, use default port range[%u, %u]",
285 : HOST_CONTROL_BASE_PORT, HOST_CONTROL_BASE_PORT + 31);
286 2 : portRanges.push_back({HOST_CONTROL_BASE_PORT, HOST_CONTROL_BASE_PORT + 31});
287 : } else {
288 0 : hostPort = devicePhysicID_ + GetExternalInputHcclIfBasePort() + TOPO_GROUPLEADER_PORT_OFFSET;
289 : }
290 : } else {
291 1 : portRanges = GetExternalInputHostSocketPortRange();
292 : }
293 :
294 3 : CHK_RET(StartGroupLeaderNetwork(whitelist, hostIP, hostPort, portRanges));
295 3 : CHK_RET(GenerateRootInfo(hostIP, hostPort, devicePhysicID_, rankHandle));
296 :
297 3 : HCCL_INFO("rank bind port complete, port[%u]", hostPort);
298 3 : return HCCL_SUCCESS;
299 3 : }
300 :
301 0 : HcclResult TopoInfoDetect::GroupLeaderAccept(HcclRankHandle &grpLeaderInfo, vector<HcclIpAddress> whitelist,
302 : std::shared_ptr<HcclSocket> grpLeaderToRoot)
303 : {
304 0 : rootInfo_ = grpLeaderInfo;
305 0 : exchangeServerThreadPtr_.reset(new (nothrow) thread(&TopoInfoDetect::SetupTopoGroupLeader, this,
306 0 : devicePhysicID_, deviceLogicID_, bootstrapHostIP_, rootInfo_.port, whitelist, serverPortCtx_, listenSocket_,
307 0 : grpLeaderToRoot, false));
308 0 : CHK_SMART_PTR_NULL(exchangeServerThreadPtr_);
309 :
310 0 : HCCL_INFO("setup group leader server complete, identifier[%s]", grpLeaderInfo.identifier);
311 0 : return HCCL_SUCCESS;
312 : }
313 :
314 13 : HcclResult TopoInfoDetect::GenerateRootInfo(const HcclIpAddress &hostIP, u32 hostPort, u32 devicePhysicID, HcclRootHandle &rootInfo)
315 : {
316 13 : u64 timestamp = 0;
317 13 : CHK_RET(SalGetCurrentTimestamp(timestamp));
318 :
319 13 : string identifier = hostIP.GetReadableAddress();
320 13 : identifier.append("_");
321 13 : identifier.append(to_string(hostPort));
322 13 : identifier.append("_");
323 13 : identifier.append(to_string(devicePhysicID));
324 13 : identifier.append("_");
325 13 : identifier.append(to_string(timestamp));
326 13 : CHK_PRT_RET((identifier.length() >= ROOTINFO_INDENTIFIER_MAX_LENGTH),
327 : HCCL_ERROR("[Setup][Server]rootinfo identifier len[%u] is invalid.", identifier.length()), HCCL_E_INTERNAL);
328 13 : s32 sret = memcpy_s(&rootInfo.identifier[0], sizeof(rootInfo.identifier), identifier.c_str(),
329 13 : (identifier.length() + 1));
330 13 : CHK_PRT_RET(sret != EOK, HCCL_ERROR("[Setup][Server]errNo[0x%016llx] memcpy failed. ret[%d], params:"\
331 : "destMaxSize[%zu],count[%zu]", HCOM_ERROR_CODE(HCCL_E_MEMORY), sret, sizeof(rootInfo.identifier),
332 : (identifier.length() + 1)), HCCL_E_MEMORY);
333 13 : s32 sRet = strncpy_s(rootInfo.ip, sizeof(rootInfo.ip), hostIP.GetReadableIP(), strlen(hostIP.GetReadableIP()));
334 13 : CHK_PRT_RET(sRet != EOK, HCCL_ERROR("[Setup][Server]str copy fail. return[%d]", sRet), HCCL_E_INTERNAL);
335 13 : rootInfo.port = hostPort;
336 13 : rootInfo.nicDeploy = NICDeployment::NIC_DEPLOYMENT_DEVICE;
337 :
338 13 : HCCL_INFO("rootInfo: ip[%s] port[%u] identifier[%s]", rootInfo.ip, rootInfo.port, rootInfo.identifier);
339 13 : return HCCL_SUCCESS;
340 13 : }
341 :
342 0 : HcclResult TopoInfoDetect::CalcGroupSizeAndRank(const u32 nRanks, const u32 rank, u32 &groupSize, u32 &groupRank)
343 : {
344 0 : u32 groupIndex = rank / TOPO_MAX_GROUP_SIZE;
345 0 : u32 groupNum = nRanks / TOPO_MAX_GROUP_SIZE;
346 0 : groupSize = groupIndex < groupNum ? TOPO_MAX_GROUP_SIZE : (nRanks - (TOPO_MAX_GROUP_SIZE * groupNum));
347 0 : groupRank = rank == 0 ? 0 : rank % TOPO_MAX_GROUP_SIZE;
348 :
349 0 : return HCCL_SUCCESS;
350 : }
351 :
352 0 : HcclResult TopoInfoDetect::SetupGroupMember(u32 rankSize, u32 myrank, const HcclRootHandle &rootInfo)
353 : {
354 0 : CHK_RET(hrtGetDevice(&deviceLogicID_));
355 :
356 0 : HcclIpAddress rootIP(rootInfo.ip);
357 0 : CHK_PRT_RET(rootIP.IsInvalid(), HCCL_ERROR("string[%s] is invalid ip", rootInfo.ip), HCCL_E_PARA);
358 0 : CHK_RET(hrtGetDevicePhyIdByIndex(deviceLogicID_, devicePhysicID_, true));
359 :
360 0 : CHK_RET(HcclNetInit(NICDeployment::NIC_DEPLOYMENT_HOST, devicePhysicID_, deviceLogicID_, true));
361 :
362 0 : HcclIpAddress hostIP = GetBootstrapHostIP();
363 0 : CHK_RET(GetLocalHostIP(hostIP, devicePhysicID_));
364 :
365 0 : SetBootstrapHostIP(hostIP);
366 :
367 0 : bool bInitDevNic = rankSize != 1 ? true : false;
368 0 : HcclResult ret = StartNetwork(hostIP, bInitDevNic);
369 0 : CHK_PRT_RET(ret != HCCL_SUCCESS,
370 : HCCL_ERROR("[Setup][GroupMember]topo detect GroupMember start network failed! rank[%u]", myrank), ret);
371 :
372 0 : u32 groupSize = 0;
373 0 : u32 groupRank = 0;
374 0 : ret = CalcGroupSizeAndRank(rankSize, myrank, groupSize, groupRank);
375 :
376 0 : ret = GenerateLocalRankInfo(rankSize, myrank, localRankInfo_);
377 0 : CHK_PRT_RET(ret != HCCL_SUCCESS,
378 : HCCL_ERROR("[Setup][GroupMember]topo detect generate local rank info failed! rank[%u]", myrank), ret);
379 :
380 : /* 首节点日志,建链失败属常见问题,在建链前记录相关信息 */
381 0 : HCCL_RUN_INFO("[HCCL_TRACE]SetupGroupMember rankNum[%u], rank[%u], rootInfo identifier[%s], server[%s], "
382 : "deviceType[%d], logicDevId[%d], phydevId[%d], deviceIp[%s]", rankSize, myrank, rootInfo.identifier,
383 : localRankInfo_.hostIP.GetReadableAddress(), localRankInfo_.deviceType, localRankInfo_.deviceLogicID,
384 : localRankInfo_.devicePhysicID, localRankInfo_.deviceIP[0].GetReadableIP()) ;
385 :
386 0 : pTopoExchangeAgent_.reset(new (nothrow) TopoInfoExchangeAgent(rootIP, rootInfo.port,
387 0 : rootInfo.identifier, agentPortCtx_, localRankInfo_, groupSize, groupRank));
388 0 : CHK_SMART_PTR_NULL(pTopoExchangeAgent_);
389 0 : CHK_RET(pTopoExchangeAgent_->SetIsInterSuperPodRetryEnable(isInterSuperPodRetryEnable_));
390 0 : CHK_RET(pTopoExchangeAgent_->SetupMember());
391 0 : CHK_RET(pTopoExchangeAgent_->GetClusterTopoInfo(clusterTopoInfo_));
392 :
393 0 : rootInfo_ = rootInfo;
394 :
395 0 : HCCL_INFO("topo detect completed. myrank[%u], totalranks[%u], myhost[%s], totalservers[%u].",
396 : myrank, rankSize, localRankInfo_.hostIP.GetReadableAddress(), clusterTopoInfo_.serverNum);
397 0 : return HCCL_SUCCESS;
398 0 : }
399 :
400 17 : HcclResult TopoInfoDetect::TeardownServer()
401 : {
402 17 : if(pTopoExchangeServer_) {
403 0 : CHK_RET(pTopoExchangeServer_->Teardown());
404 : }
405 :
406 17 : if (serverPortCtx_) {
407 13 : HcclNetCloseDev(serverPortCtx_);
408 13 : serverPortCtx_ = nullptr;
409 13 : CHK_RET(HcclNetDeInit(NICDeployment::NIC_DEPLOYMENT_HOST, devicePhysicID_, deviceLogicID_));
410 : }
411 17 : HCCL_INFO("TopoInfoDetect TeardownServer ok, identifier[%s].", rootInfo_.identifier);
412 17 : return HCCL_SUCCESS;
413 : }
414 :
415 0 : HcclResult TopoInfoDetect::WaitTopoExchangeServerCompelte(u32 idx) const
416 : {
417 0 : const auto start = chrono::steady_clock::now();
418 0 : const auto timeout = chrono::seconds(GetExternalInputHcclLinkTimeOut());
419 0 : auto iter = g_topoExchangeServerStatus_.Find(idx);
420 0 : if (!iter.second) {
421 0 : return HCCL_SUCCESS;
422 : }
423 0 : u32 status = TOPO_EXCHANGE_SERVER_STATUS_RUNING;
424 : while (true) {
425 0 : auto it = g_topoExchangeServerStatus_.Find(idx);
426 0 : if (it.second) {
427 0 : status = it.first->second;
428 : }
429 0 : if (status == TOPO_EXCHANGE_SERVER_STATUS_ERROR) {
430 0 : HCCL_ERROR("[Wait][TopoExchangeServerCompelte]topo detect failed. topoExchangeServer port[%u] failed.",
431 : idx);
432 0 : return HCCL_E_INTERNAL;
433 0 : } else if (status == TOPO_EXCHANGE_SERVER_STATUS_IDLE) {
434 0 : HCCL_INFO("topoExchangeServer[%u] completed.", idx);
435 0 : return HCCL_SUCCESS;
436 : } else {
437 : const auto elapsed =
438 0 : chrono::duration_cast<chrono::seconds>(chrono::steady_clock::now() - start);
439 0 : if (elapsed > timeout) {
440 0 : HCCL_ERROR("[Wait][TopoExchangeServerCompelte]wait topoExchangeServer[%u] complete timeout[%lld s]",
441 : idx, elapsed);
442 0 : return HCCL_E_TIMEOUT;
443 : }
444 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
445 0 : continue;
446 0 : }
447 0 : };
448 : return HCCL_SUCCESS;
449 : }
450 :
451 0 : HcclResult TopoInfoDetect::PrepareHandle(HcclRankHandle &rankHandle, std::vector<HcclIpAddress> &whitelist)
452 : {
453 0 : CHK_RET(hrtGetDevice(&deviceLogicID_));
454 :
455 0 : if (GetExternalInputHcclEnableWhitelist() == HCCL_WHITELIST_ON) {
456 0 : CHK_RET(ReadHostSocketWhitelist(whitelist));
457 : }
458 0 : HcclIpAddress hostIP = GetBootstrapHostIP();
459 0 : CHK_RET(hrtGetDevicePhyIdByIndex(deviceLogicID_, devicePhysicID_, true));
460 0 : HCCL_INFO("[Setup][hcclIfBasePort]deviceLogicID_[%u], devicePhysicID_[%u]", deviceLogicID_, devicePhysicID_);
461 :
462 : // true代表感知白名单disable配置
463 0 : CHK_RET(HcclNetInit(NICDeployment::NIC_DEPLOYMENT_HOST, devicePhysicID_, deviceLogicID_, true));
464 :
465 0 : CHK_RET(GetRootHostIP(whitelist, hostIP, devicePhysicID_));
466 0 : SetBootstrapHostIP(hostIP);
467 :
468 0 : u32 deviceNum = 0;
469 0 : CHK_RET(hrtGetDeviceCount(&deviceNum));
470 0 : CHK_PRT_RET((static_cast<u32>(deviceLogicID_) >= deviceNum),
471 : HCCL_ERROR("[Setup][GroupLeader]deviceLogicID[%d] is invalid,deviceNum[%d].", deviceLogicID_, deviceNum),
472 : HCCL_E_PARA);
473 :
474 0 : u32 hostPort = HCCL_INVALID_PORT ;
475 0 : CHK_RET(GenerateRootInfo(hostIP, hostPort, devicePhysicID_, rankHandle));
476 :
477 0 : return HCCL_SUCCESS;
478 0 : }
479 :
480 0 : HcclResult TopoInfoDetect::SetupAgent(u32 rankSize, u32 myrank, const HcclRootHandle &rootInfo,
481 : const HcclRankHandle &rankHandle, const CommConfig &commConfig)
482 : {
483 0 : commConfig_ = commConfig;
484 0 : CHK_PRT_RET((rootInfo.nicDeploy == NICDeployment::NIC_DEPLOYMENT_HOST),
485 : HCCL_ERROR("[Setup][Agent]hcclDeviceNicDisable is [%u] when nicDeploy form root is NIC_DEPLOYMENT_HOST",
486 : rootInfo.nicDeploy), HCCL_E_PARA);
487 0 : CHK_RET(hrtGetDevice(&deviceLogicID_));
488 :
489 0 : HcclIpAddress rootIP(rootInfo.ip);
490 0 : CHK_PRT_RET(rootIP.IsInvalid(), HCCL_ERROR("string[%s] is invalid ip", rootInfo.ip), HCCL_E_PARA);
491 0 : CHK_RET(hrtGetDevicePhyIdByIndex(deviceLogicID_, devicePhysicID_, true));
492 :
493 0 : CHK_RET(HcclNetInit(NICDeployment::NIC_DEPLOYMENT_HOST, devicePhysicID_, deviceLogicID_, true));
494 :
495 0 : HcclIpAddress hostIP = GetBootstrapHostIP();
496 0 : CHK_RET(GetLocalHostIP(hostIP, devicePhysicID_));
497 :
498 0 : SetBootstrapHostIP(hostIP);
499 :
500 0 : bool bInitDevNic = rankSize != 1 ? true : false;
501 0 : HcclResult ret = StartNetwork(hostIP, bInitDevNic);
502 0 : CHK_PRT_RET(ret != HCCL_SUCCESS,
503 : HCCL_ERROR("[Setup][Agent]topo detect agent start network failed! rank[%u]", myrank), ret);
504 :
505 0 : ret = GenerateLocalRankInfo(rankSize, myrank, localRankInfo_);
506 0 : CHK_PRT_RET(ret != HCCL_SUCCESS,
507 : HCCL_ERROR("[Setup][Agent]topo detect generate local rank info failed! rank[%u]", myrank), ret);
508 :
509 0 : if (rankSize > TOPO_HIERARCHICAL_ENABLE_THRESHOLD) {
510 : /* 首节点日志,建链失败属常见问题,在建链前记录相关信息 */
511 0 : HCCL_RUN_INFO("[HCCL_TRACE][Hierarchical]SetupAgent rankNum[%u], rank[%u], rootInfo identifier[%s], server[%s], serverPort[%u]"
512 : "deviceType[%d], logicDevId[%d], phydevId[%d], deviceIp[%s]", rankSize, myrank, rootInfo.identifier,
513 : localRankInfo_.hostIP.GetReadableAddress(), rootInfo.port, localRankInfo_.deviceType,
514 : localRankInfo_.deviceLogicID, localRankInfo_.devicePhysicID, localRankInfo_.deviceIP[0].GetReadableIP()) ;
515 :
516 0 : pTopoExchangeAgent_.reset(new (nothrow) TopoInfoExchangeAgent(rootIP, rootInfo.port,
517 0 : rootInfo.identifier, agentPortCtx_, localRankInfo_, rankHandle));
518 0 : CHK_SMART_PTR_NULL(pTopoExchangeAgent_);
519 0 : CHK_RET(pTopoExchangeAgent_->SetIsInterSuperPodRetryEnable(isInterSuperPodRetryEnable_));
520 0 : CHK_RET(pTopoExchangeAgent_->Setup());
521 0 : CHK_RET(pTopoExchangeAgent_->GetGroupLeader(grpLeader_));
522 : } else {
523 : /* 首节点日志,建链失败属常见问题,在建链前记录相关信息 */
524 0 : HCCL_RUN_INFO("[HCCL_TRACE][Flat]SetupAgent rankNum[%u], rank[%u], rootInfo identifier[%s], server[%s], serverPort[%u]"
525 : "deviceType[%d], logicDevId[%d], phydevId[%d], deviceIp[%s]", rankSize, myrank, rootInfo.identifier,
526 : localRankInfo_.hostIP.GetReadableAddress(), rootInfo.port, localRankInfo_.deviceType,
527 : localRankInfo_.deviceLogicID, localRankInfo_.devicePhysicID, localRankInfo_.deviceIP[0].GetReadableIP()) ;
528 :
529 0 : pTopoExchangeAgent_.reset(new (nothrow) TopoInfoExchangeAgent(rootIP, rootInfo.port,
530 0 : rootInfo.identifier, agentPortCtx_, localRankInfo_));
531 0 : CHK_SMART_PTR_NULL(pTopoExchangeAgent_);
532 0 : CHK_RET(pTopoExchangeAgent_->SetIsInterSuperPodRetryEnable(isInterSuperPodRetryEnable_));
533 0 : CHK_RET(pTopoExchangeAgent_->Setup());
534 0 : CHK_RET(pTopoExchangeAgent_->GetClusterTopoInfo(clusterTopoInfo_));
535 : }
536 :
537 0 : rootInfo_ = rootInfo;
538 :
539 0 : HCCL_INFO("topo detect completed. myrank[%u], totalranks[%u], myhost[%s], totalservers[%u].",
540 : myrank, rankSize, localRankInfo_.hostIP.GetReadableAddress(), clusterTopoInfo_.serverNum);
541 0 : return HCCL_SUCCESS;
542 0 : }
543 :
544 0 : HcclResult TopoInfoDetect::SetupRank(std::shared_ptr<HcclSocket> &agentConnRoot) {
545 0 : CHK_RET(pTopoExchangeAgent_->SetupRank(agentConnRoot));
546 0 : CHK_RET(pTopoExchangeAgent_->GetGroupLeader(grpLeader_));
547 0 : return HCCL_SUCCESS;
548 : }
549 :
550 17 : HcclResult TopoInfoDetect::TeardownAgent()
551 : {
552 17 : bool bInitDevNic = clusterTopoInfo_.rankNum != 1 ? true : false;
553 17 : HcclIpAddress hostIP = GetBootstrapHostIP();
554 :
555 17 : if (!pTopoExchangeAgent_) { // 异常处理:如果没有创建agent,对标SetupAgent函数中的网络操作,则直接StopNetwork
556 17 : auto ret = StopNetwork(hostIP, bInitDevNic);
557 17 : CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_ERROR("[Setup][Agent]topo detect agent stop network failed!"), ret);
558 17 : return HCCL_SUCCESS;
559 : }
560 0 : CHK_RET(pTopoExchangeAgent_->Teardown());
561 :
562 0 : auto ret = StopNetwork(hostIP, bInitDevNic);
563 0 : CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_ERROR("[Setup][Agent]topo detect agent stop network failed!"), ret);
564 0 : HCCL_INFO("TopoInfoDetect TeardownAgent ok, identifier[%s].", rootInfo_.identifier);
565 0 : return HCCL_SUCCESS;
566 17 : }
567 :
568 0 : HcclResult TopoInfoDetect::SetupAgentByMasterInfo(HcclIpAddress &localHostIp, const HcclRootHandle &rootInfo)
569 : {
570 0 : CHK_RET(hrtGetDevice(&deviceLogicID_));
571 0 : SetBootstrapHostIP(localHostIp);
572 0 : CHK_RET(hrtGetDevicePhyIdByIndex(deviceLogicID_, devicePhysicID_));
573 0 : rootInfo_ = rootInfo;
574 0 : bool bInitDevNic = GetExternalInputMasterInfo().rankSize != 1 ? true : false;
575 0 : HcclResult ret = StartNetwork(localHostIp, bInitDevNic);
576 0 : CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_ERROR("[Setup][Agent]topo detect agent start network failed!"), ret);
577 :
578 0 : bool errorFlag = false;
579 : do {
580 0 : HcclIpAddress rootIP(rootInfo.ip);
581 0 : CHK_PRT_BREAK(rootIP.IsInvalid(), HCCL_ERROR("[Setup][Agent]string[%s] is invalid ip", rootInfo.ip),
582 : errorFlag = true);
583 0 : ret = GenerateLocalRankInfo(GetExternalInputMasterInfo().rankSize, INVALID_VALUE_RANKID, localRankInfo_);
584 0 : CHK_PRT_BREAK(ret != HCCL_SUCCESS, HCCL_ERROR("[Setup][Agent]topo detect generate local rank info failed"),
585 : errorFlag = true);
586 :
587 0 : pTopoExchangeAgent_.reset(new (nothrow) TopoInfoExchangeAgent(rootIP, rootInfo.port,
588 0 : rootInfo.identifier, agentPortCtx_, localRankInfo_));
589 0 : if (pTopoExchangeAgent_ == nullptr) {
590 0 : HCCL_ERROR("[Setup][Agent]pTopoExchangeAgent is nullptr");
591 0 : errorFlag = true;
592 0 : ret = HCCL_E_PTR;
593 0 : break;
594 : }
595 :
596 0 : CHK_RET(pTopoExchangeAgent_->SetIsInterSuperPodRetryEnable(isInterSuperPodRetryEnable_));
597 0 : ret = pTopoExchangeAgent_->SetupByMasterInfo();
598 0 : CHK_PRT_BREAK(ret != HCCL_SUCCESS, HCCL_ERROR("[Setup][Agent]setup by masterInfo failed"),
599 : errorFlag = true);
600 0 : } while (0);
601 :
602 0 : if (errorFlag) {
603 : // 如果StartNetwork后执行有报错,则先StopNetwork,再返回
604 0 : HcclResult result = StopNetwork(localHostIp, bInitDevNic);
605 0 : CHK_PRT_RET(result != HCCL_SUCCESS, HCCL_ERROR("[Setup][Agent]topo detect agent stop network failed!"), result);
606 :
607 0 : HCCL_ERROR("[Setup][Agent]topo detect agent failed, return[%d]", ret);
608 0 : return ret;
609 : }
610 :
611 0 : CHK_RET(pTopoExchangeAgent_->GetClusterTopoInfo(clusterTopoInfo_));
612 0 : CHK_RET(pTopoExchangeAgent_->GetIdentifier(identifierNum_));
613 :
614 0 : HCCL_INFO("topo detect completed. deviceLogicID[%u] totalranks[%u], myhost[%s], totalservers[%u].",
615 : deviceLogicID_, GetExternalInputMasterInfo().rankSize, localRankInfo_.hostIP.GetReadableAddress(),
616 : clusterTopoInfo_.serverNum);
617 0 : return HCCL_SUCCESS;
618 : }
619 :
620 0 : HcclResult TopoInfoDetect::WaitComplete(const HcclRootHandle &rootInfo)
621 : {
622 0 : return WaitTopoExchangeServerCompelte(rootInfo.port);
623 : }
624 :
625 17 : HcclResult TopoInfoDetect::Teardown()
626 : {
627 17 : CHK_RET(TeardownAgent());
628 17 : CHK_RET(TeardownServer());
629 17 : return HCCL_SUCCESS;
630 : }
631 :
632 3 : HcclResult TopoInfoDetect::ReadHostSocketWhitelist(vector<HcclIpAddress> &whitelist) const
633 : {
634 19 : RPT_ENV_ERR((GetExternalInputHcclWhiteListFile().length() == 0), "EI0001",
635 : std::vector<std::string>({"value", "env", "expect"}),
636 : vector<string>({"", "HCCL_WHITELIST_FILE", "a valid file path" }));
637 :
638 3 : CHK_PRT_RET((GetExternalInputHcclWhiteListFile().length() == 0),
639 : HCCL_ERROR("[%s][%s]environmental variable HCCL_WHITELIST_DISABLE is [0], "
640 : "but HCCL_WHITELIST_FILE is not set or not exist",
641 : LOG_KEYWORDS_INIT_GROUP.c_str(),
642 : LOG_KEYWORDS_ENV_CONFIG.c_str()),
643 : HCCL_E_PARA);
644 :
645 : // 文件路径在处理外部输入时已经做过合法性判断, 无需再次校验
646 : HcclResult ret =
647 2 : HcclWhitelist::GetInstance().LoadConfigFile(GetExternalInputHcclWhiteListFile());
648 :
649 2 : RPT_ENV_ERR(ret != HCCL_SUCCESS, "EI0001",
650 : std::vector<std::string>({"value", "env", "expect"}),
651 : std::vector<std::string>({GetExternalInputHcclWhiteListFile(), "HCCL_WHITELIST_FILE",
652 : "a valid whitelist file format"}));
653 :
654 2 : CHK_PRT_RET(ret != HCCL_SUCCESS,
655 : HCCL_ERROR("[%s][%s]hccl whitelist load config file[%s] failed. ret[%u].",
656 : LOG_KEYWORDS_INIT_GROUP.c_str(),
657 : LOG_KEYWORDS_ENV_CONFIG.c_str(),
658 : GetExternalInputHcclWhiteListFile().c_str(),
659 : ret),
660 : ret);
661 2 : CHK_RET(HcclWhitelist::GetInstance().GetHostWhiteList(whitelist));
662 :
663 2 : CHK_PRT_RET(whitelist.empty(),
664 : HCCL_ERROR("[%s][%s]whitelist file[%s] have no valid host ip.",
665 : LOG_KEYWORDS_INIT_GROUP.c_str(),
666 : LOG_KEYWORDS_ENV_CONFIG.c_str(),
667 : GetExternalInputHcclWhiteListFile().c_str()),
668 : HCCL_E_UNAVAIL);
669 2 : HCCL_INFO("get host socket whitelist success. there are %zu host ip in the whitelist.", whitelist.size());
670 2 : return HCCL_SUCCESS;
671 0 : }
672 :
673 14 : HcclResult TopoInfoDetect::GetAllHostIfInfos(vector<pair<string, HcclIpAddress>> &ifInfos, u32 devPhyId) const
674 : {
675 14 : CHK_RET(hrtGetHostIf(ifInfos, devPhyId));
676 :
677 14 : return HCCL_SUCCESS;
678 : }
679 :
680 2 : HcclResult TopoInfoDetect::GetAllValidHostIfInfos(const vector<HcclIpAddress> &whitelist,
681 : vector<pair<string, HcclIpAddress>> &ifInfos, u32 devPhyId)
682 : {
683 2 : vector<pair<string, HcclIpAddress>> orginIfInfos;
684 2 : CHK_RET(GetAllHostIfInfos(orginIfInfos, devPhyId));
685 :
686 10 : for (auto &ifInfo : orginIfInfos) {
687 8 : auto iter = find(whitelist.begin(), whitelist.end(), ifInfo.second);
688 8 : if (iter != whitelist.end()) {
689 1 : ifInfos.push_back({ ifInfo.first, ifInfo.second });
690 : }
691 : }
692 :
693 2 : return HCCL_SUCCESS;
694 2 : }
695 :
696 14 : HcclResult TopoInfoDetect::GetRootHostIP(const vector<HcclIpAddress> &whitelist, HcclIpAddress &ip, u32 devPhyId)
697 : {
698 14 : if (!ip.IsInvalid()) {
699 0 : return HCCL_SUCCESS;
700 : }
701 14 : vector<pair<string, HcclIpAddress>> ifInfos;
702 :
703 14 : if (GetExternalInputHcclEnableWhitelist() == HCCL_WHITELIST_ON) {
704 2 : CHK_RET(GetAllValidHostIfInfos(whitelist, ifInfos, devPhyId));
705 2 : CHK_PRT_RET(ifInfos.empty(), HCCL_ERROR("[Get][RootHostIP]there is no valid host if in whitelist."),
706 : HCCL_E_NOT_FOUND);
707 : } else {
708 12 : CHK_RET(GetAllHostIfInfos(ifInfos, devPhyId));
709 12 : CHK_PRT_RET(ifInfos.empty(), HCCL_ERROR("[Get][RootHostIP]there is no host if."), HCCL_E_NOT_FOUND);
710 : }
711 :
712 13 : CHK_RET(FindLocalHostIP(ifInfos, ip));
713 13 : return HCCL_SUCCESS;
714 14 : }
715 :
716 0 : HcclResult TopoInfoDetect::GetGroupLeader(HcclRankHandle &rankHandle)
717 : {
718 0 : rankHandle = grpLeader_;
719 0 : return HCCL_SUCCESS;
720 : }
721 :
722 0 : HcclResult TopoInfoDetect::SetIsInterSuperPodRetryEnable(bool isRetry)
723 : {
724 0 : isInterSuperPodRetryEnable_ = isRetry;
725 0 : return HCCL_SUCCESS;
726 : }
727 :
728 13 : HcclResult TopoInfoDetect::StartRootNetwork(const HcclIpAddress& hostIP, u32 &usePort, const std::vector<HcclSocketPortRange> &portRanges)
729 : {
730 13 : CHK_RET(HcclNetOpenDev(&serverPortCtx_, NicType::HOST_NIC_TYPE, devicePhysicID_, deviceLogicID_, hostIP));
731 13 : CHK_PTR_NULL(serverPortCtx_);
732 :
733 13 : if (usePort == HCCL_INVALID_PORT) {
734 : // 通过抢占的方式获得Root节点监听的host端口
735 13 : listenSocket_.reset(new (nothrow) HcclSocket(serverPortCtx_));
736 13 : CHK_SMART_PTR_NULL(listenSocket_);
737 13 : CHK_RET(listenSocket_->Init());
738 13 : HcclResult ret = PreemptPortManager::GetInstance(deviceLogicID_).ListenPreempt(listenSocket_,
739 : portRanges, usePort);
740 13 : CHK_PRT_RET(ret != HCCL_SUCCESS,
741 : HCCL_ERROR("[TopoInfoDetect][StartRootNetwork] devPhyId[%u], devLogicId[%u], host ip[%s], "
742 : "try to preempt port on host nic fail.",
743 : devicePhysicID_, deviceLogicID_, hostIP.GetReadableAddress()), ret);
744 : } else {
745 : // 1. 使用MasterInfo初始化时,不支持抢占master节点的监听端口
746 : // 2. 未配置port range时,不支持抢占监听端口
747 0 : listenSocket_.reset(new (nothrow) HcclSocket(serverPortCtx_, usePort));
748 0 : CHK_SMART_PTR_NULL(listenSocket_);
749 0 : CHK_RET(listenSocket_->Init());
750 0 : HcclResult ret = listenSocket_->Listen();
751 0 : CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_ERROR("[%s][%s]%s failed, ret[%u]",
752 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), __func__, ret), ret);
753 : }
754 :
755 13 : HCCL_INFO("topo info exchange server start with host ip[%s] and port[%u]", hostIP.GetReadableAddress(), usePort);
756 :
757 13 : return HCCL_SUCCESS;
758 : }
759 :
760 3 : HcclResult TopoInfoDetect::StartGroupLeaderNetwork(const vector<HcclIpAddress> &whitelist, const HcclIpAddress& hostIP,
761 : u32 &bindPort, const std::vector<HcclSocketPortRange> &portRanges)
762 : {
763 3 : CHK_RET(HcclNetOpenDev(&serverPortCtx_, NicType::HOST_NIC_TYPE, devicePhysicID_, deviceLogicID_, hostIP));
764 3 : CHK_PTR_NULL(serverPortCtx_);
765 3 : if (bindPort == HCCL_INVALID_PORT) {
766 : // 通过抢占的方式获得GroupLeader节点监听的host端口
767 3 : listenSocket_.reset(new (nothrow) HcclSocket(serverPortCtx_));
768 3 : CHK_SMART_PTR_NULL(listenSocket_);
769 3 : CHK_RET(listenSocket_->Init());
770 3 : HcclResult ret = PreemptPortManager::GetInstance(deviceLogicID_).ListenPreempt(listenSocket_,
771 : portRanges, bindPort);
772 3 : CHK_PRT_RET(ret != HCCL_SUCCESS,
773 : HCCL_ERROR("[TopoInfoDetect][StartGroupLeaderNetwork] devPhyId[%u], devLogicId[%u], host ip[%s], "
774 : "try to preempt port on host nic fail.",
775 : devicePhysicID_, deviceLogicID_, hostIP.GetReadableAddress()), ret);
776 : } else {
777 : // 1. 使用MasterInfo初始化时,不支持抢占master节点的监听端口
778 : // 2. 未配置port range时,不支持抢占监听端口
779 0 : listenSocket_.reset(new (nothrow) HcclSocket(serverPortCtx_, bindPort));
780 0 : CHK_SMART_PTR_NULL(listenSocket_);
781 0 : CHK_RET(listenSocket_->Init());
782 0 : HcclResult ret = listenSocket_->Listen();
783 0 : CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_ERROR("[%s][%s]%s failed, ret[%u]",
784 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), __func__, ret), ret);
785 : }
786 :
787 3 : HCCL_INFO("group leader start with host ip[%s] and port[%u]", hostIP.GetReadableAddress(), bindPort);
788 :
789 3 : if (GetExternalInputHcclEnableWhitelist() == HCCL_WHITELIST_ON) {
790 0 : CHK_RET(AddSocketWhiteList(bindPort, whitelist));
791 : }
792 :
793 3 : return HCCL_SUCCESS;
794 : }
795 :
796 1 : HcclResult TopoInfoDetect::AddSocketWhiteList(u32 port,
797 : const vector<HcclIpAddress> &whitelist) const
798 : {
799 1 : vector<SocketWlistInfo> wlistInfosVec;
800 2 : for (auto ip : whitelist) {
801 : SocketWlistInfo wlistInfo;
802 1 : wlistInfo.connLimit = HOST_SOCKET_CONN_LIMIT;
803 1 : wlistInfo.remoteIp.addr = ip.GetBinaryAddress().addr;
804 1 : wlistInfo.remoteIp.addr6 = ip.GetBinaryAddress().addr6;
805 1 : string tag = TOPO_DETECT_TAG + "_" + rootInfo_.identifier + "_" + to_string(port);
806 1 : s32 sRet = memcpy_s(&wlistInfo.tag[0], sizeof(wlistInfo.tag), tag.c_str(), tag.size() + 1);
807 1 : if (sRet != EOK) {
808 0 : HCCL_ERROR("[Add][SocketWhiteList]memory copy failed. errorno[%d]", sRet);
809 0 : return HCCL_E_MEMORY;
810 : }
811 1 : wlistInfosVec.push_back(wlistInfo);
812 1 : }
813 :
814 1 : CHK_RET(listenSocket_->AddWhiteList(wlistInfosVec));
815 :
816 1 : HCCL_INFO("add socket white list success. total: %zu", whitelist.size());
817 1 : return HCCL_SUCCESS;
818 1 : }
819 :
820 0 : HcclResult TopoInfoDetect::StartNetwork(HcclIpAddress &hostIP, bool bInitDevNic)
821 : {
822 0 : CHK_RET(HcclNetOpenDev(&agentPortCtx_, NicType::HOST_NIC_TYPE, devicePhysicID_, deviceLogicID_, hostIP));
823 0 : CHK_PTR_NULL(agentPortCtx_);
824 :
825 0 : if (bInitDevNic) {
826 0 : CHK_RET(HcclNetInit(NICDeployment::NIC_DEPLOYMENT_DEVICE, devicePhysicID_, deviceLogicID_, false));
827 0 : CHK_RET(
828 : HcclNetOpenDev(&devNicCtx_, NicType::DEVICE_NIC_TYPE, devicePhysicID_, deviceLogicID_, HcclIpAddress(0)));
829 0 : CHK_PTR_NULL(devNicCtx_);
830 : }
831 :
832 0 : HCCL_INFO("NetworkManager start host net success! ip[%s]", hostIP.GetReadableAddress());
833 :
834 0 : return HCCL_SUCCESS;
835 : }
836 :
837 17 : HcclResult TopoInfoDetect::StopNetwork(HcclIpAddress &hostIP, bool bInitDevNic)
838 : {
839 17 : if (agentPortCtx_) {
840 0 : HcclNetCloseDev(agentPortCtx_);
841 0 : agentPortCtx_ = nullptr;
842 0 : CHK_RET(HcclNetDeInit(NICDeployment::NIC_DEPLOYMENT_HOST, devicePhysicID_, deviceLogicID_));
843 : }
844 :
845 17 : if (bInitDevNic) {
846 17 : if (devNicCtx_) {
847 0 : HcclNetCloseDev(devNicCtx_);
848 0 : devNicCtx_ = nullptr;
849 0 : CHK_RET(HcclNetDeInit(NICDeployment::NIC_DEPLOYMENT_DEVICE, devicePhysicID_, deviceLogicID_));
850 : }
851 : }
852 :
853 17 : HCCL_INFO("NetworkManager stop host net success! ip[%s] ", hostIP.GetReadableAddress());
854 17 : return HCCL_SUCCESS;
855 : }
856 :
857 3 : HcclResult TopoInfoDetect::FilterDevIPs(std::vector<HcclIpAddress> &sourceDeviceIPs,
858 : std::vector<HcclIpAddress> &targetDeviceIPs) const
859 : {
860 3 : std::vector<HcclIpAddress> deviceIPv4;
861 3 : std::vector<HcclIpAddress> deviceIPv6;
862 6 : for (auto &iter : sourceDeviceIPs) {
863 3 : if (iter.IsIPv6()) {
864 3 : deviceIPv6.push_back(iter);
865 : } else {
866 0 : deviceIPv4.push_back(iter);
867 : }
868 : }
869 : // 同时存在ipv4/ipv6时,除非指定socket family,否则ipv4优先
870 : // 只存在ipv4/ipv6单栈时,不受用户指定的socket family约束
871 3 : if ((((GetExternalInputHcclSocketFamily() == -1) ||
872 0 : (GetExternalInputHcclSocketFamily() == AF_INET)) &&
873 3 : (!deviceIPv4.empty())) || deviceIPv6.empty()) {
874 0 : targetDeviceIPs = deviceIPv4;
875 0 : HCCL_RUN_INFO("select AF_INET family as device socket family.");
876 3 : } else if (!deviceIPv6.empty()) {
877 3 : std::sort(deviceIPv6.begin(), deviceIPv6.end());
878 3 : targetDeviceIPs.push_back(deviceIPv6[0]);
879 3 : HCCL_RUN_INFO("select AF_INET6 family as device socket family.");
880 : }
881 3 : return HCCL_SUCCESS;
882 3 : }
883 :
884 0 : HcclResult TopoInfoDetect::PreemptDeviceNicPort(const u32 devPhyId, const s32 devLogicId,
885 : const HcclIpAddress &deviceIp, u32 &usePort)
886 : {
887 0 : HcclIpAddress devIp(std::string(deviceIp.GetReadableIP()));
888 0 : HcclNetDevCtx netCtx{nullptr};
889 0 : CHK_RET(HcclNetInit(NICDeployment::NIC_DEPLOYMENT_DEVICE, devPhyId, devLogicId, false, false));
890 0 : CHK_RET(HcclNetOpenDev(&netCtx, NicType::DEVICE_NIC_TYPE, devPhyId, devLogicId, devIp));
891 0 : CHK_PTR_NULL(netCtx);
892 0 : commPortConfig_.devNicListen = std::make_pair(nullptr, netCtx);
893 :
894 0 : commPortConfig_.devNicListen.first.reset(new (std::nothrow) HcclSocket(netCtx));
895 0 : CHK_SMART_PTR_NULL(commPortConfig_.devNicListen.first);
896 0 : CHK_RET(commPortConfig_.devNicListen.first->Init());
897 :
898 0 : HcclResult ret = PreemptPortManager::GetInstance(devLogicId).ListenPreempt(commPortConfig_.devNicListen.first,
899 : GetExternalInputNpuSocketPortRange(), usePort);
900 0 : CHK_PRT_RET(ret != HCCL_SUCCESS,
901 : HCCL_ERROR("[TopoInfoDetect][PreemptDeviceNicPort] devPhyId[%u], devLogicId[%u], device ip[%s], "
902 : "try to preempt port on device nic fail.", devPhyId, devLogicId, devIp.GetReadableAddress()), ret);
903 :
904 0 : HCCL_INFO("[TopoInfoDetect][PreemptDeviceNicPort]devPhyId[%u], devLogicId[%d], "
905 : "preempt port[%u] on ip[%s] success.", devPhyId, devLogicId, usePort, devIp.GetReadableAddress());
906 0 : return HCCL_SUCCESS;
907 0 : }
908 :
909 0 : HcclResult TopoInfoDetect::PreemptDeviceVnicPort(HcclBasicRankInfo &localRankInfo)
910 : {
911 0 : u32 devPhyId = localRankInfo.devicePhysicID;
912 0 : s32 devLogicId = localRankInfo.deviceLogicID;
913 0 : HcclIpAddress vnicIp(devPhyId);
914 0 : bool useSuperPodMode = false;
915 0 : if (localRankInfo.superDeviceId != INVALID_UINT) {
916 0 : CHK_RET(IsSuperPodMode(useSuperPodMode));
917 : }
918 0 : if (useSuperPodMode) {
919 0 : CHK_RET(hrtRaGetSingleSocketVnicIpInfo(
920 : devPhyId, DeviceIdType::DEVICE_ID_TYPE_SDID, localRankInfo.superDeviceId, vnicIp));
921 : } else {
922 0 : CHK_RET(hrtRaGetSingleSocketVnicIpInfo(
923 : devPhyId, DeviceIdType::DEVICE_ID_TYPE_PHY_ID, devPhyId, vnicIp));
924 : }
925 :
926 0 : HCCL_INFO("[TopoInfoDetect][PreemptDeviceVnicPort] vnicIp is [%s]", vnicIp.GetReadableAddress());
927 :
928 0 : HcclNetDevCtx netCtx{nullptr};
929 0 : CHK_RET(HcclNetInit(NICDeployment::NIC_DEPLOYMENT_DEVICE, devPhyId, devLogicId, false, false));
930 0 : CHK_RET(HcclNetOpenDev(&netCtx, NicType::VNIC_TYPE, devPhyId, devLogicId, vnicIp));
931 0 : CHK_PTR_NULL(netCtx);
932 0 : commPortConfig_.devVnicListen = std::make_pair(nullptr, netCtx);
933 :
934 0 : commPortConfig_.devVnicListen.first.reset(new (std::nothrow) HcclSocket(netCtx));
935 0 : CHK_SMART_PTR_NULL(commPortConfig_.devVnicListen.first);
936 0 : CHK_RET(commPortConfig_.devVnicListen.first->Init());
937 :
938 0 : HcclResult ret = PreemptPortManager::GetInstance(devLogicId).ListenPreempt(commPortConfig_.devVnicListen.first,
939 0 : GetExternalInputNpuSocketPortRange(), localRankInfo.deviceVnicPort);
940 0 : CHK_PRT_RET(ret != HCCL_SUCCESS,
941 : HCCL_ERROR("[TopoInfoDetect][PreemptDeviceVnicPort] devPhyId[%u], devLogicId[%u], vnicIp[%s], "
942 : "try to preempt port on vnic fail.", devPhyId, devLogicId, vnicIp.GetReadableAddress()), ret);
943 :
944 0 : HCCL_INFO("[TopoInfoDetect][PreemptDeviceVnicPort] devPhyId[%u], devLogicId[%d], "
945 : "preempt vnic on ip[%s], port[%u] success.",
946 : devPhyId, devLogicId, vnicIp.GetReadableAddress(), localRankInfo.deviceVnicPort);
947 0 : return HCCL_SUCCESS;
948 0 : }
949 :
950 0 : HcclResult TopoInfoDetect::PreemptBackupDeviceNicPort(const u32 devPhyId, const s32 devLogicId,
951 : const HcclIpAddress &deviceIp, const HcclIpAddress &backupDeviceIp, u32 &usePort)
952 : {
953 0 : HcclIpAddress devIp(std::string(deviceIp.GetReadableIP()));
954 0 : HcclIpAddress backupDevIp(std::string(backupDeviceIp.GetReadableIP()));
955 0 : HcclNetDevCtx netCtx{nullptr};
956 0 : CHK_RET(HcclNetInit(NICDeployment::NIC_DEPLOYMENT_DEVICE, devPhyId, devLogicId, false, true));
957 0 : CHK_RET(HcclNetOpenDev(&netCtx, NicType::DEVICE_NIC_TYPE, devPhyId, devLogicId, backupDevIp, devIp));
958 0 : CHK_PTR_NULL(netCtx);
959 0 : commPortConfig_.backupDevNicListen = std::make_pair(nullptr, netCtx);
960 :
961 0 : commPortConfig_.backupDevNicListen.first.reset(new (std::nothrow) HcclSocket(netCtx));
962 0 : CHK_SMART_PTR_NULL(commPortConfig_.backupDevNicListen.first);
963 0 : CHK_RET(commPortConfig_.backupDevNicListen.first->Init());
964 :
965 0 : HcclResult ret = PreemptPortManager::GetInstance(devLogicId).ListenPreempt(commPortConfig_.backupDevNicListen.first,
966 : GetExternalInputNpuSocketPortRange(), usePort);
967 0 : CHK_PRT_RET(ret != HCCL_SUCCESS,
968 : HCCL_ERROR("[TopoInfoDetect][PreemptBackupDeviceNicPort] devPhyId[%u], devLogicId[%u], device ip[%s], "
969 : "backup device ip[%s], try to preempt port on device nic fail.",
970 : devPhyId, devLogicId, devIp.GetReadableAddress(), backupDevIp.GetReadableAddress()), ret);
971 :
972 0 : HCCL_INFO("[TopoInfoDetect][PreemptBackupDeviceNicPort]devPhyId[%u], devLogicId[%d], local ip[%s]"
973 : "preempt port[%u] on backup device ip[%s] success[%u].",
974 : devPhyId, devLogicId, devIp.GetReadableAddress(), usePort, backupDevIp.GetReadableAddress());
975 0 : return HCCL_SUCCESS;
976 0 : }
977 :
978 2 : HcclResult TopoInfoDetect::GetDeviceBackupNicInfo(HcclBasicRankInfo &localRankInfo)
979 : {
980 2 : std::vector<std::vector<HcclIpAddress>> chipDeviceIPs;
981 2 : CHK_RET(hrtRaGetDeviceAllNicIP(chipDeviceIPs));
982 2 : if (chipDeviceIPs.size() != 2U) {
983 : // 910A3场景一个chip上有两组deviceIP
984 1 : HCCL_RUN_WARNING("[TopoInfoDetect][GetDeviceBackupNicInfo]Fail to load backup device ip!"
985 : "Please check the driver version!");
986 : } else {
987 : // 取到一组backup ip,按照devPhyId排序,ipv4在前,ipv6在后,每个网卡的ip顺序一一对应
988 : // 取其中对端网卡的ip作为备用网卡ip
989 1 : u32 ipIdex = 1U - (localRankInfo.devicePhysicID % 2U);
990 1 : CHK_RET(FilterDevIPs(chipDeviceIPs[ipIdex], localRankInfo.backupDeviceIP));
991 1 : HCCL_INFO("[TopoInfoDetect][GetDeviceBackupNicInfo]devicePhysicID[%u], backupDeviceIP[0]:[%s]",
992 : localRankInfo.devicePhysicID, localRankInfo.backupDeviceIP[0].GetReadableAddress());
993 : // 开启device侧端口配置,并且存在备用网卡时,抢占一个备用网卡上的端口
994 1 : if (GetExternalInputNpuPortSwitch() && localRankInfo.backupDeviceIP.size() > 0) {
995 0 : u32 backupDevPhyId = INVALID_INT;
996 0 : u32 backupDevLogicId = INVALID_INT;
997 0 : CHK_RET(hrtGetPairDevicePhyId(localRankInfo.devicePhysicID, backupDevPhyId));
998 0 : CHK_RET(hrtGetDeviceIndexByPhyId(backupDevPhyId, backupDevLogicId));
999 0 : CHK_RET(PreemptBackupDeviceNicPort(backupDevPhyId, backupDevLogicId, localRankInfo.deviceIP[0],
1000 : localRankInfo.backupDeviceIP[0], localRankInfo.backupDevicePort));
1001 : }
1002 : }
1003 2 : return HCCL_SUCCESS;
1004 2 : }
1005 :
1006 2 : HcclResult TopoInfoDetect::GenerateLocalRankInfo(u32 rankSize, u32 rankID, HcclBasicRankInfo &localRankInfo)
1007 : {
1008 2 : localRankInfo.hostIP = GetBootstrapHostIP();
1009 2 : localRankInfo.rank = rankID;
1010 2 : localRankInfo.rankSize = rankSize;
1011 2 : localRankInfo.nicDeploy = NICDeployment::NIC_DEPLOYMENT_DEVICE;
1012 :
1013 2 : if (devNicCtx_ != nullptr) {
1014 0 : HcclResult ret = HcclNetDevGetTlsStatus(devNicCtx_, &localRankInfo.tlsStatus);
1015 0 : if (ret != HCCL_SUCCESS && ret != HCCL_E_NOT_SUPPORT) {
1016 0 : HCCL_RUN_WARNING("[GenerateLocalRankInfo] HcclNetDevGetTlsStatus failed ret[%u]", ret);
1017 : }
1018 : }
1019 :
1020 2 : CHK_RET(hrtGetDeviceType(localRankInfo.deviceType));
1021 2 : CHK_RET(hrtGetDevice(reinterpret_cast<s32 *>(&localRankInfo.deviceLogicID)));
1022 2 : CHK_RET(hrtGetDevicePhyIdByIndex(static_cast<u32>(localRankInfo.deviceLogicID), localRankInfo.devicePhysicID));
1023 :
1024 2 : if (localRankInfo.deviceType == DevType::DEV_TYPE_910_93) {
1025 2 : CHK_RET(GetSuperPodInfo(localRankInfo.deviceLogicID, localRankInfo.superPodId, localRankInfo.superDeviceId));
1026 : }
1027 :
1028 2 : localRankInfo.deviceIP.clear();
1029 2 : if (localRankInfo.nicDeploy == NICDeployment::NIC_DEPLOYMENT_DEVICE && rankSize != 1) {
1030 2 : std::vector<HcclIpAddress> deviceIPs;
1031 2 : CHK_RET(hrtRaGetDeviceIP(localRankInfo.devicePhysicID, deviceIPs));
1032 2 : CHK_RET(FilterDevIPs(deviceIPs, localRankInfo.deviceIP));
1033 : // 开启device侧端口配置时,需要抢占监听端口
1034 2 : if (GetExternalInputNpuPortSwitch()) {
1035 : // 如果有device nic,则抢占device nic的port
1036 0 : if (localRankInfo.deviceIP.size() > 0) {
1037 0 : CHK_RET(PreemptDeviceNicPort(localRankInfo.devicePhysicID, localRankInfo.deviceLogicID,
1038 : localRankInfo.deviceIP[0], localRankInfo.deviceNicPort));
1039 : }
1040 : // 使用device网卡时,必定抢占vnic上的port
1041 0 : CHK_RET(PreemptDeviceVnicPort(localRankInfo));
1042 0 : commPortConfig_.devPortSwitchOn = true;
1043 : }
1044 :
1045 : // 此处不知道拓扑形态,无法判断是否需要backupIp,只能从硬件类型和重执行开关判断一下
1046 2 : bool useSuperPodMode = false;
1047 2 : CHK_RET(IsSuperPodMode(useSuperPodMode));
1048 2 : if (useSuperPodMode && commConfig_.GetConfigAicpuUnfold() && isInterSuperPodRetryEnable_) {
1049 2 : CHK_RET(GetDeviceBackupNicInfo(localRankInfo));
1050 : }
1051 2 : }
1052 :
1053 2 : if (localRankInfo.deviceIP.size() == 0) {
1054 : // 和 rank table 保持一致,如果没有device网卡时,默认填充 0。
1055 0 : HcclIpAddress invalidAddr;
1056 0 : localRankInfo.deviceIP.push_back(invalidAddr);
1057 0 : HCCL_RUN_INFO("no device ip: use 0 as device ip.");
1058 0 : }
1059 2 : if (localRankInfo.backupDeviceIP.size() == 0) {
1060 : // 如果没有 backup device ip 时,默认填充 0。
1061 1 : HcclIpAddress invalidAddr;
1062 1 : localRankInfo.backupDeviceIP.push_back(invalidAddr);
1063 1 : HCCL_RUN_INFO("no backup device ip: use 0 as device ip.");
1064 1 : }
1065 2 : return HCCL_SUCCESS;
1066 : }
1067 :
1068 2 : HcclResult TopoInfoDetect::GetSuperPodInfo(s32 deviceLogicId, std::string &superPodId, u32 &superDeviceId)
1069 : {
1070 : // 解析super_pod_id
1071 2 : superPodId = GetExternalInputLogicSuperPodId(); // 逻辑super pod id
1072 2 : if (superPodId.empty()) {
1073 2 : s64 val = 0;
1074 2 : CHK_RET(hrtGetDeviceInfo(deviceLogicId, HcclRtDeviceModuleType::HCCL_RT_MODULE_TYPE_SYSTEM,
1075 : HcclRtDeviceInfoType::HCCL_INFO_TYPE_SUPER_POD_ID, val));
1076 2 : superPodId = std::to_string(val); // 真实super pod id
1077 : }
1078 :
1079 : // 解析sdid
1080 2 : s64 sdid = 0;
1081 2 : CHK_RET(hrtGetDeviceInfo(deviceLogicId, HcclRtDeviceModuleType::HCCL_RT_MODULE_TYPE_SYSTEM,
1082 : HcclRtDeviceInfoType::HCCL_INFO_TYPE_SDID, sdid));
1083 2 : superDeviceId = static_cast<u32>(sdid);
1084 2 : HCCL_INFO("[Get][SuperPodInfo]deviceLogicID[%d], superPodId[%s], superDeviceId[%u]",
1085 : deviceLogicId, superPodId.c_str(), superDeviceId);
1086 2 : return HCCL_SUCCESS;
1087 : }
1088 :
1089 0 : HcclResult TopoInfoDetect::GetCluterInfo(RankTable_t &clusterInfo)
1090 : {
1091 0 : CHK_PRT_RET((clusterTopoInfo_.rankList.size() == 0),
1092 : HCCL_ERROR("[Get][CluterInfo]GetCluterInfo failed, topo detect has not started."), HCCL_E_INTERNAL);
1093 0 : clusterInfo = clusterTopoInfo_;
1094 0 : return HCCL_SUCCESS;
1095 : }
1096 0 : HcclResult TopoInfoDetect::GetRankId(u32 &rankId)
1097 : {
1098 0 : rankId = identifierNum_;
1099 0 : return HCCL_SUCCESS;
1100 : }
1101 :
1102 0 : HcclResult TopoInfoDetect::GetLocalRankInfo(HcclBasicRankInfo &rankInfo)
1103 : {
1104 0 : CHK_PRT_RET((localRankInfo_.rankSize == 0), HCCL_ERROR("[Get][LocalRankInfo]GetLocalRankInfo failed, topo "\
1105 : "detect has not started."), HCCL_E_INTERNAL);
1106 0 : rankInfo = localRankInfo_;
1107 0 : return HCCL_SUCCESS;
1108 : }
1109 :
1110 16 : void TopoInfoDetect::SetBootstrapHostIP(HcclIpAddress& ip)
1111 : {
1112 16 : bootstrapHostIP_ = ip;
1113 16 : }
1114 :
1115 36 : HcclIpAddress TopoInfoDetect::GetBootstrapHostIP() const
1116 : {
1117 36 : return bootstrapHostIP_;
1118 : }
1119 0 : HcclResult TopoInfoDetect::TransformRankTableStr(const RankTable_t &clusterInfo, string &ranktableStr)
1120 : {
1121 0 : nlohmann::json basicJson;
1122 0 : HcclResult ret = Struct2JsonRankTable(clusterInfo, basicJson);
1123 0 : CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_WARNING("cluster info to json failed ,ret[%d]", ret), HCCL_E_INTERNAL);
1124 0 : ranktableStr = basicJson.dump(2); // dump参数为2
1125 0 : return HCCL_SUCCESS;
1126 0 : }
1127 0 : HcclResult TopoInfoDetect::TransformDeviceList(const RankTable_t &clusterInfo,
1128 : vector<RankInfo_t> &tmpRankList, nlohmann::json &perServerJson, u32 serverIndex)
1129 : {
1130 0 : for (auto it = tmpRankList.begin(); it != tmpRankList.end();) {
1131 0 : if (it->serverId == clusterInfo.serverList[serverIndex].serverId) {
1132 0 : nlohmann::json perDeviceJson;
1133 0 : perDeviceJson[PROP_DEV_ID] = to_string(it->deviceInfo.devicePhyId);
1134 0 : perDeviceJson[PROP_RANK_ID] = to_string(it->rankId);
1135 0 : perDeviceJson[PROP_SUPER_DEVICE_ID] = to_string(it->superDeviceId);
1136 0 : if (it->deviceInfo.port != HCCL_INVALID_PORT) {
1137 0 : perDeviceJson[PROP_DEV_NIC_PORT] = to_string(it->deviceInfo.port);
1138 : }
1139 0 : if (it->deviceInfo.vnicPort != HCCL_INVALID_PORT) {
1140 0 : perDeviceJson[PROP_DEV_VNIC_PORT] = to_string(it->deviceInfo.vnicPort);
1141 : }
1142 0 : if (it->deviceInfo.backupPort != HCCL_INVALID_PORT) {
1143 0 : perDeviceJson[PROP_BACKUP_DEV_PORT] = to_string(it->deviceInfo.backupPort);
1144 : }
1145 0 : if (clusterInfo.nicDeploy == NICDeployment::NIC_DEPLOYMENT_DEVICE && it->deviceInfo.deviceIp.size() != 0 &&
1146 0 : !it->deviceInfo.deviceIp[0].IsInvalid()) {
1147 0 : perDeviceJson[PROP_DEV_IP] = std::string(it->deviceInfo.deviceIp[0].GetReadableIP());
1148 : }
1149 0 : if (clusterInfo.nicDeploy == NICDeployment::NIC_DEPLOYMENT_DEVICE &&
1150 0 : it->deviceInfo.backupDeviceIp.size() != 0 && !it->deviceInfo.backupDeviceIp[0].IsInvalid()) {
1151 0 : perDeviceJson[PROP_BACKUP_DEV_IP] = std::string(it->deviceInfo.backupDeviceIp[0].GetReadableIP());
1152 : }
1153 0 : if (!it->hostIp.IsInvalid()) {
1154 0 : perServerJson[PROP_HOST_IP] = std::string(it->hostIp.GetReadableIP());
1155 : }
1156 0 : perServerJson[PROP_DEVICE].push_back(perDeviceJson);
1157 0 : it = tmpRankList.erase(it);
1158 0 : } else {
1159 0 : it++;
1160 : }
1161 : }
1162 0 : return HCCL_SUCCESS;
1163 : }
1164 0 : HcclResult TopoInfoDetect::Struct2JsonRankTable(const RankTable_t &clusterInfo, nlohmann::json& ClusterJson)
1165 : {
1166 0 : nlohmann::json serverListJson;
1167 0 : ClusterJson[PROP_SERVER_COUNT] = to_string(clusterInfo.serverNum);
1168 0 : vector<RankInfo_t> tmpRankList = clusterInfo.rankList;
1169 0 : ClusterJson[PROP_SERVER_LIST] = serverListJson;
1170 0 : for (u32 i = 0; i < clusterInfo.serverNum; i++) {
1171 0 : nlohmann::json perServerJson;
1172 0 : perServerJson[PROP_SERVER_ID] = clusterInfo.serverList[i].serverId;
1173 0 : nlohmann::json deviceList;
1174 0 : perServerJson[PROP_DEVICE] = deviceList;
1175 0 : CHK_RET(TransformDeviceList(clusterInfo, tmpRankList, perServerJson, i));
1176 0 : ClusterJson[PROP_SERVER_LIST].push_back(perServerJson);
1177 0 : }
1178 :
1179 0 : nlohmann::json superPodListJson;
1180 0 : CHK_RET(TransformSuperPodList(clusterInfo.rankList, superPodListJson));
1181 0 : ClusterJson[PROP_SUPER_POD_LIST] = superPodListJson;
1182 :
1183 0 : ClusterJson[PROP_STATUS] = "completed";
1184 0 : ClusterJson[PROP_VERSION] = (localRankInfo_.deviceType == DevType::DEV_TYPE_910_93) ? "1.2" : "1.0";
1185 0 : return HCCL_SUCCESS;
1186 0 : }
1187 :
1188 0 : HcclResult TopoInfoDetect::TransformSuperPodList(const std::vector<RankInfo_t> &rankInfo,
1189 : nlohmann::json &superPodListJson) const
1190 : {
1191 : // 按照 <super_pod_id, <server_id>> 格式从RankInfo_t中解析super pod信息
1192 0 : std::map<std::string, std::set<std::string>> superPodMap;
1193 0 : for (u32 i = 0; i < rankInfo.size(); i++) {
1194 0 : auto iter = superPodMap.find(rankInfo[i].superPodId);
1195 0 : if (iter == superPodMap.end()) {
1196 0 : std::set<std::string> perSuperPod;
1197 0 : perSuperPod.insert(rankInfo[i].serverId);
1198 0 : superPodMap.insert(std::pair<std::string, std::set<string>>(rankInfo[i].superPodId, perSuperPod));
1199 0 : } else {
1200 : // superDeviceId在VerifyClusterSuperPodInfo中已经查重校验过
1201 0 : iter->second.insert(rankInfo[i].serverId);
1202 : }
1203 : }
1204 :
1205 0 : for (auto it = superPodMap.begin(); it != superPodMap.end(); ++it) {
1206 0 : nlohmann::json superPodIdJson;
1207 0 : superPodIdJson[PROP_SUPER_POD_ID] = it->first;
1208 0 : nlohmann::json serverListJson;
1209 0 : for (auto perServer = it->second.begin(); perServer != it->second.end(); ++perServer) {
1210 0 : nlohmann::json perServerJson;
1211 0 : perServerJson[PROP_SERVER_ID] = *perServer;
1212 0 : serverListJson.push_back(perServerJson);
1213 0 : }
1214 0 : superPodIdJson[PROP_SERVER_LIST] = serverListJson;
1215 0 : superPodListJson.push_back(superPodIdJson);
1216 0 : }
1217 0 : return HCCL_SUCCESS;
1218 0 : }
1219 : } // namespace hccl
|