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