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_exchange_base.h"
12 : #include <thread>
13 : #include <iostream>
14 : #include <fstream>
15 : #include "externalinput_pub.h"
16 : #include "mem_host_pub.h"
17 : #include "json_utils.h"
18 :
19 : namespace hccl {
20 :
21 : std::atomic<BroadcastStage> g_broadcastStage(BroadcastStage::Idle);
22 : std::mutex g_broadcast_stage_mutex;
23 : std::condition_variable g_broadcast_stage_cv;
24 :
25 38 : TopoInfoExchangeBase::TopoInfoExchangeBase() : currentStep_(0) {}
26 :
27 38 : TopoInfoExchangeBase::~TopoInfoExchangeBase() {}
28 :
29 21 : HcclResult TopoInfoExchangeBase::DisconnectSocket(std::shared_ptr<HcclSocket> socket) const
30 : {
31 21 : if (socket) {
32 0 : socket->Close();
33 : }
34 21 : return HCCL_SUCCESS;
35 : }
36 :
37 0 : HcclResult TopoInfoExchangeBase::SendClusterInfoMsg(
38 : std::shared_ptr<HcclSocket> socket, const RankTable_t& clusterInfo, const std::string buffer, const u32 msgLen)
39 : {
40 0 : HcclResult ret = socket->Send(&msgLen, sizeof(msgLen));
41 0 : CHK_PRT_RET(
42 : ret != HCCL_SUCCESS,
43 : HCCL_ERROR(
44 : "[Send][ClusterInfoMsg]errNo[0x%016llx] ra send msg length failed! "
45 : "msgLen[%u], ret[%u]",
46 : HCCL_ERROR_CODE(HCCL_E_TCP_TRANSFER), msgLen, ret),
47 : ret);
48 :
49 0 : ret = socket->Send(buffer.c_str(), msgLen);
50 0 : CHK_PRT_RET(
51 : ret != HCCL_SUCCESS,
52 : HCCL_ERROR(
53 : "[Send][ClusterInfoMsg]errNo[0x%016llx] ra send failed! size[%u], ret[%u]",
54 : HCCL_ERROR_CODE(HCCL_E_TCP_TRANSFER), msgLen, ret),
55 : ret);
56 :
57 0 : return HCCL_SUCCESS;
58 : }
59 :
60 0 : HcclResult TopoInfoExchangeBase::SendClusterInfo(std::shared_ptr<HcclSocket> socket, const RankTable_t& clusterInfo)
61 : {
62 0 : nlohmann::json basicJson;
63 0 : CHK_RET(Struct2Json(clusterInfo, basicJson));
64 0 : basicJson[PROP_STEP] = currentStep_; // add step to verify.
65 0 : std::string buffer = basicJson.dump();
66 0 : u32 msgLen = buffer.length();
67 0 : CHK_RET(SendClusterInfoMsg(socket, clusterInfo, buffer, msgLen));
68 0 : currentStep_++;
69 0 : return HCCL_SUCCESS;
70 0 : }
71 :
72 0 : void TopoInfoExchangeBase::PrintRecvFailReasons(std::shared_ptr<HcclSocket> socket, HcclResult ret)
73 : {
74 0 : HCCL_ERROR(
75 : "[%s][%s]receive msg length from fdhandle failed, ret[%d]", LOG_KEYWORDS_INIT_GROUP.c_str(),
76 : LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), ret);
77 0 : HCCL_ERROR(
78 : "Current rank get socket with server[%s] success, but wait for recv rankTable from server failed, maybe due to "
79 : "following reasons:",
80 : socket->GetRemoteIp().GetReadableIP());
81 0 : HCCL_ERROR(
82 : "1. client wait for recv timeout, please check [ERROR] info in server[%s], whether all ranks were executed to "
83 : "create the communication",
84 : socket->GetRemoteIp().GetReadableIP());
85 0 : HCCL_ERROR("2. in large-scale cluster scenarios, occasional connection failures may occur due to the maximum "
86 : "connection limit in the system configuration. ");
87 0 : HCCL_ERROR(" these issues can be resolved by modifying the system configuration in all node: `sysctl -w "
88 : "net.core.somaxconn=65535` and `sysctl -w net.ipv4.tcp_max_syn_backlog=65535`");
89 0 : }
90 :
91 0 : HcclResult TopoInfoExchangeBase::RecvClusterInfoMsg(std::shared_ptr<HcclSocket> socket, RankTable_t& clusterInfo)
92 : {
93 0 : const u32 recvBufferLimit = 100 * 1024 * 1024; // 100 * 1024 * 1024 = 100MB
94 0 : u32 msgLen = 0;
95 0 : std::string errormessage = "";
96 0 : HcclResult ret = socket->Recv(reinterpret_cast<char*>(&msgLen), sizeof(msgLen));
97 0 : if (ret == HCCL_E_TIMEOUT) {
98 : errormessage = "Receiving message from the root node timed out. Check whether node "
99 0 : + std::string(socket->GetRemoteIp().GetReadableIP()) + " reports an error";
100 0 : RPT_INPUT_ERR(
101 : true, "EI0015", std::vector<std::string>({"error_reason"}), std::vector<std::string>({errormessage}));
102 : }
103 0 : CHK_PRT_RET(ret != HCCL_SUCCESS, PrintRecvFailReasons(socket, ret), HCCL_E_INTERNAL);
104 0 : CHK_PRT_RET(
105 : ((msgLen == 0) || (msgLen > recvBufferLimit)),
106 : HCCL_ERROR(
107 : "[%s][%s]receive msg "
108 : "length[%u] from fdhandle failed, msg length is beyond [1 ~ %u].",
109 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), msgLen, recvBufferLimit),
110 : HCCL_E_INTERNAL);
111 :
112 0 : u32 recvBufferLen = msgLen + 1;
113 0 : HostMem recvMsg = HostMem::alloc(recvBufferLen);
114 0 : CHK_PTR_NULL(recvMsg.ptr());
115 0 : char* recvMsgBuf = static_cast<char*>(recvMsg.ptr());
116 :
117 0 : s32 sRet = memset_s(recvMsgBuf, recvBufferLen, 0, recvBufferLen);
118 0 : CHK_PRT_RET(
119 : sRet != EOK,
120 : HCCL_ERROR(
121 : "[%s][%s]sockBuff memset failed", LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RANKTABLE_DETECT.c_str()),
122 : HCCL_E_MEMORY);
123 0 : ret = socket->Recv(recvMsgBuf, msgLen);
124 0 : if (ret == HCCL_E_TIMEOUT) {
125 0 : RPT_INPUT_ERR(
126 : true, "EI0015", std::vector<std::string>({"error_reason"}), std::vector<std::string>({errormessage}));
127 : }
128 0 : CHK_PRT_RET(
129 : ret != HCCL_SUCCESS,
130 : HCCL_ERROR(
131 : "[%s][%s]receive from fdhandle failed ,ret[%d]", LOG_KEYWORDS_INIT_GROUP.c_str(),
132 : LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), ret),
133 : HCCL_E_INTERNAL);
134 0 : nlohmann::json jClusterJson;
135 0 : CHK_RET(parseJsonBuff(recvMsgBuf, recvBufferLen, jClusterJson));
136 :
137 : // Verify json basic info
138 : u32 step;
139 0 : CHK_RET(JsonUtils::GetJsonProperty(jClusterJson, PROP_STEP, step));
140 :
141 0 : CHK_PRT_RET(
142 : step != currentStep_,
143 : HCCL_ERROR(
144 : "[Recv][ClusterInfoMsg]RecvClusterInfo step failed "
145 : "step[%u] vs currentStep_[%u]",
146 : step, currentStep_),
147 : HCCL_E_INTERNAL);
148 :
149 0 : s32 logicDevId = 0;
150 0 : u32 devPhyId = 0;
151 0 : CHK_RET(hrtGetDevice(&logicDevId));
152 0 : CHK_RET(hrtGetDevicePhyIdByIndex(static_cast<u32>(logicDevId), devPhyId));
153 0 : HcclIpAddress localHostIp;
154 0 : CHK_RET(GetLocalHostIP(localHostIp, devPhyId));
155 :
156 : bool isRoot
157 0 : = (localHostIp == GetExternalInputMasterInfo().serverIp
158 0 : && logicDevId == static_cast<s32>(GetExternalInputMasterInfo().serverDeviceId));
159 : errormessage = "No rank in the communicator can connect to the root node within the timeout period. List of "
160 : "unconnected ranks: "
161 0 : + std::string(jClusterJson["fault_info"].dump().c_str());
162 0 : if (!isRoot && jClusterJson.find("fault_type") != jClusterJson.end()
163 0 : && jClusterJson.find("fault_info") != jClusterJson.end()) {
164 0 : RPT_INPUT_ERR(
165 : true, "EI0015", std::vector<std::string>({"error_reason"}), std::vector<std::string>({errormessage}));
166 0 : HCCL_ERROR(
167 : "[%s][%s] TopoDetect ERROR occur fault_type[%s], fault_info[%s]", LOG_KEYWORDS_INIT_GROUP.c_str(),
168 : LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), jClusterJson["fault_type"].dump().c_str(),
169 : jClusterJson["fault_info"].dump().c_str());
170 : }
171 :
172 0 : ret = Json2Struct(jClusterJson, clusterInfo);
173 0 : CHK_PRT_RET(
174 : ret != HCCL_SUCCESS, HCCL_ERROR("[Recv][ClusterInfoMsg]step[%u] json to struct failed!", currentStep_),
175 : HCCL_E_INTERNAL);
176 0 : return HCCL_SUCCESS;
177 0 : }
178 :
179 0 : HcclResult TopoInfoExchangeBase::RecvClusterInfo(std::shared_ptr<HcclSocket> socket, RankTable_t& clusterInfo)
180 : {
181 0 : CHK_RET(RecvClusterInfoMsg(socket, clusterInfo));
182 0 : std::string errormessage = "";
183 0 : if (isByMasterInfo_) {
184 0 : u32 identify = 0;
185 0 : auto ret = socket->Recv(reinterpret_cast<char*>(&identify), sizeof(identify));
186 0 : if (ret == HCCL_E_TIMEOUT) {
187 : errormessage = "Receiving message from the root node timed out. Check whether node "
188 0 : + std::string(socket->GetRemoteIp().GetReadableIP()) + " reports an error";
189 0 : RPT_INPUT_ERR(
190 : true, "EI0015", std::vector<std::string>({"error_reason"}), std::vector<std::string>({errormessage}));
191 : }
192 0 : CHK_PRT_RET(
193 : ret != HCCL_SUCCESS,
194 : HCCL_ERROR(
195 : "[%s][%s] receive identify from fdhandle failed", LOG_KEYWORDS_INIT_GROUP.c_str(),
196 : LOG_KEYWORDS_RANKTABLE_DETECT.c_str()),
197 : HCCL_E_INTERNAL);
198 0 : identifierNum_ = identify;
199 : }
200 0 : currentStep_++;
201 0 : return HCCL_SUCCESS;
202 0 : }
203 :
204 0 : HcclResult TopoInfoExchangeBase::RecvClusterJson(std::shared_ptr<HcclSocket> socket, nlohmann::json& jClusterJson)
205 : {
206 0 : const u32 recvBufferLimit = 10 * 1024 * 1024; // 10 * 1024 * 1024 = 10MB
207 0 : u32 msgLen = 0;
208 0 : std::string errormessage = "";
209 0 : HcclResult ret = socket->Recv(reinterpret_cast<char*>(&msgLen), sizeof(msgLen));
210 0 : if (ret == HCCL_E_TIMEOUT) {
211 : errormessage = "Receiving message from the root node timed out. Check whether node "
212 0 : + std::string(socket->GetRemoteIp().GetReadableIP()) + " reports an error";
213 0 : RPT_INPUT_ERR(
214 : true, "EI0015", std::vector<std::string>({"error_reason"}), std::vector<std::string>({errormessage}));
215 : }
216 0 : CHK_PRT_RET(
217 : ret != HCCL_SUCCESS,
218 : HCCL_ERROR(
219 : "[%s][%s] receive msg length from fdhandle failed, ret[%d]", LOG_KEYWORDS_INIT_GROUP.c_str(),
220 : LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), ret),
221 : HCCL_E_INTERNAL);
222 0 : CHK_PRT_RET(
223 : ((msgLen == 0) || (msgLen > recvBufferLimit)),
224 : HCCL_ERROR(
225 : "[%s][%s]receive msg length "
226 : "from fdhandle failed, msg length is beyond [1 ~ %u].",
227 : LOG_KEYWORDS_INIT_GROUP.c_str(), LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), recvBufferLimit),
228 : HCCL_E_INTERNAL);
229 :
230 0 : u32 recvBufferLen = msgLen + 1;
231 0 : HostMem recvMsg = HostMem::alloc(recvBufferLen);
232 0 : CHK_PTR_NULL(recvMsg.ptr());
233 0 : char* recvMsgBuf = static_cast<char*>(recvMsg.ptr());
234 :
235 0 : s32 sRet = memset_s(recvMsgBuf, recvBufferLen, 0, recvBufferLen);
236 0 : CHK_PRT_RET(sRet != EOK, HCCL_ERROR("[Recv][ClusterInfoMsg]sockBuff memset failed"), HCCL_E_MEMORY);
237 0 : ret = socket->Recv(recvMsgBuf, msgLen);
238 0 : if (ret == HCCL_E_TIMEOUT) {
239 : errormessage = "Receiving message from the root node timed out. Check whether node "
240 0 : + std::string(socket->GetRemoteIp().GetReadableIP()) + " reports an error";
241 0 : RPT_INPUT_ERR(
242 : true, "EI0015", std::vector<std::string>({"error_reason"}), std::vector<std::string>({errormessage}));
243 : }
244 0 : CHK_PRT_RET(
245 : ret != HCCL_SUCCESS,
246 : HCCL_ERROR(
247 : "[%s][%s] receive from fdhandle failed ,ret[%d]", LOG_KEYWORDS_INIT_GROUP.c_str(),
248 : LOG_KEYWORDS_RANKTABLE_DETECT.c_str(), ret),
249 : HCCL_E_INTERNAL);
250 0 : CHK_RET(parseJsonBuff(recvMsgBuf, recvBufferLen, jClusterJson));
251 :
252 0 : return HCCL_SUCCESS;
253 0 : }
254 :
255 0 : HcclResult TopoInfoExchangeBase::RecvGrpLeaderInfoMsg(std::shared_ptr<HcclSocket> socket, GroupLeader_t& LeaderInfo)
256 : {
257 0 : nlohmann::json jClusterJson;
258 0 : CHK_RET(RecvClusterJson(socket, jClusterJson));
259 :
260 : // Verify json basic info
261 : u32 step;
262 0 : CHK_RET(JsonUtils::GetJsonProperty(jClusterJson, PROP_STEP, step));
263 :
264 0 : CHK_PRT_RET(
265 : step != currentStep_,
266 : HCCL_ERROR(
267 : "[Recv][ClusterInfoMsg]RecvClusterInfo step failed "
268 : "step[%u] vs currentStep_[%u]",
269 : step, currentStep_),
270 : HCCL_E_INTERNAL);
271 :
272 0 : HcclResult ret = Json2GrpLeader(jClusterJson, LeaderInfo);
273 0 : CHK_PRT_RET(
274 : ret != HCCL_SUCCESS, HCCL_ERROR("[Recv][ClusterInfoMsg]step[%u] json to struct failed!", currentStep_),
275 : HCCL_E_INTERNAL);
276 :
277 0 : return HCCL_SUCCESS;
278 0 : }
279 :
280 0 : HcclResult TopoInfoExchangeBase::BlockReceive(std::shared_ptr<HcclSocket> socket, char* buff, u32 size) const
281 : {
282 0 : CHK_PTR_NULL(buff);
283 0 : CHK_RET(socket->Recv(buff, size));
284 0 : return HCCL_SUCCESS;
285 : }
286 :
287 0 : HcclResult TopoInfoExchangeBase::parseJsonBuff(const char buff[], u32 buffLen, nlohmann::json& buffJson) const
288 : {
289 0 : u32 len = strnlen(buff, buffLen);
290 0 : CHK_PRT_RET(
291 : (len > buffLen || len == 0),
292 : HCCL_ERROR("[Parse][JsonBuff]buff len invalid, buff len[%u], msgLen[%u]", len, buffLen), HCCL_E_INTERNAL);
293 :
294 0 : CHK_RET(JsonUtils::ParseInformation(buffJson, buff));
295 : u32 step;
296 0 : CHK_RET(JsonUtils::GetJsonProperty(buffJson, PROP_STEP, step));
297 0 : if (step != currentStep_) {
298 0 : HCCL_ERROR(
299 : "[Parse][JsonBuff]errNo[0x%016llx] received step[%u] is invalid , expect step is %u",
300 : HCCL_ERROR_CODE(HCCL_E_INTERNAL), step, currentStep_);
301 0 : return HCCL_E_INTERNAL;
302 : }
303 0 : return HCCL_SUCCESS;
304 : }
305 :
306 0 : HcclResult TopoInfoExchangeBase::Json2GrpLeader(const nlohmann::json& jClusterJson, GroupLeader_t& GrpLeaderInfo) const
307 : {
308 0 : GrpLeaderInfo.grpLeaderNum = jClusterJson[PROP_RANK_NUM];
309 0 : for (auto& leaderInfoJson : jClusterJson[PROP_GROUP_LEADER_LIST]) {
310 : HcclRankHandle rankHandle;
311 0 : std::string strTmp = leaderInfoJson[PROP_NETWORK_IPADDR];
312 0 : s32 sRet = memcpy_s(rankHandle.ip, IP_ADDRESS_BUFFER_LEN, strTmp.c_str(), strTmp.size());
313 0 : CHK_PRT_RET(sRet != EOK, HCCL_ERROR("[Json2GrpLeader]memcpy_s failed, errorno[%d]", sRet), HCCL_E_MEMORY);
314 0 : rankHandle.ip[strTmp.size()] = '\0';
315 0 : rankHandle.port = leaderInfoJson[PROP_NETWORK_NETWORKPORT];
316 0 : strTmp = leaderInfoJson[PROP_NETWORK_IDENTIFIER];
317 0 : sRet = memcpy_s(rankHandle.identifier, ROOTINFO_INDENTIFIER_MAX_LENGTH, strTmp.c_str(), strTmp.size());
318 0 : CHK_PRT_RET(sRet != EOK, HCCL_ERROR("[Json2GrpLeader]memcpy_s failed, errorno[%d]", sRet), HCCL_E_MEMORY);
319 0 : rankHandle.identifier[strTmp.size()] = '\0';
320 0 : rankHandle.nicDeploy = leaderInfoJson[PROP_DEPLOY_MODE];
321 0 : rankHandle.rankId = leaderInfoJson[PROP_RANK_ID];
322 0 : GrpLeaderInfo.GroupLeaderList.emplace_back(rankHandle);
323 0 : }
324 :
325 0 : return HCCL_SUCCESS;
326 : }
327 :
328 0 : HcclResult TopoInfoExchangeBase::Json2Struct(const nlohmann::json& jClusterJson, RankTable_t& clusterInfo) const
329 : {
330 0 : CHK_RET(SetClusterDeploy(jClusterJson, clusterInfo)); // deploymode为枚举类变量 需单独处理判断
331 0 : CHK_RET(JsonUtils::GetJsonProperty(jClusterJson, PROP_DEV_NUM, clusterInfo.deviceNum));
332 0 : CHK_RET(JsonUtils::GetJsonProperty(jClusterJson, PROP_SRV_NUM, clusterInfo.serverNum));
333 0 : CHK_RET(JsonUtils::GetJsonProperty(jClusterJson, PROP_SUPER_POD_NUM, clusterInfo.superPodNum));
334 0 : CHK_RET(JsonUtils::GetJsonProperty(jClusterJson, PROP_RANK_NUM, clusterInfo.rankNum));
335 0 : for (auto& rankInfoJson : jClusterJson[PROP_RANK_LIST]) {
336 0 : RankInfo_t rankInfo;
337 0 : rankInfo.rankId = rankInfoJson[PROP_RANK_ID];
338 0 : rankInfo.serverId = rankInfoJson[PROP_SERVER_ID];
339 0 : rankInfo.tlsStatus = rankInfoJson[PROP_TLS_STATUS];
340 0 : CHK_RET(rankInfo.hostIp.SetReadableAddress(rankInfoJson[PROP_HOST_IP]));
341 0 : rankInfo.deviceInfo.devicePhyId = rankInfoJson[PROP_DEV_INFO][PROP_DEV_ID];
342 0 : rankInfo.deviceInfo.deviceType = rankInfoJson[PROP_DEV_INFO][PROP_DEV_TYPE];
343 0 : CHK_PRT_RET(
344 : rankInfoJson[PROP_DEV_INFO].find(PROP_DEV_NIC_PORT) == rankInfoJson[PROP_DEV_INFO].end()
345 : || rankInfoJson[PROP_DEV_INFO].find(PROP_DEV_VNIC_PORT) == rankInfoJson[PROP_DEV_INFO].end()
346 : || rankInfoJson[PROP_DEV_INFO].find(PROP_BACKUP_DEV_PORT) == rankInfoJson[PROP_DEV_INFO].end(),
347 : HCCL_ERROR("[Json2Struct] Fail to find port infos in rank info json. "
348 : "Please make sure the CANN version is consistent within the communication."),
349 : HCCL_E_NOT_SUPPORT);
350 0 : rankInfo.deviceInfo.port = rankInfoJson[PROP_DEV_INFO][PROP_DEV_NIC_PORT];
351 0 : rankInfo.deviceInfo.vnicPort = rankInfoJson[PROP_DEV_INFO][PROP_DEV_VNIC_PORT];
352 0 : rankInfo.deviceInfo.backupPort = rankInfoJson[PROP_DEV_INFO][PROP_BACKUP_DEV_PORT];
353 0 : for (auto& devIp : rankInfoJson[PROP_DEV_INFO][PROP_DEV_IP]) {
354 0 : std::string ipStr = devIp;
355 0 : rankInfo.deviceInfo.deviceIp.emplace_back(ipStr);
356 0 : }
357 0 : if (rankInfoJson[PROP_DEV_INFO].find(PROP_BACKUP_DEV_IP) == rankInfoJson[PROP_DEV_INFO].end()) {
358 0 : HCCL_RUN_WARNING("[Json2Struct] Fail to find backup device ip in rank info json. "
359 : "Backup device ip will not be parsed. If you want to use backup device ip, "
360 : "Please make sure the CANN version is consistent within the communication.");
361 : } else {
362 0 : for (auto& backupDevIp : rankInfoJson[PROP_DEV_INFO][PROP_BACKUP_DEV_IP]) {
363 0 : std::string backupIpStr = backupDevIp;
364 0 : rankInfo.deviceInfo.backupDeviceIp.emplace_back(backupIpStr);
365 0 : }
366 : }
367 :
368 0 : rankInfo.superPodId = rankInfoJson[PROP_SUPER_POD_ID];
369 0 : rankInfo.superDeviceId = rankInfoJson[PROP_SUPER_DEVICE_ID];
370 :
371 : /* Optional: for second communication stage */
372 0 : if (rankInfoJson.find(PROP_TRANS_INFO) != rankInfoJson.end()) {
373 0 : for (auto& transInfoJson : rankInfoJson[PROP_TRANS_INFO]) {
374 : TransportInfo_t transportInfo;
375 0 : transportInfo.dstRankId = transInfoJson[PROP_DEST_RANK];
376 0 : transportInfo.transportType = transInfoJson[PROP_TRANS_TYPE];
377 0 : rankInfo.transportInfo.push_back(transportInfo);
378 : }
379 : }
380 0 : clusterInfo.rankList.push_back(rankInfo);
381 0 : }
382 0 : for (auto& serverInfoJson : jClusterJson[PROP_SERVER_LIST]) {
383 0 : ServerInfo_t serverInfo;
384 0 : serverInfo.serverId = serverInfoJson[PROP_SERVER_ID];
385 0 : for (auto& networkInfoJson : serverInfoJson[PROP_NETWORK_INFO_LIST]) {
386 0 : NetworkInfo_t networkInfo;
387 0 : networkInfo.ethName = networkInfoJson[PROP_NETWORK_ETHNAME];
388 0 : CHK_RET(networkInfo.ipAddr.SetReadableAddress(networkInfoJson[PROP_NETWORK_IPADDR]));
389 0 : networkInfo.networkPort = networkInfoJson[PROP_NETWORK_NETWORKPORT];
390 0 : CHK_RET(networkInfo.refIp.SetReadableAddress(networkInfoJson[PROP_NETWORK_REFIP]));
391 0 : networkInfo.planeID = networkInfoJson[PROP_NETWORK_PLANEID];
392 0 : serverInfo.networkInfo.push_back(networkInfo);
393 0 : }
394 0 : clusterInfo.serverList.push_back(serverInfo);
395 0 : }
396 :
397 0 : return HCCL_SUCCESS;
398 : }
399 :
400 13 : HcclResult TopoInfoExchangeBase::Struct2Json(const RankTable_t& clusterInfo, nlohmann::json& ClusterJson)
401 : {
402 13 : nlohmann::json rankListJson;
403 13 : nlohmann::json serverListJson;
404 :
405 13 : TransformRankListToJson(clusterInfo, rankListJson);
406 13 : for (auto& serverInfo : clusterInfo.serverList) {
407 0 : nlohmann::json serverJson;
408 0 : serverJson[PROP_SERVER_ID] = serverInfo.serverId;
409 0 : nlohmann::json networkInfoListJson;
410 0 : for (auto& networkInfo : serverInfo.networkInfo) {
411 0 : nlohmann::json networkInfoJson;
412 0 : networkInfoJson[PROP_NETWORK_ETHNAME] = networkInfo.ethName;
413 0 : networkInfoJson[PROP_NETWORK_IPADDR] = std::string(networkInfo.ipAddr.GetReadableIP());
414 0 : networkInfoJson[PROP_NETWORK_NETWORKPORT] = networkInfo.networkPort;
415 0 : networkInfoJson[PROP_NETWORK_REFIP] = std::string(networkInfo.refIp.GetReadableIP());
416 0 : networkInfoJson[PROP_NETWORK_PLANEID] = networkInfo.planeID;
417 0 : networkInfoListJson.push_back(networkInfoJson);
418 0 : }
419 0 : serverJson[PROP_NETWORK_INFO_LIST] = networkInfoListJson;
420 0 : serverListJson.push_back(serverJson);
421 0 : }
422 :
423 13 : ClusterJson[PROP_RANK_NUM] = clusterInfo.rankNum;
424 13 : ClusterJson[PROP_DEV_NUM] = clusterInfo.deviceNum;
425 13 : ClusterJson[PROP_SRV_NUM] = clusterInfo.serverNum;
426 13 : ClusterJson[PROP_SUPER_POD_NUM] = clusterInfo.superPodNum;
427 13 : ClusterJson[PROP_DEPLOY_MODE] = clusterInfo.nicDeploy;
428 13 : ClusterJson[PROP_RANK_LIST] = rankListJson;
429 13 : ClusterJson[PROP_SERVER_LIST] = serverListJson;
430 13 : return HCCL_SUCCESS;
431 13 : }
432 :
433 0 : HcclResult TopoInfoExchangeBase::GrpLeader2Json(const GroupLeader_t& GrpLeaderInfo, nlohmann::json& GroupLeaderJson)
434 : {
435 0 : nlohmann::json leaderListJson;
436 :
437 0 : for (auto& leaderInfo : GrpLeaderInfo.GroupLeaderList) {
438 0 : nlohmann::json leaderJson;
439 :
440 0 : leaderJson[PROP_NETWORK_IPADDR] = std::string(leaderInfo.ip);
441 0 : leaderJson[PROP_NETWORK_NETWORKPORT] = leaderInfo.port;
442 0 : leaderJson[PROP_NETWORK_IDENTIFIER] = std::string(leaderInfo.identifier);
443 0 : leaderJson[PROP_DEPLOY_MODE] = leaderInfo.nicDeploy;
444 0 : leaderJson[PROP_RANK_ID] = leaderInfo.rankId;
445 :
446 0 : leaderListJson.push_back(leaderJson);
447 0 : }
448 :
449 0 : GroupLeaderJson[PROP_RANK_NUM] = GrpLeaderInfo.grpLeaderNum;
450 0 : GroupLeaderJson[PROP_GROUP_LEADER_LIST] = leaderListJson;
451 :
452 0 : return HCCL_SUCCESS;
453 0 : }
454 :
455 : HcclResult
456 13 : TopoInfoExchangeBase::TransformRankListToJson(const RankTable_t& clusterInfo, nlohmann::json& rankListJson) const
457 : {
458 13 : for (auto& rankInfo : clusterInfo.rankList) {
459 0 : nlohmann::json deviceIp;
460 0 : for (auto& devIp : rankInfo.deviceInfo.deviceIp) {
461 0 : deviceIp.push_back(std::string(devIp.GetReadableIP()));
462 : }
463 0 : nlohmann::json backupDeviceIp;
464 0 : for (auto& backupDevIp : rankInfo.deviceInfo.backupDeviceIp) {
465 0 : backupDeviceIp.push_back(std::string(backupDevIp.GetReadableIP()));
466 : }
467 0 : nlohmann::json devInfoJson;
468 0 : devInfoJson[PROP_DEV_ID] = rankInfo.deviceInfo.devicePhyId;
469 0 : devInfoJson[PROP_DEV_TYPE] = rankInfo.deviceInfo.deviceType;
470 0 : devInfoJson[PROP_DEV_NIC_PORT] = rankInfo.deviceInfo.port;
471 0 : devInfoJson[PROP_DEV_VNIC_PORT] = rankInfo.deviceInfo.vnicPort;
472 0 : devInfoJson[PROP_BACKUP_DEV_PORT] = rankInfo.deviceInfo.backupPort;
473 0 : devInfoJson[PROP_DEV_IP] = deviceIp;
474 0 : devInfoJson[PROP_BACKUP_DEV_IP] = backupDeviceIp;
475 0 : nlohmann::json rankJson;
476 0 : rankJson[PROP_RANK_ID] = rankInfo.rankId;
477 0 : rankJson[PROP_SERVER_ID] = rankInfo.serverId;
478 0 : rankJson[PROP_HOST_IP] = std::string(rankInfo.hostIp.GetReadableIP());
479 0 : rankJson[PROP_DEV_INFO] = devInfoJson;
480 :
481 0 : rankJson[PROP_SUPER_POD_ID] = rankInfo.superPodId;
482 0 : rankJson[PROP_SUPER_DEVICE_ID] = rankInfo.superDeviceId;
483 0 : rankJson[PROP_TLS_STATUS] = rankInfo.tlsStatus;
484 :
485 : /* Optional: for second communication stage */
486 0 : nlohmann::json transInfosJson;
487 0 : for (auto& transInfo : rankInfo.transportInfo) {
488 0 : nlohmann::json transInfoJson;
489 0 : transInfoJson[PROP_TRANS_TYPE] = transInfo.transportType;
490 0 : transInfoJson[PROP_DEST_RANK] = transInfo.dstRankId;
491 0 : transInfosJson.push_back(transInfoJson);
492 0 : }
493 0 : if (!transInfosJson.empty()) {
494 0 : rankJson[PROP_TRANS_INFO] = transInfosJson;
495 : }
496 0 : rankListJson.push_back(rankJson);
497 0 : }
498 13 : return HCCL_SUCCESS;
499 : }
500 :
501 0 : HcclResult TopoInfoExchangeBase::SetClusterDeploy(const nlohmann::json& jClusterJson, RankTable_t& clusterInfo) const
502 : {
503 0 : if (!jClusterJson.contains(PROP_DEPLOY_MODE)) {
504 0 : HCCL_ERROR("[SetClusterDeploy]PROP_DEPLOY_MODE is invalid");
505 0 : return HCCL_E_INTERNAL;
506 : }
507 0 : clusterInfo.nicDeploy = jClusterJson[PROP_DEPLOY_MODE];
508 0 : return HCCL_SUCCESS;
509 : }
510 :
511 : } // namespace hccl
|