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