LCOV - code coverage report
Current view: top level - legacy/ascend910/framework/cluster_maintenance/recovery/operator_retry - opretry_server.cc (source / functions) Coverage Total Hit
Test: coverage.info Lines: 33.9 % 631 214
Test Date: 2026-08-25 19:18:03 Functions: 46.4 % 28 13

            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,
     494              :     HcclAgentRetryInfo& agentRetryInfo) const
     495              : {
     496            0 :     for (auto it = serverSockets.begin(); it != serverSockets.end(); ++it) {
     497            0 :         if (recvVaild.find(it->first) == recvVaild.end()) { // 未接收到有效数据
     498            0 :             continue;
     499              :         }
     500            0 :         auto& opInfo = it->second.retryInfo.opInfo;
     501            0 :         const char* tag = reinterpret_cast<const char*>(opInfo.opId.tag);
     502            0 :         u32 index = opInfo.opId.index;
     503            0 :         const KfcStatus& aicpuState = opInfo.execStatus.kfcStatus;
     504            0 :         if (aicpuState == KfcStatus::kEnd) { // 该rank未下发算子,或算子已执行结束
     505            0 :             HCCL_ERROR(
     506              :                 "[OpRetry][Server]OpRetryServerWaitResp dst[%u] with IpInfo[%s], hccl op not launch or "
     507              :                 "is complete, hccl aicpu can not retry",
     508              :                 it->first, it->second.retryInfo.dfxIpInfo);
     509            0 :             agentRetryInfo.retryInfo.isNeedReportOpRetryErr = true;
     510              :         }
     511            0 :         HCCL_RUN_INFO(
     512              :             "[OpRetry][Server]Print rank[%u], tag[%s], index[%u], aicpuStatus[%d]", it->first, tag, index, aicpuState);
     513              :     }
     514            0 : }
     515              : 
     516            1 : HcclResult OpRetryServerCheckOp::ProcessEvent(RetryContext* retryCtx)
     517              : {
     518            1 :     HcclResult ret = CheckRetryInfo(*retryCtx);
     519            1 :     RetryState nextState = (ret == HCCL_SUCCESS) ? RETRY_STATE_CMD_CHECK_LINK : RETRY_STATE_SERVER_RETRY_FAIL;
     520              : 
     521            1 :     if (ret == HCCL_E_OPRETRY_FAIL) {
     522            1 :         HCCL_RUN_INFO("[OpRetry][Server][CheckRetryInfo] Opname is Inconsistent, RETRY_CONSTRAINT, ret[%u]", ret);
     523            1 :         retryCtx->isNeedReportOpRetryErr = true;
     524              :     }
     525              : 
     526            1 :     HCCL_RUN_INFO("[OpRetry][Server]check op ret[%d], nextState[%s]", ret, GetReadableState(nextState));
     527            1 :     CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
     528            1 :     return HCCL_SUCCESS;
     529              : }
     530              : 
     531            0 : HcclResult OpRetryServerWaitLinkInfo::ProcessEvent(RetryContext* retryCtx)
     532              : {
     533            0 :     std::chrono::steady_clock::time_point startTime = std::chrono::steady_clock::now();
     534            0 :     const u32 timeoutValue = std::max(static_cast<u32>(GetExternalInputHcclLinkTimeOut()), OP_RETRY_SEND_RECV_TIMEOUT)
     535            0 :                              + OP_RETRY_WAIT_AICPU_TIMEOUT;
     536            0 :     const std::chrono::seconds timeout = std::chrono::seconds(timeoutValue);
     537              :     // 下一个server状态
     538            0 :     RetryState nextState = RETRY_STATE_CHECK_ALL_LINK;
     539              : 
     540            0 :     std::set<u32> recvVaild;
     541            0 :     while (recvVaild.size() < retryCtx->needRetryServerRanks_.size()) {
     542            0 :         CHK_PRT_RET(
     543              :             retryCtx->isServerStateWaitResume_,
     544              :             HCCL_RUN_INFO("[OpRetry][Server]switched state form wait link to wait resume"), HCCL_SUCCESS);
     545            0 :         std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
     546            0 :         const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - startTime);
     547            0 :         CHK_PRT_RET(
     548              :             elapsed > timeout, HCCL_ERROR("[OpRetry][Server]OpRetryServerWaitLinkInfo timeout"), HCCL_E_TIMEOUT);
     549              : 
     550            0 :         for (auto rank : retryCtx->needRetryServerRanks_) {
     551            0 :             if (recvVaild.find(rank) != recvVaild.end()) {
     552            0 :                 continue;
     553              :             }
     554            0 :             auto& agentRetryInfo = retryCtx->serverSockets_[rank];
     555              :             // 接收agent信息
     556            0 :             HcclResult ret = WaitLinkPortCheckResult(agentRetryInfo.socket, agentRetryInfo.linkPortStatus);
     557            0 :             CHK_PRT_RET(
     558              :                 ret != HCCL_SUCCESS && ret != HCCL_E_AGAIN,
     559              :                 HCCL_ERROR("[OpRetry][Server]OpRetryServerWaitLinkCheckResult fail, ret[%u]", ret), ret);
     560            0 :             if (ret == HCCL_SUCCESS) {
     561            0 :                 recvVaild.insert(rank);
     562            0 :                 HCCL_INFO("[OpRetry][Server]OpRetryServerWaitLinkCheckResult recv success from dst[%u], ", rank);
     563              :             }
     564              :         }
     565              :     }
     566            0 :     CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
     567            0 :     HCCL_INFO("[OpRetry][Server]OpRetryServerWaitLinkInfo success");
     568            0 :     return HCCL_SUCCESS;
     569            0 : }
     570              : 
     571            0 : HcclResult OpRetryServerCheckAllLink::ProcessEvent(RetryContext* retryCtx)
     572              : {
     573              :     // 收集所有rank的主备网口信息
     574            0 :     std::map<u32, std::pair<bool, bool>> allLinkInfo;
     575            0 :     for (auto rank : retryCtx->needRetryServerRanks_) {
     576            0 :         auto& linkPortStatus = retryCtx->serverSockets_[rank].linkPortStatus;
     577            0 :         allLinkInfo.insert({rank, std::make_pair(linkPortStatus.defaultPort, linkPortStatus.backupPort)});
     578              :     }
     579              : 
     580              :     // 对所有rank依次遍历
     581            0 :     for (auto rank : retryCtx->needRetryServerRanks_) {
     582            0 :         u32 remoteRankIndex = 0;
     583            0 :         auto& linkPortStatus = retryCtx->serverSockets_[rank].linkPortStatus;
     584              :         // 对rank的所有对端进行遍历
     585            0 :         for (u32 i = 0; i < linkPortStatus.rankSize; i++) {
     586            0 :             u32 remoteRank = linkPortStatus.rankList[i];
     587            0 :             retryCtx->serverSockets_[rank].changeLinkInfo.remoteRankList[remoteRankIndex] = remoteRank;
     588            0 :             if (allLinkInfo[rank].first && allLinkInfo[remoteRank].first) {
     589              :                 // 本端和对端的主网口均up,则使用主网口
     590            0 :                 retryCtx->serverSockets_[rank].changeLinkInfo.isUseDefaultPort[remoteRankIndex] = true;
     591            0 :             } else if (allLinkInfo[rank].second && allLinkInfo[remoteRank].second) {
     592              :                 // 本端和对端的备网口均up,则使用备网口
     593            0 :                 retryCtx->serverSockets_[rank].changeLinkInfo.isUseDefaultPort[remoteRankIndex] = false;
     594              :             } else {
     595              :                 // 本端和对端无可用的网口,重执行失败
     596            0 :                 HCCL_ERROR(
     597              :                     "[OpRetry][Server]rank[%u]:default[%d], backup[%d], IpInfo[%s]; rank[%u]:default[%d], "
     598              :                     "backup[%d], can not find same port, can not retry",
     599              :                     rank, allLinkInfo[rank].first, allLinkInfo[rank].second,
     600              :                     retryCtx->serverSockets_[rank].retryInfo.dfxIpInfo, remoteRank, allLinkInfo[remoteRank].first,
     601              :                     allLinkInfo[remoteRank].second);
     602            0 :                 CHK_RET(CreateOpRetryServerByState(RETRY_STATE_SERVER_RETRY_FAIL, retryCtx));
     603            0 :                 return HCCL_SUCCESS;
     604              :             }
     605            0 :             remoteRankIndex += 1;
     606              :         }
     607            0 :         retryCtx->serverSockets_[rank].changeLinkInfo.remoteRankNum = remoteRankIndex;
     608              :     }
     609              : 
     610              :     // 打印所有rank的借轨信息
     611            0 :     for (auto rank : retryCtx->needRetryServerRanks_) {
     612            0 :         auto& changeLinkInfo = retryCtx->serverSockets_[rank].changeLinkInfo;
     613            0 :         std::string changeLinkInfoStr = "rank[" + std::to_string(rank) + "]";
     614            0 :         for (u32 i = 0; i < changeLinkInfo.remoteRankNum; i++) {
     615              :             changeLinkInfoStr
     616            0 :                 += (std::to_string(changeLinkInfo.remoteRankList[i]) + ":"
     617            0 :                     + std::to_string(changeLinkInfo.isUseDefaultPort[i]) + "; ");
     618              :         }
     619            0 :         HCCL_INFO("[OpRetry][Server]changeLinkInfoStr:%s", changeLinkInfoStr.c_str());
     620            0 :     }
     621              : 
     622              :     // 所有rank网口确认成功,切换到给agent发借轨命令状态
     623            0 :     CHK_RET(CreateOpRetryServerByState(RETRY_STATE_CMD_RESUME_TRANSPORT, retryCtx));
     624            0 :     return HCCL_SUCCESS;
     625            0 : }
     626              : 
     627            0 : HcclResult OpRetryServerIssueChangeLinkAndResume::ProcessEvent(RetryContext* retryCtx)
     628              : {
     629            0 :     HcclResult ret = HCCL_SUCCESS;
     630            0 :     RetryState curState = retryCtx->GetRetryState();
     631              :     // 先将每个rank的changeLinkInfo发送至对应agent
     632            0 :     for (auto rank : retryCtx->needRetryServerRanks_) {
     633            0 :         ret = IssueChangeLink(retryCtx->serverSockets_[rank].socket, retryCtx->serverSockets_[rank].changeLinkInfo);
     634            0 :         CHK_PRT_RET(
     635              :             ret != HCCL_SUCCESS,
     636              :             HCCL_ERROR("[OpRetry][Server]OpRetryServerIssueChangeLink fail, curState[%s]", GetReadableState(curState)),
     637              :             ret);
     638            0 :         HCCL_INFO("[OpRetry][Server]OpRetryServerIssueChangeLink send to rank[%u] success", rank);
     639              :     }
     640              :     // 再发送resume transport命令至每个agent
     641            0 :     std::shared_ptr<OpRetryBase> retryPtr = nullptr;
     642            0 :     EXCEPTION_CATCH((retryPtr = std::make_shared<OpRetryServerIssueCmd>()), return HCCL_E_PTR);
     643            0 :     RetryState nextState = RETRY_STATE_CMD_RESUME_TRANSPORT;
     644            0 :     retryCtx->SetRetryState(nextState, retryPtr);
     645            0 :     HCCL_INFO("[OpRetry][Server]OpRetryServerIssueChangeLinkAndResume success");
     646            0 :     return HCCL_SUCCESS;
     647            0 : }
     648              : 
     649            0 : HcclResult OpRetryServerRetryFail::ProcessEvent(RetryContext* retryCtx)
     650              : {
     651            0 :     RetryCommandInfo commandInfo;
     652            0 :     commandInfo.command = RETRY_CMD_RETRY_FAIL;
     653            0 :     if (retryCtx->isNeedReportOpRetryErr) {
     654            0 :         commandInfo.command = RETRY_CMD_RETRY_CONSTRAINT_FAIL;
     655            0 :         HCCL_RUN_INFO(
     656              :             "[OpRetry][Server]OpRetryServerRetryFail isNeedReportOpRetryErr[%d], command[%s]",
     657              :             retryCtx->isNeedReportOpRetryErr, GetReadableCmd(commandInfo.command));
     658              :     }
     659            0 :     HCCL_INFO("[OpRetry][Server]could not retry, command all rank %s", GetReadableCmd(commandInfo.command));
     660            0 :     for (auto rank : retryCtx->needRetryServerRanks_) {
     661            0 :         HcclResult ret = IssueCommandWithOpId(retryCtx->serverSockets_[rank].socket, commandInfo);
     662            0 :         CHK_PRT_RET(
     663              :             ret != HCCL_SUCCESS,
     664              :             HCCL_ERROR("[OpRetry][Server]OpRetryServerRetryFail IssueCommandWithOpId AgentId[%u] fail", rank), ret);
     665              :     }
     666              : 
     667              :     // 重执行异常,通知心跳存在异常 -> 广播异常给整个集群
     668            0 :     Heartbeat::GetInstance(retryCtx->deviceLogicId_).SetOpretryErr();
     669              : 
     670            0 :     RetryState nextState = RETRY_STATE_SERVER_RUNNING;
     671            0 :     CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
     672            0 :     return HCCL_SUCCESS;
     673              : }
     674              : 
     675           15 : bool SwitchNicServerCheckAllSwitchRanks::CompareSwitchRankList(
     676              :     const u32* firstSwitchRankList, const u32* switchRankList, const u32 switchRankNum) const
     677              : {
     678           15 :     if (switchRankNum == 0 || switchRankNum > AICPU_MAX_RANK_NUM) {
     679            0 :         return false;
     680              :     }
     681              : 
     682           15 :     std::set<u32> switchRankSet;
     683           45 :     for (u32 i = 0; i < switchRankNum; i++) {
     684           30 :         switchRankSet.insert(firstSwitchRankList[i]);
     685              :     }
     686              : 
     687           15 :     u32 ranksNum[AICPU_MAX_RANK_NUM] = {0};
     688           43 :     for (u32 i = 0; i < switchRankNum; i++) {
     689           30 :         if (switchRankSet.find(switchRankList[i]) == switchRankSet.end()) {
     690            1 :             HCCL_ERROR("[SwitchNic][Server] rankList has error, id[%u]", switchRankList[i]);
     691            1 :             return false;
     692              :         } else {
     693           29 :             ranksNum[switchRankList[i]]++;
     694              :         }
     695           29 :         if (ranksNum[switchRankList[i]] > 1) {
     696            1 :             HCCL_ERROR(
     697              :                 "[SwitchNic][Server] rankList has error, id[%u], num[%u]", switchRankList[i],
     698              :                 ranksNum[switchRankList[i]]);
     699            1 :             return false;
     700              :         }
     701              :     }
     702           13 :     return true;
     703           15 : }
     704              : 
     705           13 : bool SwitchNicServerCheckAllSwitchRanks::CompareUseBackupLists(
     706              :     const bool* firstArray, const bool* secondArray, const u32 switchRankNum) const
     707              : {
     708           13 :     if (switchRankNum == 0 || switchRankNum > AICPU_MAX_RANK_NUM) {
     709            0 :         return false;
     710              :     }
     711           37 :     for (u32 i = 0; i < switchRankNum; i++) {
     712           25 :         if (firstArray[i] != secondArray[i]) {
     713            1 :             HCCL_ERROR(
     714              :                 "[SwitchNic][Server] backupLists has error first[%u], second[%u], index[%u]", firstArray[i],
     715              :                 secondArray[i], i);
     716            1 :             return false;
     717              :         }
     718              :     }
     719           12 :     return true;
     720              : }
     721              : 
     722           12 : bool SwitchNicServerCheckAllSwitchRanks::CheckRemotePorts(const u32 rankId, const ActiveSwitchInfo& switchRankInfo)
     723              : {
     724           28 :     for (u32 i = 0; i < switchRankInfo.remoteRankNum; i++) {
     725           18 :         if (switchRankInfo.remoteRankNicStatus[i] == CONNECT_REMOTE_DEFAULT && !switchRankInfo.defaultPortStatus) {
     726            1 :             HCCL_ERROR(
     727              :                 "[SwitchNic][Server] defaultPortStatus has error, localRank[%u], remoteRank[%u], nicStatus[%u]", rankId,
     728              :                 i, switchRankInfo.remoteRankNicStatus[i]);
     729            1 :             return false;
     730              :         }
     731           17 :         if (switchRankInfo.remoteRankNicStatus[i] == CONNECT_REMOTE_BACKUP && !switchRankInfo.backupPortStatus) {
     732            1 :             HCCL_ERROR(
     733              :                 "[SwitchNic][Server] backupPortStatus has error, localRank[%u], remoteRank[%u], nicStatus[%u]", rankId,
     734              :                 i, switchRankInfo.remoteRankNicStatus[i]);
     735            1 :             return false;
     736              :         }
     737              :     }
     738           10 :     return true;
     739              : }
     740              : 
     741            3 : HcclResult SwitchNicServerCheckAllSwitchRanks::CollectSingleAgentActiveSwitchInfo(
     742              :     RetryContext* retryCtx, const u32 rankId, HcclAgentRetryInfo& agentInfo)
     743              : {
     744            3 :     std::chrono::steady_clock::time_point startTime = std::chrono::steady_clock::now();
     745            3 :     const auto timeout = std::chrono::seconds(GetExternalInputHcclLinkTimeOut() * ACTIVE_SWITCH_TIMES);
     746            3 :     HcclResult ret = HCCL_SUCCESS;
     747              :     while (true) {
     748            6 :         CHK_PRT_RET(
     749              :             retryCtx->isServerStateWaitResume_,
     750              :             HCCL_RUN_INFO("[OpRetry][Server]switched state form check switch Nic to wait resume"), HCCL_SUCCESS);
     751            3 :         std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
     752            3 :         const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - startTime);
     753            3 :         CHK_PRT_RET(
     754              :             elapsed > timeout,
     755              :             HCCL_ERROR("[SwitchNic][Server] timeout in recv agent RetryInfo, waitime[%u s>%u s]", elapsed, timeout),
     756              :             HCCL_E_TIMEOUT);
     757            3 :         ret = WaitResponse(agentInfo.socket, agentInfo.retryInfo);
     758            3 :         if (ret == HCCL_SUCCESS) { // 成功接收到数据
     759              :             // 成功接收到retryInfo含RETRY_STATE_SEND_SWITCH_INFO,接收 ActiveSwitchInfo,否则为保活数据,忽略
     760            3 :             if (agentInfo.retryInfo.retryState == RETRY_STATE_SEND_SWITCH_INFO) {
     761            3 :                 ret = RecvActiveSwitchInfo(agentInfo.socket, rankId, agentInfo.switchInfo);
     762            3 :                 if (ret != HCCL_SUCCESS) {
     763            0 :                     return ret;
     764              :                 }
     765            3 :                 HCCL_INFO("[SwitchNic][server] recv ActiveSwitchInfo form rank[%u] while collecting", rankId);
     766            3 :                 retryCtx->switchInfoMap_[rankId] = agentInfo.switchInfo;
     767            3 :                 return HCCL_SUCCESS;
     768              :             }
     769            0 :         } else if (ret == HCCL_E_AGAIN) {
     770            0 :             RetryCommand command = RETRY_CMD_RUNNING;
     771            0 :             CHK_RET(IssueCommand(agentInfo.socket, command));
     772            0 :             HCCL_DEBUG("[SwitchNic][Server] send keeping active info to dst[%u]", rankId);
     773              :         } else {
     774            0 :             HCCL_ERROR("[SwitchNic][Server] get active switch info failed, ret[%u], dst[%u]", ret, rankId);
     775            0 :             break;
     776              :         }
     777            0 :         SaluSleep(OP_RETRY_POLL_AICPU_STATE_INTERVAL);
     778            0 :     }
     779            0 :     return ret;
     780              : }
     781              : 
     782            9 : HcclResult SwitchNicServerCheckAllSwitchRanks::CollectAgentActiveSwitchInfo(RetryContext* retryCtx)
     783              : {
     784            9 :     HCCL_RUN_INFO("[SwitchNic][Server] began to CollectAgentActiveSwitchInfo");
     785            9 :     HcclResult ret = HCCL_SUCCESS;
     786              :     // 轮询接收agent信息
     787           27 :     for (auto& it : retryCtx->serverSockets_) {
     788           18 :         const u32& rank = it.first;
     789           18 :         if (retryCtx->switchInfoMap_.find(rank) != retryCtx->switchInfoMap_.end()) {
     790           15 :             HCCL_DEBUG("[SwitchNic][Server] rank[%u] has been received", rank);
     791           15 :             continue;
     792              :         }
     793              :         // 轮询接收agent状态机信息
     794            3 :         ret = CollectSingleAgentActiveSwitchInfo(retryCtx, it.first, it.second);
     795            3 :         if (ret != HCCL_SUCCESS) {
     796            0 :             return ret;
     797              :         }
     798              :     }
     799              :     // 可不检查,理论上不会不等于
     800            9 :     if (retryCtx->switchInfoMap_.size() != retryCtx->serverSockets_.size()) {
     801            0 :         return HCCL_E_UNAVAIL;
     802              :     }
     803            9 :     return ret;
     804              : }
     805              : 
     806            9 : HcclResult SwitchNicServerCheckAllSwitchRanks::CheckAgentActiveSwitchInfo(RetryContext* retryCtx)
     807              : {
     808            9 :     HCCL_RUN_INFO("[SwitchNic][Server] began to CheckAgentActiveSwitchInfo");
     809            9 :     if (retryCtx->switchInfoMap_.empty()) {
     810            0 :         HCCL_ERROR("[SwitchNic][Server] switchInfoMap is empty");
     811            0 :         return HCCL_E_PARA;
     812              :     }
     813              : 
     814            9 :     RetryCommand command = RETRY_CMD_NOTIFY_SWITCH_SUC;
     815              : 
     816            9 :     auto firstInfo = retryCtx->switchInfoMap_.begin();
     817            9 :     auto firstRankId = firstInfo->first;
     818            9 :     ActiveSwitchInfo& firstSwitchInfo = firstInfo->second;
     819              : 
     820           19 :     for (const auto& it : retryCtx->switchInfoMap_) {
     821           17 :         const u32& rankId = it.first;
     822           17 :         const ActiveSwitchInfo& switchInfo = it.second;
     823           17 :         if (!switchInfo.refreshTransportFin) {
     824            1 :             HCCL_ERROR("[SwitchNic][Server] refreshTransportFin is false, first[%u], rank[%u]", firstRankId, rankId);
     825            1 :             command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
     826            1 :             break;
     827              :         }
     828           16 :         if (switchInfo.switchRankNum != firstSwitchInfo.switchRankNum) {
     829            1 :             HCCL_ERROR(
     830              :                 "[SwitchNic][Server] switchRankNum is not same as the first[%u:%u], rank[%u:%u]", firstRankId,
     831              :                 firstSwitchInfo.switchRankNum, rankId, switchInfo.switchRankNum);
     832            1 :             command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
     833            1 :             break;
     834              :         } else {
     835           15 :             if (!CompareSwitchRankList(
     836           15 :                     firstSwitchInfo.switchRankList, switchInfo.switchRankList, switchInfo.switchRankNum)) {
     837            2 :                 HCCL_ERROR(
     838              :                     "[SwitchNic][Server] SwitchRankList is not same as the first[%u], rank[%u], rankNum[%u]",
     839              :                     firstRankId, rankId, switchInfo.switchRankNum);
     840            2 :                 command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
     841            2 :                 break;
     842              :             }
     843           13 :             if (!CompareUseBackupLists(
     844           13 :                     firstSwitchInfo.switchUseBackup, switchInfo.switchUseBackup, switchInfo.switchRankNum)) {
     845            1 :                 HCCL_ERROR(
     846              :                     "[SwitchNic][Server] UseBackupLists is not same as the first[%u], rank[%u], rankNum[%u]",
     847              :                     firstRankId, rankId, switchInfo.switchRankNum);
     848            1 :                 command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
     849            1 :                 break;
     850              :             }
     851              :         }
     852              : 
     853           12 :         if (!switchInfo.localPortsCheckRet) {
     854            0 :             HCCL_ERROR("[SwitchNic][Server] localPortsCheckRet is false, first[%u], rank[%u]", firstRankId, rankId);
     855            0 :             command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
     856            0 :             break;
     857              :         }
     858           12 :         if (!CheckRemotePorts(rankId, switchInfo)) {
     859            2 :             command = RETRY_CMD_NOTIFY_SWITCH_FAIL;
     860            2 :             break;
     861              :         }
     862              :     }
     863              : 
     864              :     // 全部卡确认无误后或者发现错误后,通知OpRetryAgent。
     865           27 :     for (auto it : retryCtx->serverSockets_) {
     866           18 :         HcclResult ret = IssueCommand(it.second.socket, command);
     867           18 :         CHK_PRT_RET(
     868              :             ret != HCCL_SUCCESS,
     869              :             HCCL_ERROR("[SwitchNic][Server] CheckAllSwitchRanks IssueCommand AgentId[%u] fail", it.first), ret);
     870           18 :     }
     871            9 :     retryCtx->switchInfoMap_.clear();
     872            9 :     return HCCL_SUCCESS;
     873              : }
     874              : 
     875              : // OpRetryServer遍历通信域内的所有卡,接收主动借轨信息,并且校验每张卡信息一致
     876              : // 全部卡确认无误后或者发现错误后,通知OpRetryAgent
     877            9 : HcclResult SwitchNicServerCheckAllSwitchRanks::ProcessEvent(RetryContext* retryCtx)
     878              : {
     879            9 :     HCCL_RUN_INFO("[SwitchNic][Server] CheckAllSwitchRanks begin");
     880            9 :     RetryState nextState = RETRY_STATE_SERVER_RUNNING;
     881              : 
     882            9 :     CHK_RET(CollectAgentActiveSwitchInfo(retryCtx));
     883            9 :     CHK_RET(CheckAgentActiveSwitchInfo(retryCtx));
     884            9 :     CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
     885            9 :     HCCL_RUN_INFO("[SwitchNic][Server] CheckAllSwitchRanks end");
     886            9 :     return HCCL_SUCCESS;
     887              : }
     888              : 
     889            0 : HcclResult OpRetryServerWaitResume::ProcessEvent(RetryContext* retryCtx)
     890              : {
     891            0 :     if (!retryCtx->isServerStateWaitResume_ && !retryCtx->haveCommEnableBackupLink_) {
     892            0 :         CHK_RET(CreateOpRetryServerByState(RETRY_STATE_SERVER_RUNNING, retryCtx));
     893            0 :         HCCL_RUN_INFO(
     894              :             "[OpRetry][Server]OpRetryServerWaitResume, group[%s], no comm enable backup link, set state to running",
     895              :             retryCtx->group_.c_str());
     896            0 :         return HCCL_SUCCESS;
     897              :     }
     898            0 :     if (!retryCtx->isServerStateWaitResume_ && retryCtx->haveCommEnableBackupLink_) {
     899            0 :         HCCL_RUN_INFO("[OpRetry][Server]OpRetryServerWaitResume, start to send cmd");
     900            0 :         for (auto& it : retryCtx->serverSockets_) {
     901            0 :             const u32& agentId = it.first;
     902            0 :             RetryCommandInfo commandInfo;
     903            0 :             commandInfo.command = RESUME_CMD_CHECK_LINK;
     904            0 :             CHK_PRT_RET(
     905              :                 IssueCommandWithOpId(it.second.socket, commandInfo),
     906              :                 HCCL_ERROR("[OpRetry][Server][Resume]rank[%u] send resume check link fail", agentId), HCCL_E_INTERNAL);
     907            0 :             HCCL_RUN_INFO(
     908              :                 "[OpRetry][Server][Resume]rank[%u] send RESUME_CMD_CHECK_LINK, group[%s]", agentId,
     909              :                 retryCtx->group_.c_str());
     910              :         }
     911            0 :         CHK_RET(CreateOpRetryServerByState(RETRY_RESUME_STATE_SERVER_CHECK_LINK, retryCtx));
     912            0 :         HCCL_RUN_INFO("[OpRetry][Server]OpRetryServerWaitResume, set state to check link");
     913            0 :         retryCtx->isRdmaError = false;
     914            0 :         return HCCL_SUCCESS;
     915              :     }
     916              : 
     917            0 :     return HCCL_SUCCESS;
     918              : }
     919              : 
     920            0 : HcclResult ResumeServerCheckAllLink::ProcessEvent(RetryContext* retryCtx)
     921              : {
     922            0 :     HCCL_RUN_INFO("[OpRetry][Server]ResumeServerCheckAllLink begin group[%s]", retryCtx->group_.c_str());
     923            0 :     RetryState nextState = RETRY_RESUME_STATE_SERVER_CHANGE_LINK;
     924              :     // 收集所有agent的检查网口结果
     925            0 :     CHK_RET(WaitAgentCheckLinkResult(retryCtx));
     926              :     // 检查所有rank的主备链路连接情况
     927            0 :     CHK_RET(CheckAllLink(retryCtx, nextState));
     928              : 
     929            0 :     CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
     930            0 :     return HCCL_SUCCESS;
     931              : }
     932              : 
     933            0 : HcclResult ResumeServerChangeLink::ProcessEvent(RetryContext* retryCtx)
     934              : {
     935            0 :     HCCL_RUN_INFO("[OpRetry][Server]ResumeServerChangeLink begin");
     936            0 :     RetryState nextState = RETRY_STATE_SERVER_RUNNING;
     937              :     // 下发切换链路命令字到所有rank,并发送重建transport命令
     938            0 :     CHK_RET(CmdAgentChangeLink(retryCtx));
     939              : 
     940              :     // 接收所有rank的切换链路结果
     941            0 :     CHK_RET(WaitAllChangeLinkResult(retryCtx, nextState));
     942            0 :     if (nextState != RETRY_STATE_SERVER_RUNNING) {
     943            0 :         HCCL_ERROR("[OpRetry][Server]ResumeServerChangeLink fail, nextState[%s]", GetReadableState(nextState));
     944              :     }
     945            0 :     CHK_RET(CreateOpRetryServerByState(nextState, retryCtx));
     946            0 :     return HCCL_SUCCESS;
     947              : }
     948              : 
     949            1 : HcclResult ResumeServerCheckAllLink::WaitAgentCheckLinkResult(RetryContext* retryCtx)
     950              : {
     951            1 :     HCCL_RUN_INFO("[OpRetry][Server][Resume]WaitAgentCheckLinkResult begin");
     952            1 :     std::chrono::steady_clock::time_point startTime = std::chrono::steady_clock::now();
     953            1 :     const u32 timeoutValue = std::max(static_cast<u32>(GetExternalInputHcclLinkTimeOut()), OP_RETRY_SEND_RECV_TIMEOUT)
     954            1 :                              + OP_RETRY_WAIT_AICPU_TIMEOUT;
     955            1 :     const std::chrono::seconds timeout = std::chrono::seconds(timeoutValue);
     956            1 :     std::set<u32> recvVaild;
     957            2 :     while (recvVaild.size() < retryCtx->serverSockets_.size()) {
     958            1 :         std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
     959            1 :         const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - startTime);
     960            1 :         CHK_PRT_RET(
     961              :             elapsed > timeout, HCCL_ERROR("[OpRetry][Server][Resume]WaitAgentCheckLinkResult timeout"), HCCL_E_TIMEOUT);
     962              : 
     963            2 :         for (auto& rank : retryCtx->serverSockets_) {
     964            1 :             const u32& agentId = rank.first;
     965            1 :             if (recvVaild.find(agentId) != recvVaild.end()) {
     966            0 :                 continue;
     967              :             }
     968            1 :             auto& agentRetryInfo = rank.second;
     969            1 :             HcclResult ret = WaitLinkPortCheckResult(agentRetryInfo.socket, agentRetryInfo.linkPortStatus);
     970            1 :             CHK_PRT_RET(
     971              :                 ret != HCCL_SUCCESS && ret != HCCL_E_AGAIN,
     972              :                 HCCL_ERROR("[OpRetry][Server][Resume]WaitAgentCheckLinkResult WaitLinkPortCheckResult Failed"), ret);
     973            1 :             if (ret == HCCL_SUCCESS) {
     974            1 :                 recvVaild.insert(agentId);
     975            1 :                 HCCL_RUN_INFO(
     976              :                     "[OpRetry][Server][Resume]WaitAgentCheckLinkResult recv valid from agentId[%u], "
     977              :                     "group[%s],remoteRank[%u]",
     978              :                     agentId, retryCtx->group_.c_str(), rank.second.linkPortStatus.rankList[0]);
     979              :             }
     980              :         }
     981              :     }
     982            1 :     HCCL_INFO("[OpRetry][Server][Resume]WaitAgentCheckLinkResult recv all valid");
     983            1 :     return HCCL_SUCCESS;
     984            1 : }
     985              : 
     986            0 : HcclResult ResumeServerCheckAllLink::CheckAllLink(RetryContext* retryCtx, RetryState& nextState)
     987              : {
     988            0 :     std::map<u32, std::pair<bool, bool>> allLinkInfo;
     989            0 :     for (auto it : retryCtx->serverSockets_) {
     990            0 :         u32 rank = it.first;
     991            0 :         auto& linkPortStatus = retryCtx->serverSockets_[rank].linkPortStatus;
     992            0 :         HCCL_RUN_INFO(
     993              :             "[OpRetry][Server][Resume]CheckAllLink rank[%u], rankListSize[%d], rankSize[%d]", rank,
     994              :             std::end(it.second.linkPortStatus.rankList) - std::begin(it.second.linkPortStatus.rankList),
     995              :             it.second.linkPortStatus.rankSize);
     996            0 :         allLinkInfo.insert({it.first, std::make_pair(linkPortStatus.defaultPort, linkPortStatus.backupPort)});
     997            0 :     }
     998              : 
     999            0 :     for (auto it : retryCtx->serverSockets_) {
    1000            0 :         u32 rank = it.first;
    1001            0 :         u32 remoteRankIndex = 0;
    1002            0 :         auto& linkPortStatus = it.second.linkPortStatus;
    1003            0 :         for (u32 i = 0; i < linkPortStatus.rankSize; i++) {
    1004            0 :             u32 remoteRank = linkPortStatus.rankList[i];
    1005            0 :             retryCtx->serverSockets_[rank].changeLinkInfo.remoteRankList[remoteRankIndex] = remoteRank;
    1006            0 :             if (allLinkInfo[rank].first && allLinkInfo[remoteRank].first) {
    1007              :                 // 本端与对端的主网口均up, 则使用主网口
    1008            0 :                 retryCtx->serverSockets_[rank].changeLinkInfo.isUseDefaultPort[remoteRankIndex] = true;
    1009            0 :             } else if (allLinkInfo[rank].second && allLinkInfo[remoteRank].second) {
    1010            0 :                 retryCtx->serverSockets_[rank].changeLinkInfo.isUseDefaultPort[remoteRankIndex] = false;
    1011              :             } else {
    1012            0 :                 HCCL_ERROR(
    1013              :                     "[OpRetry][Server][Resume]rank[%u]:default[%d], backup[%d], IpInfo[%s]; "
    1014              :                     "remoterank[%u]:default[%d], "
    1015              :                     "backup[%d], can not find same port, can not resume",
    1016              :                     rank, allLinkInfo[rank].first, allLinkInfo[rank].second,
    1017              :                     retryCtx->serverSockets_[rank].retryInfo.dfxIpInfo, remoteRank, allLinkInfo[remoteRank].first,
    1018              :                     allLinkInfo[remoteRank].second);
    1019            0 :                 nextState = RETRY_STATE_SERVER_RETRY_FAIL;
    1020            0 :                 return HCCL_SUCCESS;
    1021              :             }
    1022            0 :             HCCL_RUN_INFO(
    1023              :                 "[OpRetry][Server][Resume]CheckAllLink remoteRank[%u], changeLinkInfo remoteRankList[%u], "
    1024              :                 "linkPortStatus rankList[%u]",
    1025              :                 remoteRank, retryCtx->serverSockets_[rank].changeLinkInfo.remoteRankList[0],
    1026              :                 linkPortStatus.rankList[0]);
    1027            0 :             remoteRankIndex++;
    1028              :         }
    1029            0 :         retryCtx->serverSockets_[rank].changeLinkInfo.remoteRankNum = remoteRankIndex;
    1030            0 :     }
    1031              : 
    1032              :     // 打印所有rank的借轨信息
    1033            0 :     for (auto it : retryCtx->serverSockets_) {
    1034            0 :         u32 rank = it.first;
    1035            0 :         auto& changeLinkInfo = it.second.changeLinkInfo;
    1036            0 :         std::string changeLinkInfoStr = "rank[" + std::to_string(rank) + "]";
    1037            0 :         for (u32 i = 0; i < changeLinkInfo.remoteRankNum; i++) {
    1038              :             changeLinkInfoStr
    1039            0 :                 += (std::to_string(changeLinkInfo.remoteRankList[i]) + ":"
    1040            0 :                     + std::to_string(changeLinkInfo.isUseDefaultPort[i]) + "; ");
    1041              :         }
    1042            0 :         HCCL_INFO("[OpRetry][Server][Resume]changeLinkInfoStr:%s", changeLinkInfoStr.c_str());
    1043            0 :     }
    1044            0 :     return HCCL_SUCCESS;
    1045            0 : }
    1046              : 
    1047            0 : HcclResult ResumeServerChangeLink::CmdAgentChangeLink(RetryContext* retryCtx)
    1048              : {
    1049            0 :     HCCL_RUN_INFO("[OpRetry][Server][Resume]CmdAgentChangeLink begin");
    1050              :     // 先将每个rank的changeLinkInfo发送至对应agent
    1051            0 :     for (auto it : retryCtx->serverSockets_) {
    1052            0 :         HcclResult ret = IssueChangeLink(it.second.socket, it.second.changeLinkInfo);
    1053            0 :         CHK_PRT_RET(
    1054              :             ret != HCCL_SUCCESS,
    1055              :             HCCL_ERROR("[OpRetry][Server][Resume]CmdAgentChangeLink IssueCommandChangeLink RankId[%u] fail", it.first),
    1056              :             ret);
    1057            0 :         HCCL_INFO(
    1058              :             "[OpRetry][Server]OpRetryServerIssueChangeLink send ChangeLinkInfo to rank[%u] success, group[%s]",
    1059              :             it.first, retryCtx->group_.c_str());
    1060            0 :     }
    1061            0 :     for (auto& it : retryCtx->serverSockets_) {
    1062            0 :         const u32& agentId = it.first;
    1063            0 :         RetryCommandInfo commandInfo;
    1064            0 :         commandInfo.command = RETRY_CMD_RESUME_TRANSPORT;
    1065            0 :         HcclResult ret = IssueCommandWithOpId(it.second.socket, commandInfo);
    1066            0 :         CHK_PRT_RET(
    1067              :             ret != HCCL_SUCCESS,
    1068              :             HCCL_ERROR("[OpRetry][Server][Resume]CmdAgentChangeLink IssueCommand AgentId[%u] fail", agentId), ret);
    1069            0 :         HCCL_INFO(
    1070              :             "[OpRetry][Server]OpRetryServerIssueChangeLink send change link command to rank[%u] success, group[%s]",
    1071              :             it.first, retryCtx->group_.c_str());
    1072              :     }
    1073            0 :     return HCCL_SUCCESS;
    1074              : }
    1075              : 
    1076            0 : HcclResult ResumeServerChangeLink::WaitAllChangeLinkResult(RetryContext* retryCtx, RetryState& nextState)
    1077              : {
    1078            0 :     HCCL_RUN_INFO("[OpRetry][Server][Resume]WaitAllChangeLinkResult begin group[%s]", retryCtx->group_.c_str());
    1079            0 :     std::chrono::steady_clock::time_point startTime = std::chrono::steady_clock::now();
    1080            0 :     const u32 timeoutValue = std::max(static_cast<u32>(GetExternalInputHcclLinkTimeOut()), OP_RETRY_SEND_RECV_TIMEOUT)
    1081            0 :                              + OP_RETRY_WAIT_AICPU_TIMEOUT;
    1082            0 :     const std::chrono::seconds timeout = std::chrono::seconds(timeoutValue);
    1083            0 :     RetryState expectAgentState = RETRY_STATE_AGENT_RUNNING;
    1084            0 :     std::set<u32> recvValid;
    1085            0 :     while (recvValid.size() < retryCtx->serverSockets_.size()) {
    1086            0 :         std::chrono::steady_clock::time_point curTime = std::chrono::steady_clock::now();
    1087            0 :         const auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(curTime - startTime);
    1088            0 :         CHK_PRT_RET(
    1089              :             elapsed > timeout, HCCL_ERROR("[OpRetry][Server][Resume]WaitAllChangeLinkResult timeout"), HCCL_E_TIMEOUT);
    1090              : 
    1091            0 :         for (auto& it : retryCtx->serverSockets_) {
    1092            0 :             u32 rank = it.first;
    1093            0 :             if (recvValid.find(rank) != recvValid.end()) {
    1094            0 :                 continue;
    1095              :             }
    1096            0 :             auto& agentRetryInfo = retryCtx->serverSockets_[rank];
    1097            0 :             HcclResult ret = WaitResponse(agentRetryInfo.socket, agentRetryInfo.retryInfo);
    1098            0 :             RetryState dstState = agentRetryInfo.retryInfo.retryState;
    1099            0 :             KfcStatus aicpuState = agentRetryInfo.retryInfo.opInfo.execStatus.kfcStatus;
    1100            0 :             if (ret == HCCL_SUCCESS && dstState == expectAgentState && aicpuState == KfcStatus::kResumeChanged) {
    1101            0 :                 recvValid.insert(rank);
    1102            0 :                 HCCL_INFO("[OpRetry][Server][Resume]WaitAllChangeLinkResult recv valid from rank[%u]", rank);
    1103            0 :             } else if (ret == HCCL_SUCCESS && dstState == RETRY_STATE_RESP_RUNNING_ERR) {
    1104            0 :                 recvValid.insert(rank);
    1105            0 :                 nextState = RETRY_STATE_SERVER_RETRY_FAIL;
    1106            0 :                 HCCL_ERROR("[OpRetry][Server][Resume]WaitAllChangeLinkResult recv err from rank[%u]", rank);
    1107            0 :                 return HCCL_SUCCESS;
    1108              :             }
    1109              :         }
    1110              :     }
    1111            0 :     return HCCL_SUCCESS;
    1112            0 : }
    1113              : } // namespace hccl
        

Generated by: LCOV version 2.0-1