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