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

Generated by: LCOV version 2.0-1