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

Generated by: LCOV version 2.0-1