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