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