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