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 "opretry_server.h"
12 : #include "externalinput_pub.h"
13 : #include "heartbeat.h"
14 : #include "comm_configer.h"
15 :
16 : namespace hccl {
17 :
18 12 : HcclResult CreateOpRetryServerByState(RetryState state, RetryContext* retryCtx)
19 : {
20 12 : HCCL_INFO("[OpRetry][Server]CreateOpRetryServerByState state[%s]", GetReadableState(state));
21 12 : std::shared_ptr<OpRetryBase> retryPtr = nullptr;
22 12 : switch (state) {
23 9 : case RETRY_STATE_SERVER_RUNNING: {
24 9 : EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerRunning>()), return HCCL_E_PTR);
25 9 : break;
26 : }
27 0 : case RETRY_STETA_HANDLE_ALL_ERR: {
28 0 : EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerHandleError>()), return HCCL_E_PTR);
29 0 : break;
30 : }
31 1 : case RETRY_STATE_SERVER_RETRY_FAIL: {
32 1 : EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerRetryFail>()), return HCCL_E_PTR);
33 1 : break;
34 : }
35 0 : case RETRY_STATE_WAIT_LINK_CHECKED:
36 0 : EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerWaitLinkInfo>()), return HCCL_E_PTR);
37 0 : break;
38 0 : case RETRY_STATE_WAIT_AICPU_STOPED:
39 : case RETRY_STATE_WAIT_STREAM_STOPED:
40 : case RETRY_STATE_WAIT_STREAM_CLEARED:
41 : case RETRY_STATE_WAIT_STOP_TRANSPORT:
42 : case RETRY_STATE_WAIT_NOTIFY_RESETED:
43 : case RETRY_STATE_WAIT_RESUME_TRANSPORT:
44 : case RETRY_STATE_WAIT_CHECK_INFO:
45 : case RETRY_STATE_WAIT_CAN_RETRY: {
46 0 : EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerWaitResp>()), return HCCL_E_PTR);
47 0 : break;
48 : }
49 0 : case RETRY_STATE_CHECK_ALL_LINK:
50 0 : EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerCheckAllLink>()), return HCCL_E_PTR);
51 0 : break;
52 0 : case RETRY_STATE_CMD_RESUME_TRANSPORT: {
53 0 : EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerIssueChangeLinkAndResume>()), return HCCL_E_PTR);
54 0 : break;
55 : }
56 1 : case RETRY_STATE_CMD_STOP_AICPU:
57 : case RETRY_STATE_CMD_STOP_STREAM:
58 : case RETRY_STATE_CMD_CLEAR_STREAM:
59 : case RETRY_STATE_CMD_STOP_TRANSPORT:
60 : case RETRY_STATE_CMD_CHECK_LINK:
61 : case RETRY_STATE_CMD_RESET_NOTIFY:
62 : case RETRY_STATE_CMD_CHECK:
63 : case RETRY_STATE_CMD_CAN_RETRY: {
64 1 : EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerIssueCmd>()), return HCCL_E_PTR);
65 1 : break;
66 : }
67 0 : case RETRY_STATE_CHECK_OP: {
68 0 : EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerCheckOp>()), return HCCL_E_PTR);
69 0 : break;
70 : }
71 : // 检查各agennt主动接轨信息,并发送cmd命令
72 1 : case RETRY_STATE_CMD_PLAN_SWITCH_NIC: {
73 1 : EXCEPTION_CATCH((retryPtr = std::make_shared<SwitchNicServerCheckAllSwitchRanks>()), return HCCL_E_PTR);
74 1 : break;
75 : }
76 : // Server下发命令给agent进行链路检查,接收到agent回复消息后,检查所有链路情况
77 0 : case RETRY_RESUME_STATE_SERVER_CHECK_LINK: {
78 0 : EXCEPTION_CATCH((retryPtr = std::make_shared<ResumeServerCheckAllLink>()), return HCCL_E_PTR);
79 0 : break;
80 : }
81 : // Server下发命令给agent进行借轨操作,接收到借轨成功消息后,切换下一状态
82 0 : case RETRY_RESUME_STATE_SERVER_CHANGE_LINK: {
83 0 : EXCEPTION_CATCH((retryPtr = std::make_shared<ResumeServerChangeLink>()), return HCCL_E_PTR);
84 0 : break;
85 : }
86 0 : default: {
87 0 : HCCL_ERROR(
88 : "[OpRetry][Server]CreateOpRetryServerByState failed, state[%s] is invalid", GetReadableState(state));
89 0 : return HCCL_E_NOT_SUPPORT;
90 : }
91 : }
92 12 : retryCtx->SetRetryState(state, retryPtr);
93 12 : return HCCL_SUCCESS;
94 12 : }
95 :
96 0 : HcclResult OpRetryServerBase::ProcessError(RetryContext* retryCtx)
97 : {
98 0 : HCCL_ERROR(
99 : "[%s]OpRetryServer run fail, rankId[%u], state[%s]", __func__, retryCtx->rankId_,
100 : retryCtx->GetReadableCtxState());
101 : // 状态切换至RETRY_STATE_SERVER_RETRY_FAIL
102 0 : CHK_RET(CreateOpRetryServerByState(RETRY_STATE_SERVER_RETRY_FAIL, retryCtx));
103 0 : return HCCL_SUCCESS;
104 : }
105 :
106 1 : HcclResult OpRetryServerRunning::ProcessEvent(RetryContext* retryCtx)
107 : {
108 1 : if (retryCtx->errorRankList_.size() > 0) {
109 : // 若当前errorRankList_中有未处理的errorRank,则先进行处理
110 0 : HCCL_RUN_INFO("[OpRetry][Server]deal rank from errorRankList_, size[%d]", retryCtx->errorRankList_.size());
111 0 : CHK_RET(CreateOpRetryServerByState(RETRY_STETA_HANDLE_ALL_ERR, retryCtx));
112 0 : return HCCL_SUCCESS;
113 : }
114 :
115 1 : const std::chrono::seconds timeout = std::chrono::seconds(OP_RETRY_KEEP_INTERVAL);
116 : // 轮询接收agent信息
117 1 : for (auto& it : retryCtx->serverSockets_) {
118 1 : const u32& agentId = it.first;
119 : // 若对端已经关闭, 则不再轮询
120 1 : if (disableAgent_.find(agentId) != disableAgent_.end()) {
121 0 : continue;
122 : }
123 :
124 : // 记录时间, 检测和对端上一次通信时间是否超过保活时间
125 1 : std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
126 1 : if (lastRecvTimes_.find(agentId) == lastRecvTimes_.end()) {
127 1 : lastRecvTimes_.insert(std::make_pair(agentId, curTime));
128 : }
129 :
130 : // 轮询接收agent状态机信息
131 1 : HcclResult ret = WaitResponse(it.second.socket, it.second.retryInfo);
132 1 : if (ret == HCCL_SUCCESS) { // 成功接收到数据
133 : // 新增逻辑:检查是否需要接收 ActiveSwitchInfo,进入主动接轨校验阶段
134 1 : if (it.second.retryInfo.retryState == RETRY_STATE_SEND_SWITCH_INFO) {
135 : // 1. 接收剩余字段
136 1 : ret = RecvActiveSwitchInfo(it.second.socket, agentId, it.second.switchInfo);
137 1 : if (ret != HCCL_SUCCESS) {
138 0 : disableAgent_.insert(agentId);
139 : } else {
140 1 : retryCtx->switchInfoMap_[agentId] = it.second.switchInfo;
141 1 : HCCL_INFO("[SwitchNic][Server] recv first ActiveSwitchInfo from rank[%u] while running", agentId);
142 : // 2. 此时 activeInfo 包含完整数据,可用于后续逻辑
143 2 : CHK_RET(CreateOpRetryServerByState(RETRY_STATE_CMD_PLAN_SWITCH_NIC, retryCtx));
144 1 : return HCCL_SUCCESS;
145 : }
146 : }
147 0 : RetryState nextState = RETRY_STATE_SERVER_RUNNING;
148 0 : CHK_RET(ParaseErrorCode(retryCtx, it.second, nextState));
149 0 : if (nextState != RETRY_STATE_SERVER_RUNNING) {
150 : // 收到第一个报错后加入errorRankList_中,并切换到RETRY_STETA_HANDLE_ALL_ERR状态
151 0 : HCCL_RUN_INFO(
152 : "[OpRetry][Server]agent[%u] tag[%s] index[%u] find error, insert to errorRankList_", agentId,
153 : it.second.retryInfo.opInfo.opId.tag, it.second.retryInfo.opInfo.opId.index);
154 0 : retryCtx->errorRankList_.insert(std::make_pair(agentId, it.second.retryInfo.opInfo.opId));
155 0 : CHK_RET(CreateOpRetryServerByState(RETRY_STETA_HANDLE_ALL_ERR, retryCtx));
156 0 : return HCCL_SUCCESS;
157 : }
158 0 : lastRecvTimes_[agentId] = curTime;
159 0 : } else if (ret == HCCL_E_AGAIN) { // 未接收到数据
160 : // 校验是否超时
161 0 : const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - lastRecvTimes_[agentId]);
162 0 : if (elapsed > timeout) {
163 0 : HCCL_WARNING(
164 : "[OpRetry][Server]OpRetryServerRunning recv Retry Frame from agentId[%u] timeout", agentId);
165 0 : lastRecvTimes_[agentId] = curTime;
166 : }
167 : } else { // 接收数据失败
168 0 : disableAgent_.insert(agentId);
169 0 : HCCL_RUN_INFO("[OpRetry][Server]WaitResponse from agentId[%u] fail, ret[%u]", agentId, ret);
170 : }
171 : }
172 :
173 : // 轮询间隔
174 0 : SaluSleep(OP_RETRY_RUNNING_POLL_INTERVAL);
175 0 : return HCCL_SUCCESS;
176 : }
177 :
178 1 : HcclResult OpRetryServerHandleError::ProcessEvent(RetryContext* retryCtx)
179 : {
180 1 : const u32 timeoutValue = std::max(static_cast<u32>(GetExternalInputHcclLinkTimeOut()), OP_RETRY_SEND_RECV_TIMEOUT)
181 1 : + OP_RETRY_WAIT_AGENT_AICPU_TIMEOUT;
182 1 : const std::chrono::seconds timeout = std::chrono::seconds(timeoutValue);
183 1 : std::chrono::steady_clock::time_point startTime = std::chrono::steady_clock::now();
184 1 : u32 waitTime = CommConfiger::GetInstance().GetCommConfigRetryHoldTime(retryCtx->group_);
185 : while (true) {
186 2 : CHK_PRT_RET(
187 : retryCtx->isServerStateWaitResume_,
188 : HCCL_RUN_INFO("[OpRetry][Server]switched state form wait handle error to wait resume"), HCCL_SUCCESS);
189 : // 判断是否超时
190 1 : std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
191 1 : const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - startTime);
192 1 : if (elapsed > timeout) {
193 0 : HCCL_ERROR("[OpRetry][Server] OpRetryServerHandleError timeout");
194 0 : for (auto& it : retryCtx->serverSockets_) {
195 0 : auto tag = std::string(reinterpret_cast<const char*>(it.second.retryInfo.opInfo.opId.tag));
196 0 : HCCL_ERROR(
197 : "[OpRetry][Server]OpRetryHandle retryinfo rank[%u] tag[%s] index[%u] IpInfo[%s]", it.first,
198 : tag.c_str(), it.second.retryInfo.opInfo.opId.index, it.second.retryInfo.dfxIpInfo);
199 0 : }
200 0 : return HCCL_E_TIMEOUT;
201 : }
202 :
203 : // 轮询接收agent信息,只期望收上开故障信息
204 2 : for (auto& it : retryCtx->serverSockets_) {
205 1 : const u32& agentId = it.first;
206 1 : if (retryCtx->errorRankList_.find(agentId) != retryCtx->errorRankList_.end()) {
207 : // 若当前rank已在errorRankList_,则不进行轮训
208 1 : continue;
209 : }
210 : // 轮询接收agent状态机信息
211 0 : HcclResult ret = WaitResponse(it.second.socket, it.second.retryInfo);
212 0 : if (ret == HCCL_SUCCESS) { // 成功接收到数据
213 0 : RetryState nextState = RETRY_STATE_SERVER_RUNNING;
214 0 : CHK_RET(ParaseErrorCode(retryCtx, it.second, nextState));
215 0 : if (nextState != RETRY_STATE_SERVER_RUNNING) {
216 : // 当前rank报错,收集到errorRankList_后统一处理
217 0 : HCCL_RUN_INFO(
218 : "[OpRetry][Server]agent[%u] tag[%s] index[%u] find error, insert to errorRankList_", agentId,
219 : it.second.retryInfo.opInfo.opId.tag, it.second.retryInfo.opInfo.opId.index);
220 0 : retryCtx->errorRankList_.insert(std::make_pair(agentId, it.second.retryInfo.opInfo.opId));
221 0 : continue;
222 : }
223 0 : } else if (ret == HCCL_E_AGAIN) {
224 : // 未收到数据,则发送一个保活数据给agent
225 0 : RetryCommandInfo commandInfo;
226 0 : commandInfo.command = RETRY_CMD_RUNNING;
227 0 : CHK_RET(IssueCommandWithOpId(it.second.socket, commandInfo));
228 : }
229 : }
230 :
231 1 : bool isFoundSendRecv = false;
232 1 : std::set<u32> errorRank;
233 2 : for (auto iter = retryCtx->errorRankList_.begin(); iter != retryCtx->errorRankList_.end(); ++iter) {
234 1 : errorRank.insert(iter->first);
235 : }
236 : // 对errorRankList_中rank进行遍历
237 2 : for (auto rank : errorRank) {
238 1 : if (retryCtx->errorRankList_[rank].isSendRecv) {
239 : // 当前报错rank中存在send/recv算子,优先处理send/recv算子
240 0 : isFoundSendRecv = true;
241 0 : auto curOpId = retryCtx->errorRankList_[rank];
242 0 : uint32_t remoteRank = (rank == curOpId.detRank) ? curOpId.srcRank : curOpId.detRank;
243 0 : auto& remoteOpId = retryCtx->serverSockets_[remoteRank].retryInfo.opInfo.opId;
244 0 : std::string curTag = std::string(reinterpret_cast<const char*>(curOpId.tag));
245 0 : std::string remoteTag = std::string(reinterpret_cast<const char*>(remoteOpId.tag));
246 : // sendrecv没有下边那两字段
247 0 : auto remoteSendTag = std::string(reinterpret_cast<const char*>(remoteOpId.bsrInfo[HCCL_SEND].bsrTag));
248 0 : auto remoteRecvTag = std::string(reinterpret_cast<const char*>(remoteOpId.bsrInfo[HCCL_RECV].bsrTag));
249 0 : HCCL_RUN_INFO(
250 : "[OpRetry][Server]curRank[%u], tag[%s], index[%u], startTaskComplete[%d]"
251 : "remoteRank[%u], remotetag[%s], remoteindex[%u], remoteStartTaskComplete[%d]"
252 : "Sendtag[%s], sendindex[%u]"
253 : "Recvtag[%s], recvindex[%u]",
254 : rank, curTag.c_str(), curOpId.index, curOpId.isBsrTaskStart, remoteRank, remoteTag.c_str(),
255 : remoteOpId.index, remoteOpId.isBsrTaskStart, remoteSendTag.c_str(),
256 : remoteOpId.bsrInfo[HCCL_SEND].index, remoteRecvTag.c_str(), remoteOpId.bsrInfo[HCCL_RECV].index);
257 0 : if (curOpId.opType == HcclCMDType::HCCL_CMD_BATCH_SEND_RECV && !remoteOpId.isBsrTaskStart) {
258 0 : continue;
259 : }
260 : // 如果对端也停在同一个send/recv算子,则触发该算子的重执行
261 0 : if ((curTag == remoteSendTag && curOpId.index == remoteOpId.bsrInfo[HCCL_SEND].index)
262 0 : || (curTag == remoteRecvTag && curOpId.index == remoteOpId.bsrInfo[HCCL_RECV].index)
263 0 : || (curTag == remoteTag && curOpId.index == remoteOpId.index)) {
264 : // 从errorRankList_中清除本端和对端rank
265 0 : retryCtx->errorRankList_.erase(rank);
266 0 : if (retryCtx->errorRankList_.find(remoteRank) != retryCtx->errorRankList_.end()) {
267 0 : if (curTag == remoteTag && curOpId.index == remoteOpId.index) {
268 0 : HCCL_RUN_INFO("[OpRetry][Server]delete remoteRank[%u] from errorRankList_", remoteRank);
269 0 : retryCtx->errorRankList_.erase(remoteRank);
270 : }
271 : }
272 : // 触发重执行
273 0 : HCCL_RUN_INFO(
274 : "[OpRetry][Server]begin to exec retry of tag[%s] from rank[%u] and rank[%u]", curTag.c_str(),
275 : rank, remoteRank);
276 0 : retryCtx->needRetryServerRanks_.clear();
277 0 : CHK_PRT(SetNeedRetryServerRank(retryCtx, curOpId));
278 0 : CHK_RET(CreateOpRetryServerByState(RETRY_STATE_CMD_STOP_AICPU, retryCtx));
279 0 : return HCCL_SUCCESS;
280 : }
281 0 : }
282 : }
283 :
284 1 : if (!isFoundSendRecv) {
285 1 : u32 firstErrorRank = *(errorRank.begin());
286 1 : auto curOpId = retryCtx->errorRankList_[firstErrorRank];
287 1 : auto curTag = std::string(reinterpret_cast<const char*>(curOpId.tag));
288 :
289 1 : retryCtx->errorRankList_.clear();
290 : // 开始重执行
291 1 : HCCL_RUN_INFO(
292 : "[OpRetry][Server]begin to exec retry of tag[%s] from rank[%u]", curTag.c_str(), firstErrorRank);
293 1 : retryCtx->needRetryServerRanks_.clear();
294 1 : CHK_PRT(SetNeedRetryServerRank(retryCtx, curOpId));
295 1 : CHK_RET(CreateOpRetryServerByState(RETRY_STATE_CMD_STOP_AICPU, retryCtx));
296 1 : return HCCL_SUCCESS;
297 1 : }
298 0 : errorRank.clear();
299 0 : SaluSleep(waitTime * TIME_MS_TO_US);
300 0 : HCCL_INFO("[OpRetry][Server]no rank can retry, wait for [%u]ms for collect all error rank", waitTime);
301 1 : }
302 : }
303 :
304 1 : HcclResult OpRetryServerHandleError::SetNeedRetryServerRank(RetryContext* retryCtx, const HcclOpIdentifier& opId)
305 : {
306 1 : if (opId.isSendRecv) {
307 : // 在send/recv场景下,仅需对本端和对端进行重执行即可
308 0 : if (retryCtx->serverSockets_.find(opId.srcRank) == retryCtx->serverSockets_.end()
309 0 : || retryCtx->serverSockets_.find(opId.detRank) == retryCtx->serverSockets_.end()) {
310 0 : HCCL_ERROR(
311 : "[OpRetry][Server]srcRank[%u] or detRank[%u] isn't in serverSockets_", opId.srcRank, opId.detRank);
312 0 : return HCCL_E_INTERNAL;
313 : }
314 0 : retryCtx->needRetryServerRanks_.push_back(opId.srcRank);
315 0 : retryCtx->needRetryServerRanks_.push_back(opId.detRank);
316 0 : retryCtx->curFaultOpId = opId;
317 0 : HCCL_INFO(
318 : "[OpRetry][Server]set needRetryServerRank[%u] for send/recv success: srcRank=[%u],detRank=[%u],"
319 : "tag =[%s], streamid =%u",
320 : retryCtx->needRetryServerRanks_.size(), opId.srcRank, opId.detRank, opId.tag, opId.streamId);
321 : } else {
322 : // 其余场景下需要对所有rank进行重执行
323 2 : for (auto& it : retryCtx->serverSockets_) {
324 1 : retryCtx->curFaultOpId = opId;
325 1 : retryCtx->needRetryServerRanks_.push_back(it.first);
326 : }
327 1 : HCCL_DEBUG("[OpRetry][Server]set needRetryServerRank[%u] success", retryCtx->needRetryServerRanks_.size());
328 : }
329 1 : return HCCL_SUCCESS;
330 : }
331 :
332 : HcclResult
333 0 : OpRetryServerRunning::ParaseErrorCode(RetryContext* retryCtx, HcclAgentRetryInfo& agentInfo, RetryState& nextState)
334 : {
335 : // 处理接收到的数据
336 0 : KfcError errorCode = agentInfo.retryInfo.opInfo.execStatus.kfcError;
337 0 : switch (errorCode) {
338 0 : case KfcError::kNone: { // 发送保活数据
339 : // 保活数据携带一个空的opid
340 0 : RetryCommandInfo commandInfo;
341 0 : commandInfo.command = RETRY_CMD_RUNNING;
342 0 : CHK_RET(IssueCommandWithOpId(agentInfo.socket, commandInfo));
343 0 : break;
344 : }
345 0 : case KfcError::kRdma:
346 0 : retryCtx->isRdmaError = true;
347 : [[fallthrough]];
348 0 : case KfcError::kExecConstraint:
349 : case KfcError::kSdma: { // 处理ERROR
350 0 : nextState = RETRY_STATE_CMD_STOP_AICPU;
351 0 : HCCL_RUN_INFO(
352 : "[OpRetry][Server]OpRetryServerRunning recv ErrorCode[%d] from rank[%u]", errorCode,
353 : agentInfo.retryInfo.rankId);
354 0 : break;
355 : }
356 0 : default: { // 不支持的ErrorCode
357 0 : HCCL_ERROR(
358 : "[OpRetry][Server]OpRetryServerRunning recv invalid ErrorCode[%d] from rank[%u]", errorCode,
359 : agentInfo.retryInfo.rankId);
360 0 : break;
361 : }
362 : }
363 0 : return HCCL_SUCCESS;
364 : }
365 :
366 0 : HcclResult OpRetryServerIssueCmd::ProcessEvent(RetryContext* retryCtx)
367 : {
368 0 : HcclResult ret = HCCL_SUCCESS;
369 0 : RetryState curState = retryCtx->GetRetryState();
370 : // 获取下一个状态
371 0 : auto itState = RETRY_SERVER_STATE_TRANSFER_LABEL.find(curState);
372 0 : CHK_PRT_RET(
373 : itState == RETRY_SERVER_STATE_TRANSFER_LABEL.end(),
374 : HCCL_ERROR(
375 : "[OpRetry][Server]OpRetryServerIssueCmd fail, state[%s] is not in RETRY_SERVER_STATE_TRANSFER_LABEL",
376 : GetReadableState(curState)),
377 : HCCL_E_INTERNAL);
378 0 : RetryState nextState = itState->second;
379 :
380 : // 发送命令
381 0 : auto itCommand = RETRY_SERVER_STATE_TO_CMD_LABEL.find(curState);
382 0 : CHK_PRT_RET(
383 : itCommand == RETRY_SERVER_STATE_TO_CMD_LABEL.end(),
384 : HCCL_ERROR(
385 : "[OpRetry][Server]OpRetryServerIssueCmd fail, state[%s] is not in RETRY_SERVER_STATE_TO_CMD_LABEL",
386 : GetReadableState(curState)),
387 : HCCL_E_INTERNAL);
388 0 : RetryCommand command = itCommand->second;
389 0 : HCCL_INFO(
390 : "[OpRetry][Server]OpRetryServerIssueCmd curState[%s], command[%s]", GetReadableState(curState),
391 : GetReadableCmd(command));
392 :
393 0 : for (auto rank : retryCtx->needRetryServerRanks_) {
394 0 : RetryCommandInfo commandInfo;
395 0 : commandInfo.command = command;
396 0 : commandInfo.opId = retryCtx->curFaultOpId;
397 0 : HCCL_INFO(
398 : "[OpRetry][Server]IssueCommandWithOpId tag[%s], index[%u], srcRank[%u], detRank[%u], isSendRecv[%d],"
399 : "streamid[%u]",
400 : commandInfo.opId.tag, commandInfo.opId.index, commandInfo.opId.srcRank, commandInfo.opId.detRank,
401 : commandInfo.opId.isSendRecv, commandInfo.opId.streamId);
402 0 : ret = IssueCommandWithOpId(retryCtx->serverSockets_[rank].socket, commandInfo);
403 0 : CHK_PRT_RET(
404 : ret != HCCL_SUCCESS,
405 : HCCL_ERROR(
406 : "[OpRetry][Server]OpRetryServerIssueCmd IssueCommand fail, curState[%s], command[%s]",
407 : GetReadableState(curState), GetReadableCmd(command)),
408 : ret);
409 : }
410 0 : CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
411 0 : return HCCL_SUCCESS;
412 : }
413 :
414 0 : HcclResult OpRetryServerWaitResp::ProcessEvent(RetryContext* retryCtx)
415 : {
416 0 : std::chrono::steady_clock::time_point startTime = std::chrono::steady_clock::now();
417 0 : const u32 timeoutValue = std::max(static_cast<u32>(GetExternalInputHcclLinkTimeOut()), OP_RETRY_SEND_RECV_TIMEOUT)
418 0 : + OP_RETRY_WAIT_AICPU_TIMEOUT;
419 0 : const std::chrono::seconds timeout = std::chrono::seconds(timeoutValue);
420 0 : RetryState curState = retryCtx->GetRetryState();
421 :
422 : // 获取预期的下一个server状态
423 0 : auto serverTransferIt = RETRY_SERVER_STATE_TRANSFER_LABEL.find(curState);
424 0 : CHK_PRT_RET(
425 : serverTransferIt == RETRY_SERVER_STATE_TRANSFER_LABEL.end(),
426 : HCCL_ERROR(
427 : "[OpRetry][Server]OpRetryServerWaitResp fail, state[%s] is not in RETRY_SERVER_STATE_TRANSFER_LABEL",
428 : GetReadableState(curState)),
429 : HCCL_E_INTERNAL);
430 0 : RetryState expectNextState = serverTransferIt->second;
431 :
432 : // 获取预期的对端agent状态
433 0 : auto agentStateIt = RETRY_SERVER_WAIT_AGENT_STATE_LABEL.find(curState);
434 0 : CHK_PRT_RET(
435 : agentStateIt == RETRY_SERVER_WAIT_AGENT_STATE_LABEL.end(),
436 : HCCL_ERROR(
437 : "[OpRetry][Server]OpRetryServerWaitResp fail, state[%s] is not in RETRY_SERVER_WAIT_AGENT_STATE_LABEL",
438 : GetReadableState(curState)),
439 : HCCL_E_INTERNAL);
440 0 : RetryState expectagentState = agentStateIt->second;
441 0 : HCCL_DEBUG(
442 : "[OpRetry][Server]OpRetryServerWaitResp state[%s], expect next state[%s], expect peer state[%s]",
443 : GetReadableState(curState), GetReadableState(expectNextState), GetReadableState(expectagentState));
444 :
445 0 : std::set<u32> recvVaild;
446 0 : while (recvVaild.size() < retryCtx->needRetryServerRanks_.size()) {
447 0 : CHK_PRT_RET(
448 : retryCtx->isServerStateWaitResume_,
449 : HCCL_RUN_INFO("[OpRetry][Server]switched state form wait resp to wait resume"), HCCL_SUCCESS);
450 0 : std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
451 0 : const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - startTime);
452 0 : CHK_PRT_RET(elapsed > timeout, HCCL_ERROR("[OpRetry][Server]OpRetryServerWaitResp timeout"), HCCL_E_TIMEOUT);
453 :
454 0 : for (auto rank : retryCtx->needRetryServerRanks_) {
455 0 : if (recvVaild.find(rank) != recvVaild.end()) {
456 0 : continue;
457 : }
458 0 : auto& agentRetryInfo = retryCtx->serverSockets_[rank];
459 : // 接收agent信息
460 0 : HcclResult ret = WaitResponse(agentRetryInfo.socket, agentRetryInfo.retryInfo);
461 0 : CHK_PRT_RET(
462 : ret != HCCL_SUCCESS && ret != HCCL_E_AGAIN,
463 : HCCL_ERROR("[OpRetry][Server]OpRetryServerWaitResp WaitResponse fail, ret[%u]", ret), ret);
464 :
465 0 : RetryState dstState = agentRetryInfo.retryInfo.retryState;
466 0 : if (ret == HCCL_SUCCESS && dstState == expectagentState) { // 接收到对端信息且状态有效
467 0 : recvVaild.insert(rank);
468 0 : HCCL_INFO(
469 : "[OpRetry][Server]OpRetryServerWaitResp recv success from dst[%u], state[%s]", rank,
470 : GetReadableState(dstState));
471 0 : } else if (ret == HCCL_SUCCESS && dstState == RETRY_STATE_RESP_RUNNING_ERR) { // 对端重执行失败
472 0 : recvVaild.insert(rank);
473 0 : PrintAgentInfoAfterFail(retryCtx->serverSockets_, recvVaild, agentRetryInfo);
474 0 : HCCL_ERROR(
475 : "[OpRetry][Server]OpRetryServerWaitResp dst rank[%u] with IpInfo[%s] retry fail, "
476 : "command all rank retry fail",
477 : rank, agentRetryInfo.retryInfo.dfxIpInfo);
478 0 : retryCtx->isNeedReportOpRetryErr = agentRetryInfo.retryInfo.isNeedReportOpRetryErr;
479 0 : HCCL_RUN_INFO(
480 : "[OpRetry][Server]OpRetryServerWaitResp retry fail, isNeedReportOpRetryErr[%d]",
481 : retryCtx->isNeedReportOpRetryErr);
482 0 : CHK_RET(CreateOpRetryServerByState(RETRY_STATE_SERVER_RETRY_FAIL, retryCtx));
483 0 : return HCCL_SUCCESS;
484 : }
485 : }
486 : }
487 :
488 0 : CHK_RET(CreateOpRetryServerByState(expectNextState, retryCtx));
489 0 : return HCCL_SUCCESS;
490 0 : }
491 :
492 0 : void OpRetryServerWaitResp::PrintAgentInfoAfterFail(
493 : std::map<u32, HcclAgentRetryInfo>& serverSockets, std::set<u32>& recvVaild, HcclAgentRetryInfo& agentRetryInfo)
494 : {
495 0 : for (auto it = serverSockets.begin(); it != serverSockets.end(); ++it) {
496 0 : if (recvVaild.find(it->first) == recvVaild.end()) { // 未接收到有效数据
497 0 : continue;
498 : }
499 0 : auto& opInfo = it->second.retryInfo.opInfo;
500 0 : const char* tag = reinterpret_cast<const char*>(opInfo.opId.tag);
501 0 : u32 index = opInfo.opId.index;
502 0 : const KfcStatus& aicpuState = opInfo.execStatus.kfcStatus;
503 0 : if (aicpuState == KfcStatus::kEnd) { // 该rank未下发算子,或算子已执行结束
504 0 : HCCL_ERROR(
505 : "[OpRetry][Server]OpRetryServerWaitResp dst[%u] with IpInfo[%s], hccl op not launch or "
506 : "is complete, hccl aicpu can not retry",
507 : it->first, it->second.retryInfo.dfxIpInfo);
508 0 : agentRetryInfo.retryInfo.isNeedReportOpRetryErr = true;
509 : }
510 0 : HCCL_RUN_INFO(
511 : "[OpRetry][Server]Print rank[%u], tag[%s], index[%u], aicpuStatus[%d]", it->first, tag, index, aicpuState);
512 : }
513 0 : }
514 :
515 1 : HcclResult OpRetryServerCheckOp::ProcessEvent(RetryContext* retryCtx)
516 : {
517 1 : HcclResult ret = CheckRetryInfo(*retryCtx);
518 1 : RetryState nextState = (ret == HCCL_SUCCESS) ? RETRY_STATE_CMD_CHECK_LINK : RETRY_STATE_SERVER_RETRY_FAIL;
519 :
520 1 : if (ret == HCCL_E_OPRETRY_FAIL) {
521 1 : HCCL_RUN_INFO("[OpRetry][Server][CheckRetryInfo] Opname is Inconsistent, RETRY_CONSTRAINT, ret[%u]", ret);
522 1 : retryCtx->isNeedReportOpRetryErr = true;
523 : }
524 :
525 1 : HCCL_RUN_INFO("[OpRetry][Server]check op ret[%d], nextState[%s]", ret, GetReadableState(nextState));
526 1 : CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
527 1 : return HCCL_SUCCESS;
528 : }
529 :
530 0 : HcclResult OpRetryServerWaitLinkInfo::ProcessEvent(RetryContext* retryCtx)
531 : {
532 0 : std::chrono::steady_clock::time_point startTime = std::chrono::steady_clock::now();
533 0 : const u32 timeoutValue = std::max(static_cast<u32>(GetExternalInputHcclLinkTimeOut()), OP_RETRY_SEND_RECV_TIMEOUT)
534 0 : + OP_RETRY_WAIT_AICPU_TIMEOUT;
535 0 : const std::chrono::seconds timeout = std::chrono::seconds(timeoutValue);
536 : // 下一个server状态
537 0 : RetryState nextState = RETRY_STATE_CHECK_ALL_LINK;
538 :
539 0 : std::set<u32> recvVaild;
540 0 : while (recvVaild.size() < retryCtx->needRetryServerRanks_.size()) {
541 0 : CHK_PRT_RET(
542 : retryCtx->isServerStateWaitResume_,
543 : HCCL_RUN_INFO("[OpRetry][Server]switched state form wait link to wait resume"), HCCL_SUCCESS);
544 0 : std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
545 0 : const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - startTime);
546 0 : CHK_PRT_RET(
547 : elapsed > timeout, HCCL_ERROR("[OpRetry][Server]OpRetryServerWaitLinkInfo timeout"), HCCL_E_TIMEOUT);
548 :
549 0 : for (auto rank : retryCtx->needRetryServerRanks_) {
550 0 : if (recvVaild.find(rank) != recvVaild.end()) {
551 0 : continue;
552 : }
553 0 : auto& agentRetryInfo = retryCtx->serverSockets_[rank];
554 : // 接收agent信息
555 0 : HcclResult ret = WaitLinkPortCheckResult(agentRetryInfo.socket, agentRetryInfo.linkPortStatus);
556 0 : CHK_PRT_RET(
557 : ret != HCCL_SUCCESS && ret != HCCL_E_AGAIN,
558 : HCCL_ERROR("[OpRetry][Server]OpRetryServerWaitLinkCheckResult fail, ret[%u]", ret), ret);
559 0 : if (ret == HCCL_SUCCESS) {
560 0 : recvVaild.insert(rank);
561 0 : HCCL_INFO("[OpRetry][Server]OpRetryServerWaitLinkCheckResult recv success from dst[%u], ", rank);
562 : }
563 : }
564 : }
565 0 : CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
566 0 : HCCL_INFO("[OpRetry][Server]OpRetryServerWaitLinkInfo success");
567 0 : return HCCL_SUCCESS;
568 0 : }
569 :
570 0 : HcclResult OpRetryServerCheckAllLink::ProcessEvent(RetryContext* retryCtx)
571 : {
572 : // 收集所有rank的主备网口信息
573 0 : std::map<u32, std::pair<bool, bool>> allLinkInfo;
574 0 : for (auto rank : retryCtx->needRetryServerRanks_) {
575 0 : auto& linkPortStatus = retryCtx->serverSockets_[rank].linkPortStatus;
576 0 : allLinkInfo.insert({rank, std::make_pair(linkPortStatus.defaultPort, linkPortStatus.backupPort)});
577 : }
578 :
579 : // 对所有rank依次遍历
580 0 : for (auto rank : retryCtx->needRetryServerRanks_) {
581 0 : u32 remoteRankIndex = 0;
582 0 : auto& linkPortStatus = retryCtx->serverSockets_[rank].linkPortStatus;
583 : // 对rank的所有对端进行遍历
584 0 : for (u32 i = 0; i < linkPortStatus.rankSize; i++) {
585 0 : u32 remoteRank = linkPortStatus.rankList[i];
586 0 : retryCtx->serverSockets_[rank].changeLinkInfo.remoteRankList[remoteRankIndex] = remoteRank;
587 0 : if (allLinkInfo[rank].first && allLinkInfo[remoteRank].first) {
588 : // 本端和对端的主网口均up,则使用主网口
589 0 : retryCtx->serverSockets_[rank].changeLinkInfo.isUseDefaultPort[remoteRankIndex] = true;
590 0 : } else if (allLinkInfo[rank].second && allLinkInfo[remoteRank].second) {
591 : // 本端和对端的备网口均up,则使用备网口
592 0 : retryCtx->serverSockets_[rank].changeLinkInfo.isUseDefaultPort[remoteRankIndex] = false;
593 : } else {
594 : // 本端和对端无可用的网口,重执行失败
595 0 : HCCL_ERROR(
596 : "[OpRetry][Server]rank[%u]:default[%d], backup[%d], IpInfo[%s]; rank[%u]:default[%d], "
597 : "backup[%d], can not find same port, can not retry",
598 : rank, allLinkInfo[rank].first, allLinkInfo[rank].second,
599 : retryCtx->serverSockets_[rank].retryInfo.dfxIpInfo, remoteRank, allLinkInfo[remoteRank].first,
600 : allLinkInfo[remoteRank].second);
601 0 : CHK_RET(CreateOpRetryServerByState(RETRY_STATE_SERVER_RETRY_FAIL, retryCtx));
602 0 : return HCCL_SUCCESS;
603 : }
604 0 : remoteRankIndex += 1;
605 : }
606 0 : retryCtx->serverSockets_[rank].changeLinkInfo.remoteRankNum = remoteRankIndex;
607 : }
608 :
609 : // 打印所有rank的借轨信息
610 0 : for (auto rank : retryCtx->needRetryServerRanks_) {
611 0 : auto& changeLinkInfo = retryCtx->serverSockets_[rank].changeLinkInfo;
612 0 : std::string changeLinkInfoStr = "rank[" + std::to_string(rank) + "]";
613 0 : for (u32 i = 0; i < changeLinkInfo.remoteRankNum; i++) {
614 : changeLinkInfoStr
615 0 : += (std::to_string(changeLinkInfo.remoteRankList[i]) + ":"
616 0 : + std::to_string(changeLinkInfo.isUseDefaultPort[i]) + "; ");
617 : }
618 0 : HCCL_INFO("[OpRetry][Server]changeLinkInfoStr:%s", changeLinkInfoStr.c_str());
619 0 : }
620 :
621 : // 所有rank网口确认成功,切换到给agent发借轨命令状态
622 0 : CHK_RET(CreateOpRetryServerByState(RETRY_STATE_CMD_RESUME_TRANSPORT, retryCtx));
623 0 : return HCCL_SUCCESS;
624 0 : }
625 :
626 0 : HcclResult OpRetryServerIssueChangeLinkAndResume::ProcessEvent(RetryContext* retryCtx)
627 : {
628 0 : HcclResult ret = HCCL_SUCCESS;
629 0 : RetryState curState = retryCtx->GetRetryState();
630 : // 先将每个rank的changeLinkInfo发送至对应agent
631 0 : for (auto rank : retryCtx->needRetryServerRanks_) {
632 0 : ret = IssueChangeLink(retryCtx->serverSockets_[rank].socket, retryCtx->serverSockets_[rank].changeLinkInfo);
633 0 : CHK_PRT_RET(
634 : ret != HCCL_SUCCESS,
635 : HCCL_ERROR("[OpRetry][Server]OpRetryServerIssueChangeLink fail, curState[%s]", GetReadableState(curState)),
636 : ret);
637 0 : HCCL_INFO("[OpRetry][Server]OpRetryServerIssueChangeLink send to rank[%u] success", rank);
638 : }
639 : // 再发送resume transport命令至每个agent
640 0 : std::shared_ptr<OpRetryBase> retryPtr = nullptr;
641 0 : EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerIssueCmd>()), return HCCL_E_PTR);
642 0 : RetryState nextState = RETRY_STATE_CMD_RESUME_TRANSPORT;
643 0 : retryCtx->SetRetryState(nextState, retryPtr);
644 0 : HCCL_INFO("[OpRetry][Server]OpRetryServerIssueChangeLinkAndResume success");
645 0 : return HCCL_SUCCESS;
646 0 : }
647 :
648 0 : HcclResult OpRetryServerRetryFail::ProcessEvent(RetryContext* retryCtx)
649 : {
650 0 : RetryCommandInfo commandInfo;
651 0 : commandInfo.command = RETRY_CMD_RETRY_FAIL;
652 0 : if (retryCtx->isNeedReportOpRetryErr) {
653 0 : commandInfo.command = RETRY_CMD_RETRY_CONSTRAINT_FAIL;
654 0 : HCCL_RUN_INFO(
655 : "[OpRetry][Server]OpRetryServerRetryFail isNeedReportOpRetryErr[%d], command[%s]",
656 : retryCtx->isNeedReportOpRetryErr, GetReadableCmd(commandInfo.command));
657 : }
658 0 : HCCL_INFO("[OpRetry][Server]could not retry, command all rank %s", GetReadableCmd(commandInfo.command));
659 0 : for (auto rank : retryCtx->needRetryServerRanks_) {
660 0 : HcclResult ret = IssueCommandWithOpId(retryCtx->serverSockets_[rank].socket, commandInfo);
661 0 : CHK_PRT_RET(
662 : ret != HCCL_SUCCESS,
663 : HCCL_ERROR("[OpRetry][Server]OpRetryServerRetryFail IssueCommandWithOpId AgentId[%u] fail", rank), ret);
664 : }
665 :
666 : // 重执行异常,通知心跳存在异常 -> 广播异常给整个集群
667 0 : Heartbeat::GetInstance(retryCtx->deviceLogicId_).SetOpretryErr();
668 :
669 0 : RetryState nextState = RETRY_STATE_SERVER_RUNNING;
670 0 : CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
671 0 : return HCCL_SUCCESS;
672 : }
673 :
674 15 : bool SwitchNicServerCheckAllSwitchRanks::CompareSwitchRankList(
675 : const u32* firstSwitchRankList, const u32* switchRankList, const u32 switchRankNum)
676 : {
677 15 : if (switchRankNum == 0 || switchRankNum > AICPU_MAX_RANK_NUM) {
678 0 : return false;
679 : }
680 :
681 15 : std::set<u32> switchRankSet;
682 45 : for (u32 i = 0; i < switchRankNum; i++) {
683 30 : switchRankSet.insert(firstSwitchRankList[i]);
684 : }
685 :
686 15 : u32 ranksNum[AICPU_MAX_RANK_NUM] = {0};
687 43 : for (u32 i = 0; i < switchRankNum; i++) {
688 30 : if (switchRankSet.find(switchRankList[i]) == switchRankSet.end()) {
689 1 : HCCL_ERROR("[SwitchNic][Server] rankList has error, id[%u]", switchRankList[i]);
690 1 : return false;
691 : } else {
692 29 : ranksNum[switchRankList[i]]++;
693 : }
694 29 : if (ranksNum[switchRankList[i]] > 1) {
695 1 : HCCL_ERROR(
696 : "[SwitchNic][Server] rankList has error, id[%u], num[%u]", switchRankList[i],
697 : ranksNum[switchRankList[i]]);
698 1 : return false;
699 : }
700 : }
701 13 : return true;
702 15 : }
703 :
704 13 : bool SwitchNicServerCheckAllSwitchRanks::CompareUseBackupLists(
705 : const bool* firstArray, const bool* secondArray, const u32 switchRankNum)
706 : {
707 13 : if (switchRankNum == 0 || switchRankNum > AICPU_MAX_RANK_NUM) {
708 0 : return false;
709 : }
710 37 : for (u32 i = 0; i < switchRankNum; i++) {
711 25 : if (firstArray[i] != secondArray[i]) {
712 1 : HCCL_ERROR(
713 : "[SwitchNic][Server] backupLists has error first[%u], second[%u], index[%u]", firstArray[i],
714 : secondArray[i], i);
715 1 : return false;
716 : }
717 : }
718 12 : return true;
719 : }
720 :
721 12 : bool SwitchNicServerCheckAllSwitchRanks::CheckRemotePorts(const u32 rankId, const ActiveSwitchInfo& switchRankInfo)
722 : {
723 28 : for (u32 i = 0; i < switchRankInfo.remoteRankNum; i++) {
724 18 : if (switchRankInfo.remoteRankNicStatus[i] == CONNECT_REMOTE_DEFAULT && !switchRankInfo.defaultPortStatus) {
725 1 : HCCL_ERROR(
726 : "[SwitchNic][Server] defaultPortStatus has error, localRank[%u], remoteRank[%u], nicStatus[%u]", rankId,
727 : i, switchRankInfo.remoteRankNicStatus[i]);
728 1 : return false;
729 : }
730 17 : if (switchRankInfo.remoteRankNicStatus[i] == CONNECT_REMOTE_BACKUP && !switchRankInfo.backupPortStatus) {
731 1 : HCCL_ERROR(
732 : "[SwitchNic][Server] backupPortStatus has error, localRank[%u], remoteRank[%u], nicStatus[%u]", rankId,
733 : i, switchRankInfo.remoteRankNicStatus[i]);
734 1 : return false;
735 : }
736 : }
737 10 : return true;
738 : }
739 :
740 3 : HcclResult SwitchNicServerCheckAllSwitchRanks::CollectSingleAgentActiveSwitchInfo(
741 : RetryContext* retryCtx, const u32 rankId, HcclAgentRetryInfo& agentInfo)
742 : {
743 3 : std::chrono::steady_clock::time_point startTime = std::chrono::steady_clock::now();
744 3 : const auto timeout = std::chrono::seconds(GetExternalInputHcclLinkTimeOut() * ACTIVE_SWITCH_TIMES);
745 3 : HcclResult ret = HCCL_SUCCESS;
746 : while (true) {
747 6 : CHK_PRT_RET(
748 : retryCtx->isServerStateWaitResume_,
749 : HCCL_RUN_INFO("[OpRetry][Server]switched state form check switch Nic to wait resume"), HCCL_SUCCESS);
750 3 : std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
751 3 : const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - startTime);
752 3 : CHK_PRT_RET(
753 : elapsed > timeout,
754 : HCCL_ERROR("[SwitchNic][Server] timeout in recv agent RetryInfo, waitime[%u s>%u s]", elapsed, timeout),
755 : HCCL_E_TIMEOUT);
756 3 : ret = WaitResponse(agentInfo.socket, agentInfo.retryInfo);
757 3 : if (ret == HCCL_SUCCESS) { // 成功接收到数据
758 : // 成功接收到retryInfo含RETRY_STATE_SEND_SWITCH_INFO,接收 ActiveSwitchInfo,否则为保活数据,忽略
759 3 : if (agentInfo.retryInfo.retryState == RETRY_STATE_SEND_SWITCH_INFO) {
760 3 : ret = RecvActiveSwitchInfo(agentInfo.socket, rankId, agentInfo.switchInfo);
761 3 : if (ret != HCCL_SUCCESS) {
762 0 : return ret;
763 : }
764 3 : HCCL_INFO("[SwitchNic][server] recv ActiveSwitchInfo form rank[%u] while collecting", rankId);
765 3 : retryCtx->switchInfoMap_[rankId] = agentInfo.switchInfo;
766 3 : return HCCL_SUCCESS;
767 : }
768 0 : } else if (ret == HCCL_E_AGAIN) {
769 0 : RetryCommand command = RETRY_CMD_RUNNING;
770 0 : CHK_RET(IssueCommand(agentInfo.socket, command));
771 0 : HCCL_DEBUG("[SwitchNic][Server] send keeping active info to dst[%u]", rankId);
772 : } else {
773 0 : HCCL_ERROR("[SwitchNic][Server] get active switch info failed, ret[%u], dst[%u]", ret, rankId);
774 0 : break;
775 : }
776 0 : SaluSleep(OP_RETRY_POLL_AICPU_STATE_INTERVAL);
777 0 : }
778 0 : return ret;
779 : }
780 :
781 9 : HcclResult SwitchNicServerCheckAllSwitchRanks::CollectAgentActiveSwitchInfo(RetryContext* retryCtx)
782 : {
783 9 : HCCL_RUN_INFO("[SwitchNic][Server] began to CollectAgentActiveSwitchInfo");
784 9 : HcclResult ret = HCCL_SUCCESS;
785 : // 轮询接收agent信息
786 27 : for (auto& it : retryCtx->serverSockets_) {
787 18 : const u32& rank = it.first;
788 18 : if (retryCtx->switchInfoMap_.find(rank) != retryCtx->switchInfoMap_.end()) {
789 15 : HCCL_DEBUG("[SwitchNic][Server] rank[%u] has been received", rank);
790 15 : continue;
791 : }
792 : // 轮询接收agent状态机信息
793 3 : ret = CollectSingleAgentActiveSwitchInfo(retryCtx, it.first, it.second);
794 3 : if (ret != HCCL_SUCCESS) {
795 0 : return ret;
796 : }
797 : }
798 : // 可不检查,理论上不会不等于
799 9 : if (retryCtx->switchInfoMap_.size() != retryCtx->serverSockets_.size()) {
800 0 : return HCCL_E_UNAVAIL;
801 : }
802 9 : return ret;
803 : }
804 :
805 9 : HcclResult SwitchNicServerCheckAllSwitchRanks::CheckAgentActiveSwitchInfo(RetryContext* retryCtx)
806 : {
807 9 : HCCL_RUN_INFO("[SwitchNic][Server] began to CheckAgentActiveSwitchInfo");
808 9 : if (retryCtx->switchInfoMap_.empty()) {
809 0 : HCCL_ERROR("[SwitchNic][Server] switchInfoMap is empty");
810 0 : return HCCL_E_PARA;
811 : }
812 :
813 9 : RetryCommand command = RETRY_CMD_NOTIFY_SWITCH_SUC;
814 :
815 9 : auto firstInfo = retryCtx->switchInfoMap_.begin();
816 9 : auto firstRankId = firstInfo->first;
817 9 : ActiveSwitchInfo& firstSwitchInfo = firstInfo->second;
818 :
819 19 : for (const auto& it : retryCtx->switchInfoMap_) {
820 17 : const u32& rankId = it.first;
821 17 : const ActiveSwitchInfo& switchInfo = it.second;
822 17 : if (!switchInfo.refreshTransportFin) {
823 1 : HCCL_ERROR("[SwitchNic][Server] refreshTransportFin is false, first[%u], rank[%u]", firstRankId, rankId);
824 1 : command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
825 1 : break;
826 : }
827 16 : if (switchInfo.switchRankNum != firstSwitchInfo.switchRankNum) {
828 1 : HCCL_ERROR(
829 : "[SwitchNic][Server] switchRankNum is not same as the first[%u:%u], rank[%u:%u]", firstRankId,
830 : firstSwitchInfo.switchRankNum, rankId, switchInfo.switchRankNum);
831 1 : command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
832 1 : break;
833 : } else {
834 15 : if (!CompareSwitchRankList(
835 15 : firstSwitchInfo.switchRankList, switchInfo.switchRankList, switchInfo.switchRankNum)) {
836 2 : HCCL_ERROR(
837 : "[SwitchNic][Server] SwitchRankList is not same as the first[%u], rank[%u], rankNum[%u]",
838 : firstRankId, rankId, switchInfo.switchRankNum);
839 2 : command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
840 2 : break;
841 : }
842 13 : if (!CompareUseBackupLists(
843 13 : firstSwitchInfo.switchUseBackup, switchInfo.switchUseBackup, switchInfo.switchRankNum)) {
844 1 : HCCL_ERROR(
845 : "[SwitchNic][Server] UseBackupLists is not same as the first[%u], rank[%u], rankNum[%u]",
846 : firstRankId, rankId, switchInfo.switchRankNum);
847 1 : command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
848 1 : break;
849 : }
850 : }
851 :
852 12 : if (!switchInfo.localPortsCheckRet) {
853 0 : HCCL_ERROR("[SwitchNic][Server] localPortsCheckRet is false, first[%u], rank[%u]", firstRankId, rankId);
854 0 : command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
855 0 : break;
856 : }
857 12 : if (!CheckRemotePorts(rankId, switchInfo)) {
858 2 : command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
859 2 : break;
860 : }
861 : }
862 :
863 : // 全部卡确认无误后或者发现错误后,通知OpRetryAgent。
864 27 : for (auto it : retryCtx->serverSockets_) {
865 18 : HcclResult ret = IssueCommand(it.second.socket, command);
866 18 : CHK_PRT_RET(
867 : ret != HCCL_SUCCESS,
868 : HCCL_ERROR("[SwitchNic][Server] CheckAllSwitchRanks IssueCommand AgentId[%u] fail", it.first), ret);
869 18 : }
870 9 : retryCtx->switchInfoMap_.clear();
871 9 : return HCCL_SUCCESS;
872 : }
873 :
874 : // OpRetryServer遍历通信域内的所有卡,接收主动借轨信息,并且校验每张卡信息一致
875 : // 全部卡确认无误后或者发现错误后,通知OpRetryAgent
876 9 : HcclResult SwitchNicServerCheckAllSwitchRanks::ProcessEvent(RetryContext* retryCtx)
877 : {
878 9 : HCCL_RUN_INFO("[SwitchNic][Server] CheckAllSwitchRanks begin");
879 9 : RetryState nextState = RETRY_STATE_SERVER_RUNNING;
880 :
881 9 : CHK_RET(CollectAgentActiveSwitchInfo(retryCtx));
882 9 : CHK_RET(CheckAgentActiveSwitchInfo(retryCtx));
883 9 : CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
884 9 : HCCL_RUN_INFO("[SwitchNic][Server] CheckAllSwitchRanks end");
885 9 : return HCCL_SUCCESS;
886 : }
887 :
888 0 : HcclResult OpRetryServerWaitResume::ProcessEvent(RetryContext* retryCtx)
889 : {
890 0 : if (!retryCtx->isServerStateWaitResume_ && !retryCtx->haveCommEnableBackupLink_) {
891 0 : CHK_RET(CreateOpRetryServerByState(RETRY_STATE_SERVER_RUNNING, retryCtx));
892 0 : HCCL_RUN_INFO(
893 : "[OpRetry][Server]OpRetryServerWaitResume, group[%s], no comm enable backup link, set state to running",
894 : retryCtx->group_.c_str());
895 0 : return HCCL_SUCCESS;
896 : }
897 0 : if (!retryCtx->isServerStateWaitResume_ && retryCtx->haveCommEnableBackupLink_) {
898 0 : HCCL_RUN_INFO("[OpRetry][Server]OpRetryServerWaitResume, start to send cmd");
899 0 : for (auto& it : retryCtx->serverSockets_) {
900 0 : const u32& agentId = it.first;
901 0 : RetryCommandInfo commandInfo;
902 0 : commandInfo.command = RESUME_CMD_CHECK_LINK;
903 0 : CHK_PRT_RET(
904 : IssueCommandWithOpId(it.second.socket, commandInfo),
905 : HCCL_ERROR("[OpRetry][Server][Resume]rank[%u] send resume check link fail", agentId), HCCL_E_INTERNAL);
906 0 : HCCL_RUN_INFO(
907 : "[OpRetry][Server][Resume]rank[%u] send RESUME_CMD_CHECK_LINK, group[%s]", agentId,
908 : retryCtx->group_.c_str());
909 : }
910 0 : CHK_RET(CreateOpRetryServerByState(RETRY_RESUME_STATE_SERVER_CHECK_LINK, retryCtx));
911 0 : HCCL_RUN_INFO("[OpRetry][Server]OpRetryServerWaitResume, set state to check link");
912 0 : retryCtx->isRdmaError = false;
913 0 : return HCCL_SUCCESS;
914 : }
915 :
916 0 : return HCCL_SUCCESS;
917 : }
918 :
919 0 : HcclResult ResumeServerCheckAllLink::ProcessEvent(RetryContext* retryCtx)
920 : {
921 0 : HCCL_RUN_INFO("[OpRetry][Server]ResumeServerCheckAllLink begin group[%s]", retryCtx->group_.c_str());
922 0 : RetryState nextState = RETRY_RESUME_STATE_SERVER_CHANGE_LINK;
923 : // 收集所有agent的检查网口结果
924 0 : CHK_RET(WaitAgentCheckLinkResult(retryCtx));
925 : // 检查所有rank的主备链路连接情况
926 0 : CHK_RET(CheckAllLink(retryCtx, nextState));
927 :
928 0 : CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
929 0 : return HCCL_SUCCESS;
930 : }
931 :
932 0 : HcclResult ResumeServerChangeLink::ProcessEvent(RetryContext* retryCtx)
933 : {
934 0 : HCCL_RUN_INFO("[OpRetry][Server]ResumeServerChangeLink begin");
935 0 : RetryState nextState = RETRY_STATE_SERVER_RUNNING;
936 : // 下发切换链路命令字到所有rank,并发送重建transport命令
937 0 : CHK_RET(CmdAgentChangeLink(retryCtx));
938 :
939 : // 接收所有rank的切换链路结果
940 0 : CHK_RET(WaitAllChangeLinkResult(retryCtx, nextState));
941 0 : if (nextState != RETRY_STATE_SERVER_RUNNING) {
942 0 : HCCL_ERROR("[OpRetry][Server]ResumeServerChangeLink fail, nextState[%s]", GetReadableState(nextState));
943 : }
944 0 : CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
945 0 : return HCCL_SUCCESS;
946 : }
947 :
948 1 : HcclResult ResumeServerCheckAllLink::WaitAgentCheckLinkResult(RetryContext* retryCtx)
949 : {
950 1 : HCCL_RUN_INFO("[OpRetry][Server][Resume]WaitAgentCheckLinkResult begin");
951 1 : std::chrono::steady_clock::time_point startTime = std::chrono::steady_clock::now();
952 1 : const u32 timeoutValue = std::max(static_cast<u32>(GetExternalInputHcclLinkTimeOut()), OP_RETRY_SEND_RECV_TIMEOUT)
953 1 : + OP_RETRY_WAIT_AICPU_TIMEOUT;
954 1 : const std::chrono::seconds timeout = std::chrono::seconds(timeoutValue);
955 1 : std::set<u32> recvVaild;
956 2 : while (recvVaild.size() < retryCtx->serverSockets_.size()) {
957 1 : std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
958 1 : const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - startTime);
959 1 : CHK_PRT_RET(
960 : elapsed > timeout, HCCL_ERROR("[OpRetry][Server][Resume]WaitAgentCheckLinkResult timeout"), HCCL_E_TIMEOUT);
961 :
962 2 : for (auto& rank : retryCtx->serverSockets_) {
963 1 : const u32& agentId = rank.first;
964 1 : if (recvVaild.find(agentId) != recvVaild.end()) {
965 0 : continue;
966 : }
967 1 : auto& agentRetryInfo = rank.second;
968 1 : HcclResult ret = WaitLinkPortCheckResult(agentRetryInfo.socket, agentRetryInfo.linkPortStatus);
969 1 : CHK_PRT_RET(
970 : ret != HCCL_SUCCESS && ret != HCCL_E_AGAIN,
971 : HCCL_ERROR("[OpRetry][Server][Resume]WaitAgentCheckLinkResult WaitLinkPortCheckResult Failed"), ret);
972 1 : if (ret == HCCL_SUCCESS) {
973 1 : recvVaild.insert(agentId);
974 1 : HCCL_RUN_INFO(
975 : "[OpRetry][Server][Resume]WaitAgentCheckLinkResult recv valid from agentId[%u], "
976 : "group[%s],remoteRank[%u]",
977 : agentId, retryCtx->group_.c_str(), rank.second.linkPortStatus.rankList[0]);
978 : }
979 : }
980 : }
981 1 : HCCL_INFO("[OpRetry][Server][Resume]WaitAgentCheckLinkResult recv all valid");
982 1 : return HCCL_SUCCESS;
983 1 : }
984 :
985 0 : HcclResult ResumeServerCheckAllLink::CheckAllLink(RetryContext* retryCtx, RetryState& nextState)
986 : {
987 0 : std::map<u32, std::pair<bool, bool>> allLinkInfo;
988 0 : for (auto it : retryCtx->serverSockets_) {
989 0 : u32 rank = it.first;
990 0 : auto& linkPortStatus = retryCtx->serverSockets_[rank].linkPortStatus;
991 0 : HCCL_RUN_INFO(
992 : "[OpRetry][Server][Resume]CheckAllLink rank[%u], rankListSize[%d], rankSize[%d]", rank,
993 : std::end(it.second.linkPortStatus.rankList) - std::begin(it.second.linkPortStatus.rankList),
994 : it.second.linkPortStatus.rankSize);
995 0 : allLinkInfo.insert({it.first, std::make_pair(linkPortStatus.defaultPort, linkPortStatus.backupPort)});
996 0 : }
997 :
998 0 : for (auto it : retryCtx->serverSockets_) {
999 0 : u32 rank = it.first;
1000 0 : u32 remoteRankIndex = 0;
1001 0 : auto& linkPortStatus = it.second.linkPortStatus;
1002 0 : for (u32 i = 0; i < linkPortStatus.rankSize; i++) {
1003 0 : u32 remoteRank = linkPortStatus.rankList[i];
1004 0 : retryCtx->serverSockets_[rank].changeLinkInfo.remoteRankList[remoteRankIndex] = remoteRank;
1005 0 : if (allLinkInfo[rank].first && allLinkInfo[remoteRank].first) {
1006 : // 本端与对端的主网口均up, 则使用主网口
1007 0 : retryCtx->serverSockets_[rank].changeLinkInfo.isUseDefaultPort[remoteRankIndex] = true;
1008 0 : } else if (allLinkInfo[rank].second && allLinkInfo[remoteRank].second) {
1009 0 : retryCtx->serverSockets_[rank].changeLinkInfo.isUseDefaultPort[remoteRankIndex] = false;
1010 : } else {
1011 0 : HCCL_ERROR(
1012 : "[OpRetry][Server][Resume]rank[%u]:default[%d], backup[%d], IpInfo[%s]; "
1013 : "remoterank[%u]:default[%d], "
1014 : "backup[%d], can not find same port, can not resume",
1015 : rank, allLinkInfo[rank].first, allLinkInfo[rank].second,
1016 : retryCtx->serverSockets_[rank].retryInfo.dfxIpInfo, remoteRank, allLinkInfo[remoteRank].first,
1017 : allLinkInfo[remoteRank].second);
1018 0 : nextState = RETRY_STATE_SERVER_RETRY_FAIL;
1019 0 : return HCCL_SUCCESS;
1020 : }
1021 0 : HCCL_RUN_INFO(
1022 : "[OpRetry][Server][Resume]CheckAllLink remoteRank[%u], changeLinkInfo remoteRankList[%u], "
1023 : "linkPortStatus rankList[%u]",
1024 : remoteRank, retryCtx->serverSockets_[rank].changeLinkInfo.remoteRankList[0],
1025 : linkPortStatus.rankList[0]);
1026 0 : remoteRankIndex++;
1027 : }
1028 0 : retryCtx->serverSockets_[rank].changeLinkInfo.remoteRankNum = remoteRankIndex;
1029 0 : }
1030 :
1031 : // 打印所有rank的借轨信息
1032 0 : for (auto it : retryCtx->serverSockets_) {
1033 0 : u32 rank = it.first;
1034 0 : auto& changeLinkInfo = it.second.changeLinkInfo;
1035 0 : std::string changeLinkInfoStr = "rank[" + std::to_string(rank) + "]";
1036 0 : for (u32 i = 0; i < changeLinkInfo.remoteRankNum; i++) {
1037 : changeLinkInfoStr
1038 0 : += (std::to_string(changeLinkInfo.remoteRankList[i]) + ":"
1039 0 : + std::to_string(changeLinkInfo.isUseDefaultPort[i]) + "; ");
1040 : }
1041 0 : HCCL_INFO("[OpRetry][Server][Resume]changeLinkInfoStr:%s", changeLinkInfoStr.c_str());
1042 0 : }
1043 0 : return HCCL_SUCCESS;
1044 0 : }
1045 :
1046 0 : HcclResult ResumeServerChangeLink::CmdAgentChangeLink(RetryContext* retryCtx)
1047 : {
1048 0 : HCCL_RUN_INFO("[OpRetry][Server][Resume]CmdAgentChangeLink begin");
1049 : // 先将每个rank的changeLinkInfo发送至对应agent
1050 0 : for (auto it : retryCtx->serverSockets_) {
1051 0 : HcclResult ret = IssueChangeLink(it.second.socket, it.second.changeLinkInfo);
1052 0 : CHK_PRT_RET(
1053 : ret != HCCL_SUCCESS,
1054 : HCCL_ERROR("[OpRetry][Server][Resume]CmdAgentChangeLink IssueCommandChangeLink RankId[%u] fail", it.first),
1055 : ret);
1056 0 : HCCL_INFO(
1057 : "[OpRetry][Server]OpRetryServerIssueChangeLink send ChangeLinkInfo to rank[%u] success, group[%s]",
1058 : it.first, retryCtx->group_.c_str());
1059 0 : }
1060 0 : for (auto& it : retryCtx->serverSockets_) {
1061 0 : const u32& agentId = it.first;
1062 0 : RetryCommandInfo commandInfo;
1063 0 : commandInfo.command = RETRY_CMD_RESUME_TRANSPORT;
1064 0 : HcclResult ret = IssueCommandWithOpId(it.second.socket, commandInfo);
1065 0 : CHK_PRT_RET(
1066 : ret != HCCL_SUCCESS,
1067 : HCCL_ERROR("[OpRetry][Server][Resume]CmdAgentChangeLink IssueCommand AgentId[%u] fail", agentId), ret);
1068 0 : HCCL_INFO(
1069 : "[OpRetry][Server]OpRetryServerIssueChangeLink send change link command to rank[%u] success, group[%s]",
1070 : it.first, retryCtx->group_.c_str());
1071 : }
1072 0 : return HCCL_SUCCESS;
1073 : }
1074 :
1075 0 : HcclResult ResumeServerChangeLink::WaitAllChangeLinkResult(RetryContext* retryCtx, RetryState& nextState)
1076 : {
1077 0 : HCCL_RUN_INFO("[OpRetry][Server][Resume]WaitAllChangeLinkResult begin group[%s]", retryCtx->group_.c_str());
1078 0 : std::chrono::steady_clock::time_point startTime = std::chrono::steady_clock::now();
1079 0 : const u32 timeoutValue = std::max(static_cast<u32>(GetExternalInputHcclLinkTimeOut()), OP_RETRY_SEND_RECV_TIMEOUT)
1080 0 : + OP_RETRY_WAIT_AICPU_TIMEOUT;
1081 0 : const std::chrono::seconds timeout = std::chrono::seconds(timeoutValue);
1082 0 : RetryState expectAgentState = RETRY_STATE_AGENT_RUNNING;
1083 0 : std::set<u32> recvValid;
1084 0 : while (recvValid.size() < retryCtx->serverSockets_.size()) {
1085 0 : std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
1086 0 : const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - startTime);
1087 0 : CHK_PRT_RET(
1088 : elapsed > timeout, HCCL_ERROR("[OpRetry][Server][Resume]WaitAllChangeLinkResult timeout"), HCCL_E_TIMEOUT);
1089 :
1090 0 : for (auto& it : retryCtx->serverSockets_) {
1091 0 : u32 rank = it.first;
1092 0 : if (recvValid.find(rank) != recvValid.end()) {
1093 0 : continue;
1094 : }
1095 0 : auto& agentRetryInfo = retryCtx->serverSockets_[rank];
1096 0 : HcclResult ret = WaitResponse(agentRetryInfo.socket, agentRetryInfo.retryInfo);
1097 0 : RetryState dstState = agentRetryInfo.retryInfo.retryState;
1098 0 : KfcStatus aicpuState = agentRetryInfo.retryInfo.opInfo.execStatus.kfcStatus;
1099 0 : if (ret == HCCL_SUCCESS && dstState == expectAgentState && aicpuState == KfcStatus::kResumeChanged) {
1100 0 : recvValid.insert(rank);
1101 0 : HCCL_INFO("[OpRetry][Server][Resume]WaitAllChangeLinkResult recv valid from rank[%u]", rank);
1102 0 : } else if (ret == HCCL_SUCCESS && dstState == RETRY_STATE_RESP_RUNNING_ERR) {
1103 0 : recvValid.insert(rank);
1104 0 : nextState = RETRY_STATE_SERVER_RETRY_FAIL;
1105 0 : HCCL_ERROR("[OpRetry][Server][Resume]WaitAllChangeLinkResult recv err from rank[%u]", rank);
1106 0 : return HCCL_SUCCESS;
1107 : }
1108 : }
1109 : }
1110 0 : return HCCL_SUCCESS;
1111 0 : }
1112 : } // namespace hccl
|