LCOV - code coverage report
Current view: top level - legacy/ascend910/framework/communicator/impl/zero_copy - zero_copy_memory_agent.cc (source / functions) Coverage Total Hit
Test: coverage.info Lines: 1.0 % 628 6
Test Date: 2026-08-04 10:52:23 Functions: 4.0 % 50 2

            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 "zero_copy_memory_agent.h"
      12              : #include <string>
      13              : #include "acl/acl_rt.h"
      14              : #include "hccl_network_pub.h"
      15              : #include "adapter_hccp_common.h"
      16              : #include "adapter_rts_common.h"
      17              : #include "snapshot_control.h"
      18              : 
      19              : namespace hccl {
      20              : using namespace std;
      21              : 
      22              : const string STR_IPC_MEM_EXCHANGE = "IpcMemExchange";
      23              : constexpr u32 IPC_MEMORY_EXCHANGE_LENGTH = 64;  // Bytes
      24              : constexpr u32 USLEEP_ONE_THOUSAND = 1000;
      25              : constexpr int INNER_THREAD_LOOP_US = 500;
      26              : 
      27              : std::unique_ptr<ZeroCopyAddressMgr> ZeroCopyMemoryAgent::addressMgr_ = nullptr;
      28              : 
      29              : template <typename T>
      30            0 : HcclResult ConstructData(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize, T& value)
      31              : {
      32            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &value, sizeof(T)));
      33            0 :     exchangeDataPtr += sizeof(T);
      34            0 :     exchangeDataBlankSize -= sizeof(T);
      35            0 :     return HCCL_SUCCESS;
      36              : }
      37              : 
      38              : /* copy 变长数据 */
      39            0 : HcclResult ConstructData(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize, void *ptr, size_t len)
      40              : {
      41            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, ptr, len));
      42            0 :     exchangeDataPtr += len;
      43            0 :     exchangeDataBlankSize -= len;
      44            0 :     return HCCL_SUCCESS;
      45              : }
      46              : 
      47              : 
      48              : template <typename T>
      49            0 : HcclResult ParseData(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize, T& value)
      50              : {
      51            0 :     CHK_PRT_RET(exchangeDataBlankSize < sizeof(T),
      52              :         HCCL_ERROR("[ParseData] blankSize is [%u] less than [%lu]", exchangeDataBlankSize, sizeof(T)), HCCL_E_INTERNAL);
      53              : 
      54            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(&value, sizeof(T), exchangeDataPtr, sizeof(T)));
      55            0 :     exchangeDataPtr += sizeof(T);
      56            0 :     exchangeDataBlankSize -= sizeof(T);
      57            0 :     return HCCL_SUCCESS;
      58              : }
      59              : 
      60            0 : ZeroCopyMemoryAgent::ZeroCopyMemoryAgent(const std::unique_ptr<HcclSocketManager> &socketManager, u32 devicePhyId,
      61              :     s32 deviceLogicId, const HcclIpAddress &localVnicIp, const std::vector<RankInfo> &rankInfoList, RankId userRank,
      62            0 :     bool useSuperPodMode, const std::string &identifier)
      63            0 :     : initiated_(false), socketManager_(socketManager), devicePhyId_(devicePhyId), deviceLogicId_(deviceLogicId),
      64            0 :       localVnicIp_(localVnicIp), rankInfoList_(rankInfoList), userRank_(userRank), rankSize_(rankInfoList.size()),
      65            0 :       useSuperPodMode_(useSuperPodMode), identifier_(identifier)
      66            0 : {}
      67              : 
      68              : // 创建vnic socket连接,启动recv 接收线程
      69              : // 每个rank 都启动listen,并且都和对端connect
      70            0 : HcclResult ZeroCopyMemoryAgent::Init()
      71              : {
      72            0 :     isSingleRank_ = (rankInfoList_.size() == 1);
      73            0 :     CHK_PRT_RET(isSingleRank_, HCCL_INFO("[ZeroCopyMemoryAgent][Init] single rank communicator"), HCCL_SUCCESS);
      74            0 :     std::unique_lock<std::mutex> lock(commRefCntLock_);
      75              : 
      76            0 :     if (!ZeroCopyMemoryAgent::IsAddressMgrInited()) {
      77            0 :         addressMgr_ = std::make_unique<ZeroCopyAddressMgr>();
      78            0 :         HCCL_RUN_INFO("[ZeroCopyMemoryAgent][%s]init addressMgr_ success.", __func__);
      79              :     }
      80            0 :     CHK_RET(addressMgr_->IncreCommRefCnt());
      81              : 
      82            0 :     CHK_RET(EstablishSockets());
      83              : 
      84            0 :     exchangeDataForSend_.resize(IPC_MEMORY_EXCHANGE_LENGTH * ZERO_COPY_MEMORY_AGENT_SEND_QUEUE_SIZE, 0);
      85            0 :     for (const auto& kv : mapDevPhyIdconnectedSockets_) {
      86            0 :         exchangeDataForAck_[kv.first].resize(IPC_MEMORY_EXCHANGE_LENGTH, 0);
      87            0 :         sendMgrs_[kv.first].reqDataSize_ = IPC_MEMORY_EXCHANGE_LENGTH;
      88            0 :         recvMgrs_[kv.first].receivedData_.resize(ZERO_COPY_MEMORY_AGENT_RECV_QUEUE_SIZE,
      89            0 :             std::vector<u8>(IPC_MEMORY_EXCHANGE_LENGTH, 0));
      90              :     }
      91              : 
      92            0 :     CHK_RET(InitInnerThread());
      93              : 
      94            0 :     return HCCL_SUCCESS;
      95            0 : }
      96              : 
      97            0 : HcclResult ZeroCopyMemoryAgent::InitInnerThread()
      98              : {
      99            0 :     threadRun_ = true;
     100            0 :     innerThread_.reset(new (std::nothrow) std::thread(&ZeroCopyMemoryAgent::InnerThread, std::ref(*this)));
     101            0 :     CHK_SMART_PTR_NULL(innerThread_);
     102            0 :     return HCCL_SUCCESS;
     103              : }
     104              : 
     105            0 : HcclResult ZeroCopyMemoryAgent::EstablishSockets()
     106              : {
     107            0 :     CHK_PRT_RET((vnicPortCtx_ != nullptr),
     108              :         HCCL_ERROR("[ZeroCopyMemoryAgent][Init] already initd"), HCCL_E_PARA);
     109            0 :     CHK_RET(HcclNetOpenDev(&vnicPortCtx_, NicType::VNIC_TYPE, devicePhyId_, deviceLogicId_, localVnicIp_));
     110            0 :     CHK_PTR_NULL(vnicPortCtx_);
     111              : 
     112            0 :     isSocketSupportAsync_ = HcclSocket::IsSupportAsync();
     113            0 :     HCCL_RUN_INFO("[ZeroCopyMemoryAgent][Init] isSocketSupportAsync[%d]", isSocketSupportAsync_);
     114              : 
     115            0 :     for (size_t i = 0; i < rankInfoList_.size(); i++) {
     116            0 :         if (rankInfoList_[i].devicePhyId == static_cast<s32>(devicePhyId_)) {
     117            0 :             continue;
     118              :         }
     119            0 :         HcclRankLinkInfo remoteLinkInfo;
     120            0 :         RankInfo dstRankInfo = rankInfoList_[i];
     121            0 :         remoteLinkInfo.userRank = dstRankInfo.userRank;
     122            0 :         remoteLinkInfo.devicePhyId = dstRankInfo.devicePhyId;
     123            0 :         remoteLinkInfo.ip = HcclIpAddress(dstRankInfo.devicePhyId);
     124            0 :         if (useSuperPodMode_) {
     125            0 :             CHK_RET(hrtRaGetSingleSocketVnicIpInfo(devicePhyId_, DeviceIdType::DEVICE_ID_TYPE_SDID,
     126              :                 dstRankInfo.superDeviceId, remoteLinkInfo.ip));
     127              :         } else {
     128            0 :             CHK_RET(hrtRaGetSingleSocketVnicIpInfo(devicePhyId_, DeviceIdType::DEVICE_ID_TYPE_PHY_ID,
     129              :                 dstRankInfo.devicePhyId, remoteLinkInfo.ip));
     130              :         }
     131              :         // 通信域未分配端口则使用默认端口
     132            0 :         remoteLinkInfo.port =
     133            0 :             dstRankInfo.deviceVnicPort == HCCL_INVALID_PORT ? HETEROG_CCL_PORT : dstRankInfo.deviceVnicPort;
     134            0 :         remoteLinkInfo.socketsPerLink = 1;
     135            0 :         string newTag = GenerateSocketTag(devicePhyId_, rankInfoList_[i].devicePhyId);
     136            0 :         std::vector<std::shared_ptr<HcclSocket> > tmpSockets;
     137            0 :         HcclResult ret = socketManager_->CreateSingleLinkSocket(
     138              :             newTag, vnicPortCtx_, remoteLinkInfo, tmpSockets, false, true);
     139            0 :         CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_ERROR("[Create][DestSockets]Create single link sockets failed, "
     140              :             "local rank[%u], remote rank[%u]", userRank_, i), ret);
     141            0 :         if (tmpSockets.size() != 1) {
     142            0 :             HCCL_ERROR("[ZeroCopyMemoryAgent][CreateVnic] socket number[%llu] is not 1 as expected!", tmpSockets.size());
     143            0 :             return HCCL_E_INTERNAL;
     144              :         }
     145              :         // 设置强制断链为关闭,避免进程退出时recv失败
     146            0 :         tmpSockets[0]->SetForceClose(false);
     147            0 :         mapDevPhyIdconnectedSockets_[remoteLinkInfo.devicePhyId] = (tmpSockets[0]);
     148            0 :         mapDevPhyId2RankId_[remoteLinkInfo.devicePhyId] = remoteLinkInfo.userRank;
     149            0 :     }
     150              : 
     151            0 :     for (const auto& kv : mapDevPhyIdconnectedSockets_) {
     152            0 :         CHK_PRT_RET(socketManager_->WaitLinkEstablish(kv.second) != HCCL_SUCCESS,
     153              :             HCCL_ERROR("[ZeroCopyMemoryAgent][EstablishSockets] tag[%s] socket establish failed", kv.second->GetTag().c_str()),
     154              :             HCCL_E_INTERNAL);
     155              :     }
     156            0 :     return HCCL_SUCCESS;
     157              : }
     158              : 
     159            0 : std::string ZeroCopyMemoryAgent::GenerateSocketTag(u32 localRank, u32 remoteRank)
     160              : {
     161            0 :     u32 small = localRank;
     162            0 :     u32 large = remoteRank;
     163              : 
     164            0 :     if (localRank > remoteRank) {
     165            0 :         small = remoteRank;
     166            0 :         large = localRank;
     167              :     }
     168              : 
     169              :     // Socket构造规则:前缀 + identifier + small + large
     170            0 :     std::string tag = STR_IPC_MEM_EXCHANGE + "_" + identifier_ 
     171            0 :         + "_" + std::to_string(small) + ":" + std::to_string(large);
     172            0 :     return tag;
     173              : }
     174              : 
     175            0 : HcclResult ZeroCopyMemoryAgent::SendRequestSync(RequestType requestType, const std::vector<u8>& req, u32 remoteDevPhyId)
     176              : {
     177              :     HcclResult ret;
     178            0 :     if (remoteDevPhyId != INVALID_VALUE_RANKID) {
     179            0 :         std::unique_lock<std::mutex> lock(sendMutex_);  // send 存在多线调用,需要锁保护
     180            0 :         ret = mapDevPhyIdconnectedSockets_[remoteDevPhyId]->Send(req.data(), IPC_MEMORY_EXCHANGE_LENGTH); 
     181            0 :         CHK_PRT_RET(ret != HCCL_SUCCESS, 
     182              :             HCCL_ERROR("[ZeroCopyMemoryAgent][SendRequestSync] Send %s to remote[%u] failed",
     183              :                         GetReadableRequestType(requestType), remoteDevPhyId),
     184              :             HCCL_E_INTERNAL);
     185            0 :         return HCCL_SUCCESS;
     186            0 :     }
     187              : 
     188            0 :     std::unique_lock<std::mutex> lock(sendMutex_);
     189            0 :     for (const auto& kv : mapDevPhyIdconnectedSockets_) {
     190            0 :         CHK_PRT_RET(kv.second->Send(req.data(), IPC_MEMORY_EXCHANGE_LENGTH) != HCCL_SUCCESS,
     191              :             HCCL_ERROR("[ZeroCopyMemoryAgent][SendRequestSync] Send %s to remote[%u] failed",
     192              :                         GetReadableRequestType(requestType), kv.first),
     193              :             HCCL_E_INTERNAL);
     194              :     }
     195            0 :     return HCCL_SUCCESS;
     196            0 : }
     197              : 
     198            0 : HcclResult ZeroCopyMemoryAgent::SendRequest(RequestType requestType, const std::vector<u8>& req, u32 remoteDevPhyId)
     199              : {
     200            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][SendRequest] requestType[%s] remote[%u]",
     201              :         GetReadableRequestType(requestType), remoteDevPhyId);
     202              : 
     203            0 :     if (!isSocketSupportAsync_) {  // socket不支持异步收发的场景
     204            0 :         return SendRequestSync(requestType, req, remoteDevPhyId);
     205              :     }
     206              : 
     207            0 :     bool isAck = IsAckRequestType(requestType);
     208            0 :     if (remoteDevPhyId != INVALID_VALUE_RANKID) {
     209            0 :         sendMgrs_[remoteDevPhyId].AddRequest(isAck, req);
     210              :     } else {
     211            0 :         for (auto& kv : sendMgrs_) {
     212            0 :             kv.second.AddRequest(isAck, req);
     213              :         }
     214              :     }
     215              : 
     216              :     // 唤醒内部io线程
     217            0 :     std::unique_lock<std::mutex> lock(sendMutex_);
     218            0 :     hasSendRequest_ = true;
     219            0 :     sendCv_.notify_all();
     220            0 :     return HCCL_SUCCESS;
     221            0 : }
     222              : 
     223            0 : void ZeroCopyMemoryAgent::RequestBatchSendAsync()
     224              : {
     225              :     HcclResult ret;
     226            0 :     for (auto &kv : sendMgrs_) {
     227            0 :         auto &sendMgr = kv.second;
     228            0 :         if ((sendMgr.lastSendHandle_ != nullptr) || (!sendMgr.hasReq_[0] && !sendMgr.hasReq_[1])) {
     229              :             // 前回发送未完成 或者 没有待发送的数据
     230            0 :             continue;
     231              :         }
     232              : 
     233            0 :         if (mapDevPhyIdconnectedSockets_.find(kv.first) == mapDevPhyIdconnectedSockets_.end()) {
     234            0 :             HCCL_ERROR("[ZeroCopyMemoryAgent][RequestBatchSendAsync] remote[%u] not found in"
     235              :                 "mapDevPhyIdconnectedSockets_", kv.first);
     236            0 :             continue;
     237              :         }
     238            0 :         auto &socket = mapDevPhyIdconnectedSockets_[kv.first];
     239            0 :         if (sendMgr.sentSize_ == 0) {  // 非断点续传
     240            0 :             if (sendMgr.hasReq_[0] && sendMgr.hasReq_[1]) {  // 合并发送
     241            0 :                 u8 *ptr = const_cast<u8 *>(sendMgr.reqDatas_[1]->data()) + IPC_MEMORY_EXCHANGE_LENGTH;
     242            0 :                 u32 leftSize = IPC_MEMORY_EXCHANGE_LENGTH;
     243            0 :                 if (ConstructData(ptr, leftSize, const_cast<u8 *>(sendMgr.reqDatas_[0]->data()),
     244            0 :                         IPC_MEMORY_EXCHANGE_LENGTH) == HCCL_SUCCESS) {
     245            0 :                     sendMgr.hasReq_[0] = false;
     246            0 :                     sendMgr.currIndex_ = 1;
     247            0 :                     sendMgr.reqDataSize_ = IPC_MEMORY_EXCHANGE_LENGTH + IPC_MEMORY_EXCHANGE_LENGTH;
     248              :                 } else {
     249            0 :                     sendMgr.currIndex_ = 0;
     250            0 :                     sendMgr.reqDataSize_ = IPC_MEMORY_EXCHANGE_LENGTH;
     251              :                 }
     252              :             } else {
     253            0 :                 sendMgr.currIndex_ = sendMgr.hasReq_[0] ? 0 : 1;
     254            0 :                 sendMgr.reqDataSize_ = IPC_MEMORY_EXCHANGE_LENGTH;
     255              :             }
     256              :         }
     257            0 :         const std::vector<u8> *req = sendMgr.reqDatas_[sendMgr.currIndex_];
     258            0 :         sendMgr.lastSendSize_ = 0;  // 用于ra上报发送的数据量
     259            0 :         ret = socket->SendAsync(req->data() + sendMgr.sentSize_,
     260            0 :                                 sendMgr.reqDataSize_ - sendMgr.sentSize_,
     261              :                                 &sendMgr.lastSendSize_, &sendMgr.lastSendHandle_);
     262            0 :         if (ret != HCCL_SUCCESS && ret != HCCL_E_AGAIN) {  // 发送失败的场景
     263            0 :             RequestType requestType = *reinterpret_cast<const RequestType *>(req->data());
     264            0 :             HCCL_ERROR("[ZeroCopyMemoryAgent][RequestBatchSendAsync] failed, ret[%d] remote[%u] requestType[%s] sentSize[%llu]",
     265              :                        ret, kv.first, GetReadableRequestType(requestType), sendMgr.sentSize_);
     266              :         }
     267              :     }
     268            0 : }
     269              : 
     270            0 : void ZeroCopyMemoryAgent::CheckBatchSendAsyncResult()
     271              : {
     272              :     HcclResult ret;
     273              :     HcclResult lastSendRet;
     274            0 :     for (auto &kv : sendMgrs_) {
     275            0 :         auto &sendMgr = kv.second;
     276            0 :         if (sendMgr.lastSendHandle_ == nullptr) {  // 没有正在执行的异步send
     277            0 :             continue;
     278              :         }
     279              : 
     280            0 :         if (mapDevPhyIdconnectedSockets_.find(kv.first) == mapDevPhyIdconnectedSockets_.end()) {
     281            0 :             HCCL_ERROR("[ZeroCopyMemoryAgent][CheckBatchSendAsyncResult] remote[%u] not found in"
     282              :                 "mapDevPhyIdconnectedSockets_", kv.first);
     283            0 :             continue;
     284              :         }
     285            0 :         auto &socket = mapDevPhyIdconnectedSockets_[kv.first];
     286            0 :         ret = socket->GetAsyncReqResult(sendMgr.lastSendHandle_, lastSendRet);
     287            0 :         if (ret != HCCL_SUCCESS) {
     288            0 :             CHK_PRT_CONT(ret != HCCL_E_AGAIN,
     289              :                 HCCL_ERROR("[ZeroCopyMemoryAgent][CheckBatchSendAsyncResult]GetAsyncReqResult failed, ret[%d] remote[%u]",
     290              :                             ret, kv.first));
     291            0 :             continue;
     292              :         }
     293              : 
     294            0 :         sendMgr.lastSendHandle_ = nullptr;
     295            0 :         if ((lastSendRet != HCCL_SUCCESS) && (sendMgr.lastSendSize_ == 0)) {
     296            0 :             CHK_PRT_CONT(lastSendRet != HCCL_E_AGAIN,
     297              :                 HCCL_ERROR("[ZeroCopyMemoryAgent][CheckBatchSendAsyncResult]SendAsync failed, result[%d] remote[%u] sentSize[%llu]",
     298              :                             lastSendRet, kv.first, sendMgr.sentSize_));
     299            0 :             continue;
     300              :         }
     301              : 
     302            0 :         sendMgr.sentSize_ += sendMgr.lastSendSize_;  // 下次从中断的地方开始重发
     303            0 :         if (sendMgr.sentSize_ == sendMgr.reqDataSize_) {  // request发送完成
     304            0 :             sendMgr.sentSize_ = 0;
     305            0 :             sendMgr.hasReq_[sendMgr.currIndex_] = false;
     306            0 :             HCCL_DEBUG("[ZeroCopyMemoryAgent][CheckBatchSendAsyncResult]SendAsync success, requestType[%s] remote[%u]",
     307              :                     GetReadableRequestType(*reinterpret_cast<const RequestType *>(sendMgr.reqDatas_[sendMgr.currIndex_]->data())), kv.first);
     308              :         }
     309              :     }
     310            0 : }
     311              : 
     312            0 : void ZeroCopyMemoryAgent::RequestBatchRecvAsync()
     313              : {
     314              :     HcclResult ret;
     315            0 :     for (auto &kv : recvMgrs_) {
     316            0 :         auto &recvMgr = kv.second;
     317            0 :         if ((recvMgr.lastRecvHandle_ != nullptr) ||  // 前回接收未完成
     318            0 :             ((receivedBarrierClose_.count(kv.first) != 0) && (receivedBarrierCloseAck_.count(kv.first) != 0))) {
     319              :             // 该socket已经收到BarrierClose与BarrierCloseAck报文,因此不允许再进行其他数据接收了
     320            0 :             continue;
     321              :         }
     322              : 
     323            0 :         if (mapDevPhyIdconnectedSockets_.find(kv.first) == mapDevPhyIdconnectedSockets_.end()) {
     324            0 :             HCCL_ERROR("[ZeroCopyMemoryAgent][RequestBatchRecvAsync] remote[%u] not found in"
     325              :                 "mapDevPhyIdconnectedSockets_", kv.first);
     326            0 :             continue;
     327              :         }
     328            0 :         auto &socket = mapDevPhyIdconnectedSockets_[kv.first];
     329            0 :         std::vector<u8> &req = recvMgr.receivedData_[recvMgr.recvIndex_];
     330            0 :         recvMgr.lastRecvSize_ = 0;  // 用于ra上报接收的数据量
     331            0 :         ret = socket->RecvAsync(req.data() + recvMgr.receivedSize_,
     332            0 :             IPC_MEMORY_EXCHANGE_LENGTH - recvMgr.receivedSize_, &recvMgr.lastRecvSize_, &recvMgr.lastRecvHandle_);
     333            0 :         CHK_PRT_CONT((ret != HCCL_SUCCESS) && (ret != HCCL_E_AGAIN),
     334              :             HCCL_ERROR("[ZeroCopyMemoryAgent][RequestBatchRecvAsync] RecvAsync failed, ret[%d] remote[%u] receivedSize[%llu]",
     335              :                 ret, kv.first, recvMgr.receivedSize_));
     336              :     }
     337            0 : }
     338              : 
     339            0 : void ZeroCopyMemoryAgent::CheckBatchRecvAsyncResult()
     340              : {
     341              :     HcclResult ret;
     342              :     HcclResult lastRecvRet;
     343            0 :     for (auto &kv : recvMgrs_) {
     344            0 :         auto &recvMgr = kv.second;
     345            0 :         if (recvMgr.lastRecvHandle_ == nullptr) {  // 没有正在异步接收
     346            0 :             continue;
     347              :         }
     348              : 
     349            0 :         if (mapDevPhyIdconnectedSockets_.find(kv.first) == mapDevPhyIdconnectedSockets_.end()) {
     350            0 :             HCCL_ERROR("[ZeroCopyMemoryAgent][CheckBatchRecvAsyncResult] remote[%u] not found in"
     351              :                 "mapDevPhyIdconnectedSockets_", kv.first);
     352            0 :             continue;
     353              :         }
     354            0 :         auto &socket = mapDevPhyIdconnectedSockets_[kv.first];
     355            0 :         ret = socket->GetAsyncReqResult(recvMgr.lastRecvHandle_, lastRecvRet);
     356            0 :         if (ret != HCCL_SUCCESS) {
     357            0 :             CHK_PRT_CONT(ret != HCCL_E_AGAIN,
     358              :                 HCCL_ERROR("[ZeroCopyMemoryAgent][CheckBatchRecvAsyncResult] GetAsyncReqResult failed, ret[%d] remote[%u]",
     359              :                             ret, kv.first));
     360            0 :             continue;
     361              :         }
     362              : 
     363            0 :         recvMgr.lastRecvHandle_ = nullptr;
     364            0 :         if ((lastRecvRet != HCCL_SUCCESS) && (recvMgr.lastRecvSize_ == 0)) {
     365            0 :             CHK_PRT_CONT(lastRecvRet != HCCL_E_AGAIN,
     366              :                 HCCL_WARNING("[ZeroCopyMemoryAgent][CheckBatchRecvAsyncResult] RecvAsync failed, result[%d] remote[%u] lastRecvSize[%llu]",
     367              :                     lastRecvRet, kv.first, recvMgr.lastRecvSize_));
     368            0 :             continue;
     369              :         }
     370              : 
     371            0 :         recvMgr.receivedSize_ += recvMgr.lastRecvSize_;
     372            0 :         if (recvMgr.receivedSize_ == IPC_MEMORY_EXCHANGE_LENGTH) {
     373            0 :             recvMgr.receivedSize_ = 0;
     374            0 :             RecvRequest(recvMgr, kv.first);
     375            0 :             ioRecvWaiting_ = true;  // 后面高概率还有数据要收(ack与request合并场景),loop不等待
     376              :         } else {
     377              :             // request没收全,loop不等待
     378            0 :             ioRecvWaiting_ = (recvMgr.receivedSize_ > 0);
     379              :         }
     380              :     }
     381            0 : }
     382              : 
     383            0 : inline void ZeroCopyMemoryAgent::RecvRequest(ZeroCopyMemoryAgentRecvMgr &recvMgr, u32 remoteDevicePhyId)
     384              : {
     385            0 :     std::vector<u8> &req = recvMgr.receivedData_[recvMgr.recvIndex_];
     386            0 :     RequestType requestType = *reinterpret_cast<RequestType *>(req.data());
     387            0 :     HCCL_DEBUG("[ZeroCopyMemoryAgent][RecvRequest] recv requestType[%s] remote[%u]",
     388              :         GetReadableRequestType(requestType), remoteDevicePhyId);
     389              : 
     390            0 :     if (IsAckRequestType(requestType)) {  // 收到ACK时,直接优先处理
     391            0 :         u32 remoteRank = mapDevPhyId2RankId_[remoteDevicePhyId];
     392            0 :         CHK_PRT_CONT(ParseReceivedRequest(req, remoteRank) != HCCL_SUCCESS,
     393              :                 HCCL_ERROR("[ZeroCopyMemoryAgent][ParseReceivedRequest] failed requestType[%s] remote[%u]",
     394              :                     GetReadableRequestType(requestType), remoteDevicePhyId));
     395            0 :         return;
     396              :     }
     397              : 
     398            0 :     recvMgr.recvIndex_ = (recvMgr.recvIndex_ + 1) % ZERO_COPY_MEMORY_AGENT_RECV_QUEUE_SIZE;  // 准备下一次接收
     399            0 :     hasReceivedRequest_ = true;
     400              : }
     401              : 
     402            0 : void ZeroCopyMemoryAgent::ParseReceivedRequests()
     403              : {
     404            0 :     if (!hasReceivedRequest_) {
     405            0 :         return;
     406              :     }
     407            0 :     hasReceivedRequest_ = false;
     408              : 
     409            0 :     for (auto &kv : recvMgrs_) {
     410            0 :         u32 remoteRank = mapDevPhyId2RankId_[kv.first];
     411            0 :         auto &recvMgr = kv.second;
     412            0 :         while (recvMgr.praseIndex_ != recvMgr.recvIndex_) {
     413            0 :             std::vector<u8> &req = recvMgr.receivedData_[recvMgr.praseIndex_];
     414            0 :             CHK_PRT_CONT(ParseReceivedRequest(req, remoteRank) != HCCL_SUCCESS,
     415              :                     HCCL_ERROR("[ZeroCopyMemoryAgent][ParseReceivedRequest] failed prase requestType[%s] remote[%u]",
     416              :                         GetReadableRequestType(*reinterpret_cast<RequestType *>(req.data())), kv.first));
     417            0 :             recvMgr.praseIndex_++;
     418            0 :             if (recvMgr.praseIndex_ == ZERO_COPY_MEMORY_AGENT_RECV_QUEUE_SIZE) {
     419            0 :                 recvMgr.praseIndex_ = 0;
     420              :             }
     421              :         }
     422              :     }
     423              : }
     424              : 
     425            0 : void ZeroCopyMemoryAgent::RequestBatchRecvSync()
     426              : {
     427              :     HcclResult ret;
     428            0 :     for (auto &kv : recvMgrs_) {
     429            0 :         auto &recvMgr = kv.second;
     430            0 :         if ((receivedBarrierClose_.count(kv.first) != 0) && (receivedBarrierCloseAck_.count(kv.first) != 0)) {
     431              :             // 该socket已经收到BarrierClose与BarrierCloseAck报文,因此不允许再进行其他数据接收了
     432            0 :             continue;
     433              :         }
     434              : 
     435            0 :         if (mapDevPhyIdconnectedSockets_.find(kv.first) == mapDevPhyIdconnectedSockets_.end()) {
     436            0 :             HCCL_ERROR("[ZeroCopyMemoryAgent][RequestBatchRecvSync] remote[%u] not found in"
     437              :                 "mapDevPhyIdconnectedSockets_", kv.first);
     438            0 :             continue;
     439              :         }
     440            0 :         auto &socket = mapDevPhyIdconnectedSockets_[kv.first];
     441            0 :         std::vector<u8> &req = recvMgr.receivedData_[0];
     442            0 :         recvMgr.lastRecvSize_ = 0;
     443            0 :         ret = socket->IRecv(req.data() + recvMgr.receivedSize_, IPC_MEMORY_EXCHANGE_LENGTH - recvMgr.receivedSize_,
     444            0 :                             recvMgr.lastRecvSize_);
     445            0 :         CHK_PRT_CONT((ret != HCCL_SUCCESS) && (ret != HCCL_E_AGAIN),
     446              :             HCCL_ERROR("[ZeroCopyMemoryAgent][RequestBatchRecvSync] IRecv failed, ret[%d] remote[%u] receivedSize[%llu]",
     447              :                     ret, kv.first, recvMgr.receivedSize_));
     448              : 
     449            0 :         recvMgr.receivedSize_ += recvMgr.lastRecvSize_;
     450            0 :         if (recvMgr.receivedSize_ == IPC_MEMORY_EXCHANGE_LENGTH) {
     451            0 :             recvMgr.receivedSize_ = 0;
     452            0 :             ret = ParseReceivedRequest(req, mapDevPhyId2RankId_[kv.first]);
     453            0 :             CHK_PRT_CONT(ret != HCCL_SUCCESS, HCCL_ERROR("[ZeroCopyMemoryAgent][ParseReceivedRequest] failed"));
     454              :         }
     455              :     }
     456            0 : }
     457              : 
     458            0 : void ZeroCopyMemoryAgent::InnerThread()
     459              : {
     460              :     // 新线程,更新一下使用的设备
     461            0 :     if (hrtSetDevice(deviceLogicId_) != HCCL_SUCCESS) {
     462            0 :         HCCL_ERROR("[ZeroCopyMemoryAgent][InnerThread] set device failed");
     463            0 :         return;
     464              :     }
     465              : 
     466            0 :     while (threadRun_) {
     467            0 :         CheckSnapshotStatus();
     468            0 :         if (isPaused_) {
     469            0 :             SaluSleep(USLEEP_ONE_THOUSAND);
     470            0 :             continue;
     471              :         }
     472              : 
     473            0 :         if (isSocketSupportAsync_) {
     474            0 :             CheckBatchSendAsyncResult();
     475            0 :             RequestBatchSendAsync();
     476              : 
     477            0 :             CheckBatchRecvAsyncResult();
     478            0 :             RequestBatchRecvAsync();
     479              : 
     480            0 :             ParseReceivedRequests();
     481              : 
     482            0 :             std::unique_lock<std::mutex> lock(sendMutex_);
     483            0 :             if (!ioRecvWaiting_ && !hasSendRequest_) {
     484            0 :                 sendCv_.wait_for(lock, std::chrono::microseconds(INNER_THREAD_LOOP_US));
     485              :             }
     486            0 :             hasSendRequest_ = false;
     487            0 :             ioRecvWaiting_ = false;
     488            0 :         } else {
     489            0 :             RequestBatchRecvSync();
     490            0 :             SaluSleep(USLEEP_ONE_THOUSAND);
     491              :         }
     492              :     }
     493              : 
     494            0 :     if (hrtResetDevice(deviceLogicId_) != HCCL_SUCCESS) {
     495            0 :         HCCL_ERROR("[ZeroCopyMemoryAgent][InnerThread] reset device failed");
     496            0 :         return;
     497              :     }
     498              : }
     499              : 
     500            0 : HcclResult ZeroCopyMemoryAgent::SetRemoteTgid()
     501              : {
     502            0 :     if (remotePids_.size() == mapDevPhyIdconnectedSockets_.size()) {
     503            0 :         HCCL_INFO("[ZeroCopyMemoryAgent][SetRemoteTgid] tgid exchange is ok");
     504            0 :         return HCCL_SUCCESS;
     505              :     }
     506            0 :     remotePids_.clear();
     507              : 
     508            0 :     u8 *exchangeDataPtr = exchangeDataForSend_.data();
     509            0 :     u32 exchangeDataBlankSize = IPC_MEMORY_EXCHANGE_LENGTH;
     510              : 
     511            0 :     RequestType requestType = RequestType::SET_REMOTE_BARE_TGID;
     512              : 
     513            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, requestType));
     514              : 
     515            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId_));
     516              : 
     517            0 :     CHK_RET(SendRequest(requestType, exchangeDataForSend_));
     518              : 
     519            0 :     CHK_RET(WaitForAllRemoteComplete(RequestType::SET_REMOTE_BARE_TGID_ACK));
     520            0 :     if (remotePids_.size() != mapDevPhyIdconnectedSockets_.size()) {
     521            0 :         HCCL_ERROR("[ZeroCopyMemoryAgent][SetRemoteTgid] tgid exchange failed recv pids count[%lu]", remotePids_.size());
     522            0 :         return HCCL_E_INTERNAL;
     523              :     }
     524            0 :     return HCCL_SUCCESS;
     525              : }
     526              : 
     527            0 : HcclResult ZeroCopyMemoryAgent::DeInit()
     528              : {
     529            0 :     CHK_PRT_RET(isSingleRank_, HCCL_INFO("[ZeroCopyMemoryAgent][DeInit] single rank communicator"), HCCL_SUCCESS);
     530            0 :     std::unique_lock<std::mutex> lock(commRefCntLock_);
     531            0 :     if (!ZeroCopyMemoryAgent::IsAddressMgrInited()) {
     532            0 :         HCCL_ERROR("[ZeroCopyMemoryAgent][%s]addressMgr_ is nullptr, no need to deinit. local rank[u32]", __func__,
     533              :             userRank_);
     534            0 :         return HCCL_E_INTERNAL;
     535              :     }
     536            0 :     threadRun_ = false;
     537            0 :     if (innerThread_) {
     538            0 :         if (innerThread_->joinable()) {
     539            0 :             innerThread_->join();  // 等待线程执行后释放资源
     540              :         }
     541              :     }
     542            0 :     innerThread_ = nullptr;
     543              : 
     544            0 :     if (vnicPortCtx_ != nullptr) {
     545            0 :         HcclNetCloseDev(vnicPortCtx_);
     546            0 :         vnicPortCtx_ = nullptr;
     547              :     }
     548            0 :     CHK_RET(addressMgr_->DecreCommRefCnt());
     549            0 :     if (addressMgr_->GetCommRefCnt() == 0) {
     550            0 :         addressMgr_.reset();
     551            0 :         HCCL_RUN_INFO("[ZeroCopyMemoryAgent][%s]Release addressMgr_", __func__);
     552              :     }
     553            0 :     return HCCL_SUCCESS;
     554            0 : }
     555              : 
     556            0 : HcclResult ZeroCopyMemoryAgent::SetMemoryRange(void *virPtr, size_t size, size_t alignment, uint64_t flags)
     557              : {
     558            0 :     CHK_PRT_RET(isSingleRank_, HCCL_INFO("[ZeroCopyMemoryAgent][SetMemoryRange] single rank communicator"), HCCL_SUCCESS);
     559            0 :     CHK_PRT_RET(!ZeroCopyMemoryAgent::IsAddressMgrInited(), HCCL_ERROR("[ZeroCopyMemoryAgent][%s]ZeroCopyMemoryAgent "
     560              :         "is not init.", __func__), HCCL_E_INTERNAL);
     561            0 :     CHK_PRT_RET(addressMgr_->SetMemoryRange(devicePhyId_, virPtr, size) != HCCL_SUCCESS,
     562              :         HCCL_ERROR("[ZeroCopyMemoryAgent][SetMemoryRange] invalid set ptr[%p] size[%lu] alignment[%lu] flags[%lu]",
     563              :         virPtr, size, alignment, flags), HCCL_E_PARA);
     564              : 
     565            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][SetMemoryRange] basePtr[%p] size[%lu] alignment[%lu] flag[%lu]",
     566              :         virPtr, size, alignment, flags);
     567            0 :     u8 *exchangeDataPtr = exchangeDataForSend_.data();
     568            0 :     u32 exchangeDataBlankSize = IPC_MEMORY_EXCHANGE_LENGTH;
     569              : 
     570            0 :     RequestType requestType = RequestType::SET_MEMORY_RANGE;
     571              : 
     572            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, requestType));
     573              : 
     574            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId_));
     575              : 
     576            0 :     u64 addr = reinterpret_cast<u64>(virPtr);
     577            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, addr));
     578              : 
     579            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, size));
     580              : 
     581            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, alignment));
     582              : 
     583            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, flags));
     584              : 
     585            0 :     CHK_RET(SendRequest(requestType, exchangeDataForSend_));
     586              : 
     587            0 :     CHK_RET(WaitForAllRemoteComplete(RequestType::SET_MEMORY_RANGE_ACK));
     588            0 :     return HCCL_SUCCESS;
     589              : }
     590              : 
     591            0 : HcclResult ZeroCopyMemoryAgent::UnsetMemoryRange(void *virPtr)
     592              : {
     593            0 :     CHK_PRT_RET(isSingleRank_, HCCL_INFO("[ZeroCopyMemoryAgent][UnsetMemoryRange] single rank communicator"), HCCL_SUCCESS);
     594            0 :     CHK_PRT_RET(!ZeroCopyMemoryAgent::IsAddressMgrInited(), HCCL_ERROR("[ZeroCopyMemoryAgent][%s]ZeroCopyMemoryAgent "
     595              :         "is not init.", __func__), HCCL_E_INTERNAL);
     596            0 :     CHK_PRT_RET(!addressMgr_->IsAddressSet(devicePhyId_, virPtr),
     597              :         HCCL_ERROR("[ZeroCopyMemoryAgent][UnsetMemoryRange] ptr[%p] is not set memory", virPtr), HCCL_E_PARA);
     598            0 :     CHK_RET(addressMgr_->UnsetMemoryRange(devicePhyId_, virPtr));
     599              : 
     600            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][UnsetMemoryRange] basePtr[%p]", virPtr);
     601            0 :     u8 *exchangeDataPtr = exchangeDataForSend_.data();
     602            0 :     u32 exchangeDataBlankSize = IPC_MEMORY_EXCHANGE_LENGTH;
     603              : 
     604            0 :     RequestType requestType = RequestType::UNSET_MEMORY_RANGE;
     605            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, requestType));
     606              : 
     607            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId_));
     608              : 
     609            0 :     u64 addr = reinterpret_cast<u64>(virPtr);
     610            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, addr));
     611              : 
     612            0 :     CHK_RET(SendRequest(requestType, exchangeDataForSend_));
     613              : 
     614            0 :     CHK_RET(WaitForAllRemoteComplete(RequestType::UNSET_MEMORY_RANGE_ACK));
     615            0 :     return HCCL_SUCCESS;
     616              : }
     617              : 
     618            0 : HcclResult ZeroCopyMemoryAgent::ActivateCommMemory(void *virPtr, size_t size, size_t offset, void *memHandle, uint64_t flags)
     619              : {
     620            0 :     CHK_PRT_RET(isSingleRank_, HCCL_INFO("[ZeroCopyMemoryAgent][ActivateCommMemory] single rank communicator"), HCCL_SUCCESS);
     621            0 :     CHK_PRT_RET(!ZeroCopyMemoryAgent::IsAddressMgrInited(), HCCL_ERROR("[ZeroCopyMemoryAgent][%s]ZeroCopyMemoryAgent "
     622              :         "is not init.", __func__), HCCL_E_INTERNAL);
     623            0 :     CHK_PRT_RET(!addressMgr_->IsInSetAddressRange(devicePhyId_, virPtr, size),
     624              :         HCCL_ERROR("[ZeroCopyMemoryAgent][ActivateCommMemory] input ptr[%p] size[%lu] is not in set address range", virPtr, size), HCCL_E_PARA);
     625            0 :     CHK_PRT_RET(addressMgr_->IsOverlapWithActivateAddr(virPtr, size),
     626              :         HCCL_ERROR("[ZeroCopyMemoryAgent][ActivateCommMemory] input ptr[%p] size[%lu] overlap with activate memory", virPtr, size), HCCL_E_PARA);
     627              : 
     628            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][ActivateCommMemory] virPtr[%p] size[%lu] offset[%lu] memHandle[%p], flags[%lu]",
     629              :         virPtr, size, offset, memHandle, flags);
     630            0 :     CHK_RET(SetRemoteTgid());
     631              : 
     632              :     uint64_t shareableHandle;
     633            0 :     aclrtMemHandleType handleType = ACL_MEM_HANDLE_TYPE_NONE;
     634            0 :     aclError ret = ACL_SUCCESS;
     635            0 :     ret = aclrtMemExportToShareableHandle(memHandle, handleType, 0, &shareableHandle);
     636            0 :     CHK_PRT_RET(ret != ACL_SUCCESS, HCCL_ERROR("[ZeroCopyMemoryAgent][ActivateCommMemory] aclrtMemExportToShareableHandle handle[%p] type[%d] flags[%llu] failed, ret[%d]",
     637              :         memHandle, handleType, 0, ret), HCCL_E_RUNTIME);
     638            0 :     ret = aclrtMemSetPidToShareableHandle(shareableHandle, remotePids_.data(), remotePids_.size());
     639            0 :     CHK_PRT_RET(ret != ACL_SUCCESS, HCCL_ERROR("[ZeroCopyMemoryAgent][ActivateCommMemory] aclrtMemSetPidToShareableHandle shareableHandl[%llu]",
     640              :         " failed, ret[%d]", shareableHandle, ret), HCCL_E_RUNTIME);
     641              : 
     642            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][ActivateCommMemory] dev[%u] export shareableHandle[%lu]", devicePhyId_, shareableHandle);
     643            0 :     u8 *exchangeDataPtr = exchangeDataForSend_.data();
     644            0 :     u32 exchangeDataBlankSize = IPC_MEMORY_EXCHANGE_LENGTH;
     645              : 
     646            0 :     RequestType requestType = RequestType::ACTIVATE_COMM_MEMORY;
     647            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, requestType));
     648              : 
     649            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId_));
     650              : 
     651            0 :     u64 addr = reinterpret_cast<u64>(virPtr);
     652            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, addr));
     653              : 
     654            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, size));
     655              : 
     656            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, offset));
     657              : 
     658            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, shareableHandle));
     659              : 
     660            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, flags));
     661              : 
     662            0 :     CHK_RET(SendRequest(requestType, exchangeDataForSend_));
     663              : 
     664            0 :     CHK_RET(WaitForAllRemoteComplete(RequestType::ACTIVATE_COMM_MEMORY_ACK));
     665            0 :     CHK_RET(addressMgr_->ActivateCommMemoryAddr(virPtr, size));
     666              : 
     667            0 :     return HCCL_SUCCESS;
     668              : }
     669              : 
     670            0 : HcclResult ZeroCopyMemoryAgent::DeactivateCommMemory(void *virPtr)
     671              : {
     672            0 :     CHK_PRT_RET(isSingleRank_, HCCL_INFO("[ZeroCopyMemoryAgent][DeactivateCommMemory] single rank communicator"), HCCL_SUCCESS);
     673            0 :     CHK_PRT_RET(!ZeroCopyMemoryAgent::IsAddressMgrInited(), HCCL_ERROR("[ZeroCopyMemoryAgent][%s]ZeroCopyMemoryAgent "
     674              :         "is not init.", __func__), HCCL_E_INTERNAL);
     675            0 :     CHK_PRT_RET(!addressMgr_->IsActivateCommMemoryAddr(virPtr, 1),
     676              :         HCCL_ERROR("[ZeroCopyMemoryAgent][DeactivateCommMemory] input ptr[%p] is not activate", virPtr), HCCL_E_PARA);
     677              : 
     678            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][DeactivateCommMemory] virPtr[%p]", virPtr);
     679            0 :     CHK_RET(addressMgr_->DeactivateCommMemoryAddr(virPtr));
     680              : 
     681            0 :     u8 *exchangeDataPtr = exchangeDataForSend_.data();
     682            0 :     u32 exchangeDataBlankSize = IPC_MEMORY_EXCHANGE_LENGTH;
     683              : 
     684            0 :     RequestType requestType = RequestType::DEACTIVATE_COMM_MEMORY;
     685            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, requestType));
     686              : 
     687            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId_));
     688              : 
     689            0 :     u64 addr = reinterpret_cast<u64>(virPtr);
     690            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, addr));
     691              : 
     692            0 :     CHK_RET(SendRequest(requestType, exchangeDataForSend_));
     693              : 
     694            0 :     CHK_RET(WaitForAllRemoteComplete(RequestType::DEACTIVATE_COMM_MEMORY_ACK));
     695            0 :     return HCCL_SUCCESS;
     696              : }
     697              : 
     698            0 : HcclResult ZeroCopyMemoryAgent::BarrierClose()
     699              : {
     700            0 :     CHK_PRT_RET(isSingleRank_, HCCL_INFO("[ZeroCopyMemoryAgent][BarrierClose] single rank communicator"), HCCL_SUCCESS);
     701              : 
     702            0 :     HCCL_RUN_INFO("[ZeroCopyMemoryAgent][BarrierClose] [%s] ready to barrier close", identifier_.c_str());
     703            0 :     u8 *exchangeDataPtr = exchangeDataForSend_.data();
     704            0 :     u32 exchangeDataBlankSize = IPC_MEMORY_EXCHANGE_LENGTH;
     705              : 
     706            0 :     RequestType requestType = RequestType::BARRIER_CLOSE;
     707            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, requestType));
     708            0 :     CHK_RET(ConstructData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId_));
     709              : 
     710            0 :     CHK_RET(SendRequest(requestType, exchangeDataForSend_));
     711              : 
     712            0 :     CHK_RET(WaitForAllRemoteComplete(RequestType::BARRIER_CLOSE_ACK));
     713              : 
     714            0 :     return HCCL_SUCCESS;
     715              : }
     716              : 
     717           40 : bool ZeroCopyMemoryAgent::IsActivateCommMemoryAddr(void *virPtr, u64 length)
     718              : {
     719           40 :     if (!ZeroCopyMemoryAgent::IsAddressMgrInited()) {
     720           40 :         HCCL_INFO("[ZeroCopyMemoryAgent][%s]ZeroCopyMemoryAgent is not init.", __func__);
     721           40 :         return false;
     722              :     }
     723            0 :     return addressMgr_->IsActivateCommMemoryAddr(virPtr, length);
     724              : }
     725              : 
     726            0 : HcclResult ZeroCopyMemoryAgent::GetRingBufferAddr(u64 &bufferPtr, u64 &headPtr, u64 &tailPtr)
     727              : {
     728            0 :     CHK_PRT_RET(!ZeroCopyMemoryAgent::IsAddressMgrInited(), HCCL_ERROR("[ZeroCopyMemoryAgent][%s]ZeroCopyMemoryAgent "
     729              :         "is not init.", __func__), HCCL_E_INTERNAL);
     730            0 :     addressMgr_->GetRingBufferAddr(bufferPtr, headPtr, tailPtr);
     731            0 :     return HCCL_SUCCESS;
     732              : }
     733              : 
     734           40 : bool ZeroCopyMemoryAgent::IsAddressMgrInited()
     735              : {
     736           40 :     return addressMgr_ != nullptr;
     737              : }
     738              : 
     739            0 : HcclResult ZeroCopyMemoryAgent::WaitForAllRemoteComplete(RequestType requestType)
     740              : {
     741            0 :     bool useBarrier = NeedBarrier(requestType);
     742            0 :     if (useBarrier) {
     743            0 :         reqMsgDeliverCnt_++;
     744              :     }
     745              : 
     746            0 :     u32 expectedNum = mapDevPhyIdconnectedSockets_.size();
     747            0 :     auto timeout = std::chrono::seconds(GetExternalInputHcclLinkTimeOut());
     748            0 :     std::unique_lock<std::mutex> lock(dfxMutex_);
     749            0 :     waitCompleteCv_.wait_for(lock, timeout);
     750            0 :     if ((reqMsgCounter_[static_cast<int>(requestType)] == expectedNum) &&
     751            0 :         (!useBarrier || (useBarrier && reqMsgDeliverCnt_ <= reqMsgFinishCnt_))) {
     752            0 :         reqMsgCounter_[static_cast<int>(requestType)] = 0;
     753            0 :         reqMsgFinishedRanks_[static_cast<int>(requestType)].clear();
     754            0 :         return HCCL_SUCCESS;
     755              :     }
     756              : 
     757            0 :     HCCL_ERROR("[Wait][RemoteComplete %s] dev[%u] errNo[0x%016llx] timeout[%d s] completeCount[%u] %s",
     758              :             GetReadableRequestType(requestType), devicePhyId_,
     759              :             HCCL_ERROR_CODE(HCCL_E_TCP_TRANSFER), timeout, reqMsgCounter_[static_cast<int>(requestType)].load(),
     760              :             DumpFinishInfo(requestType).c_str());
     761            0 :     reqMsgCounter_[static_cast<int>(requestType)] = 0;
     762            0 :     reqMsgFinishedRanks_[static_cast<int>(requestType)].clear();
     763            0 :     return HCCL_E_TCP_TRANSFER;
     764            0 : }
     765              : 
     766            0 : HcclResult ZeroCopyMemoryAgent::ParseSetMemoryRange(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize)
     767              : {
     768            0 :     CHK_PRT_RET(!ZeroCopyMemoryAgent::IsAddressMgrInited(), HCCL_ERROR("[ZeroCopyMemoryAgent][%s]ZeroCopyMemoryAgent "
     769              :         "is not init.", __func__), HCCL_E_INTERNAL);
     770              :     u32 devicePhyId;
     771            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId));
     772              : 
     773              :     u64 addr;
     774            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, addr));
     775              : 
     776              :     size_t size;
     777            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, size));
     778              : 
     779              :     size_t alignment;
     780            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, alignment));
     781              : 
     782              :     uint64_t flags;
     783            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, flags));
     784              : 
     785              :     u32 maxDeviceNum;
     786            0 :     CHK_RET(GetMaxDevNum(maxDeviceNum));
     787            0 :     CHK_PRT_RET(devicePhyId >= maxDeviceNum,
     788              :         HCCL_ERROR("[ZeroCopyMemoryAgent][ParseSetMemoryRange] devicePhyId[%u] is exceed max device num[%u]", devicePhyId, maxDeviceNum),
     789              :         HCCL_E_PARA);
     790              : 
     791            0 :     void *remoteAddrBase = reinterpret_cast<void *>(addr);
     792            0 :     CHK_PRT_RET(addressMgr_->IsAddressSet(devicePhyId, remoteAddrBase),
     793              :         HCCL_ERROR("[ZeroCopyMemoryAgent][ParseSetMemoryRange] devicePhyId[%u] had set addr [%p]", devicePhyId, remoteAddrBase), HCCL_E_PARA);
     794              : 
     795            0 :     void* devPtr = nullptr;
     796            0 :     void* devAddr = nullptr;
     797            0 :     aclError ret = aclrtReserveMemAddress(&devPtr, size, alignment, devAddr, flags);
     798            0 :     CHK_PRT_RET(ret != ACL_SUCCESS, HCCL_ERROR("[ZeroCopyMemoryAgent][ParseSetMemoryRange] rtReserve Memory failed, "
     799              :         "return[%d], devPtr[%p] size[%llu] alignment[%llu] devAddr[%p] flags[%llu]",
     800              :         ret, devPtr, size, alignment, devAddr, flags), HCCL_E_RUNTIME);
     801              : 
     802            0 :     CHK_RET(addressMgr_->AddLocalIpc2RemoteAddr(devicePhyId, devPtr, reinterpret_cast<void *>(addr), size));
     803              : 
     804            0 :     CHK_RET(SendAckAfterParse(RequestType::SET_MEMORY_RANGE, RequestType::SET_MEMORY_RANGE_ACK, devicePhyId));
     805              : 
     806            0 :     return HCCL_SUCCESS;
     807              : }
     808              : 
     809            0 : HcclResult ZeroCopyMemoryAgent::SendAckAfterParse(RequestType requestType, RequestType ackType, u32 remoteDevicePhyId,
     810              :     void *extraData, u64 extraDataLen)
     811              : {
     812            0 :     u8 *exchangeDataAckPtr = exchangeDataForAck_[remoteDevicePhyId].data();
     813            0 :     u32 exchangeDataAckBlankSize = IPC_MEMORY_EXCHANGE_LENGTH;
     814              : 
     815            0 :     CHK_RET(ConstructData(exchangeDataAckPtr, exchangeDataAckBlankSize, ackType));
     816              : 
     817            0 :     CHK_RET(ConstructData(exchangeDataAckPtr, exchangeDataAckBlankSize, devicePhyId_));
     818              : 
     819            0 :     if (extraData != nullptr && extraDataLen != 0) {
     820            0 :         CHK_RET(ConstructData(exchangeDataAckPtr, exchangeDataAckBlankSize, extraData, extraDataLen));
     821              :     }
     822              : 
     823              :     // 不需要进行barrier,那么我们每处理一个请求就回复一个请求
     824            0 :     if (!NeedBarrier(requestType)) {
     825            0 :         CHK_PRT_RET(SendRequest(ackType, exchangeDataForAck_[remoteDevicePhyId], remoteDevicePhyId) != HCCL_SUCCESS,
     826              :             HCCL_WARNING("[ZeroCopyMemoryAgent][SendAckAfterParse] failed, remote[%u]", remoteDevicePhyId),
     827              :             HCCL_E_INTERNAL);
     828            0 :         return HCCL_SUCCESS;
     829              :     }
     830              : 
     831              :     // 需要进行barrier的请求,我们先统计一下收到的请求数目,等于链接数才算收完所有
     832            0 :     u32 expectedNum = mapDevPhyIdconnectedSockets_.size();
     833            0 :     u32 counter = ++reqMsgCounter_[static_cast<int>(requestType)];
     834            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][SendAckAfterParse] requestType[%d] counter %u expect %u", requestType, counter, expectedNum);
     835            0 :     if (counter < expectedNum) {
     836            0 :         return HCCL_SUCCESS;
     837              :     } else {
     838            0 :         reqMsgCounter_[static_cast<int>(requestType)] = 0;
     839            0 :         reqMsgFinishCnt_++;
     840              : 
     841              :         // 我们统一将所有的请求一次性都发送过去
     842            0 :         CHK_PRT_RET(SendRequest(ackType, exchangeDataForAck_[remoteDevicePhyId]) != HCCL_SUCCESS,
     843              :             HCCL_WARNING("[ZeroCopyMemoryAgent][SendAckAfterParse] failed, remote[all]"), HCCL_E_INTERNAL);
     844              :     }
     845              : 
     846            0 :     return HCCL_SUCCESS;
     847              : }
     848              : 
     849              : 
     850            0 : HcclResult ZeroCopyMemoryAgent::ParseRemoteAck(RequestType requestType, u32 remoteRank)
     851              : {
     852            0 :     bool useBarrier = NeedBarrier(requestType);
     853            0 :     std::unique_lock<std::mutex> dfxLock(dfxMutex_);
     854            0 :     reqMsgFinishedRanks_[static_cast<int>(requestType)].insert(remoteRank);
     855            0 :     u32 counter = ++reqMsgCounter_[static_cast<int>(requestType)];
     856            0 :     if ((counter == mapDevPhyIdconnectedSockets_.size()) &&
     857            0 :         (!useBarrier || (useBarrier && reqMsgDeliverCnt_ <= reqMsgFinishCnt_))) {
     858            0 :         waitCompleteCv_.notify_all();
     859              :     }
     860            0 :     return HCCL_SUCCESS;
     861            0 : }
     862              : 
     863            0 : HcclResult ZeroCopyMemoryAgent::ParseUnsetMemoryRange(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize)
     864              : {
     865            0 :     CHK_PRT_RET(!ZeroCopyMemoryAgent::IsAddressMgrInited(), HCCL_ERROR("[ZeroCopyMemoryAgent][%s]ZeroCopyMemoryAgent "
     866              :         "is not init.", __func__), HCCL_E_INTERNAL);
     867              :     u32 devicePhyId;
     868            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId));
     869              : 
     870              :     u64 addr;
     871            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, addr));
     872              : 
     873            0 :     LocalIpc2RemoteAddr mapAddr;
     874            0 :     void *remoteAddr = reinterpret_cast<void *>(addr);
     875            0 :     CHK_PRT_RET(addressMgr_->GetLocalIpc2RemoteAddr(devicePhyId, remoteAddr, mapAddr) != HCCL_SUCCESS,
     876              :         HCCL_ERROR("[ZeroCopyMemoryAgent][ParseUnsetMemoryRange] device[%u] not set addr [%p]", devicePhyId, remoteAddr), HCCL_E_PARA);
     877            0 :     CHK_RET(addressMgr_->DelLocalIpc2RemoteAddr(devicePhyId, reinterpret_cast<void *>(mapAddr.remoteAddr)));
     878              : 
     879            0 :     void *devPtr = reinterpret_cast<void *>(mapAddr.localIpcAddr);
     880            0 :     aclError ret = aclrtReleaseMemAddress(devPtr);
     881            0 :     CHK_PRT_RET(ret != ACL_SUCCESS, HCCL_ERROR("[ZeroCopyMemoryAgent][ParseUnsetMemoryRange]rtRelease Memory failed, "\
     882              :         "return[%d], devPtr[%p]", ret, devPtr), HCCL_E_RUNTIME);
     883              : 
     884            0 :     CHK_RET(SendAckAfterParse(RequestType::UNSET_MEMORY_RANGE, RequestType::UNSET_MEMORY_RANGE_ACK, devicePhyId));
     885            0 :     return HCCL_SUCCESS;
     886              : }
     887              : 
     888            0 : HcclResult ZeroCopyMemoryAgent::ParseBareTgid(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize)
     889              : {
     890              :     u32 devicePhyId;
     891            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId));
     892              : 
     893              :     // 获取本端的ack,然后通过ack返回给对端
     894            0 :     int32_t tgid = 0;
     895            0 :     aclError ret = aclrtDeviceGetBareTgid(&tgid);
     896            0 :     CHK_PRT_RET(ret != ACL_SUCCESS, HCCL_ERROR("[ZeroCopyMemoryAgent][ParseBareTgid] get tgid failed, ret[%d]", ret), HCCL_E_RUNTIME);
     897              : 
     898            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][ParseBareTgid] dev[%u] tgid[%d] to remoteDev[%u]", devicePhyId_, tgid, devicePhyId);
     899            0 :     CHK_RET(SendAckAfterParse(RequestType::SET_REMOTE_BARE_TGID, RequestType::SET_REMOTE_BARE_TGID_ACK, devicePhyId,
     900              :         &tgid, sizeof(tgid)));
     901            0 :     return HCCL_SUCCESS;
     902              : }
     903              : 
     904            0 : HcclResult ZeroCopyMemoryAgent::ParseBareTgidAck(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize)
     905              : {
     906              :     u32 devicePhyId;
     907            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId));
     908              : 
     909              :     u32 tgid;
     910            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, tgid));
     911              : 
     912            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][ParseBareTgidAck] recv dev[%u] tgid[%u]", devicePhyId, tgid);
     913            0 :     remotePids_.emplace_back(tgid);
     914            0 :     return HCCL_SUCCESS;
     915              : }
     916              : 
     917            0 : HcclResult ZeroCopyMemoryAgent::ParseBarrierCloseAck(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize)
     918              : {
     919              :     u32 devicePhyId;
     920            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId));
     921              : 
     922              :     u32 tgid;
     923            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, tgid));
     924              : 
     925            0 :     receivedBarrierCloseAck_.insert(devicePhyId);
     926            0 :     HCCL_RUN_INFO("[ZeroCopyMemoryAgent][ParseBarrierCloseAck] [%s] recv dev[%u] barrier close ack, so we stop this socket's recv",
     927              :         identifier_.c_str(), devicePhyId, tgid);
     928            0 :     return HCCL_SUCCESS;
     929              : }
     930              : 
     931            0 : HcclResult ZeroCopyMemoryAgent::ParseActivateCommMemory(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize)
     932              : {
     933            0 :     CHK_PRT_RET(!ZeroCopyMemoryAgent::IsAddressMgrInited(), HCCL_ERROR("[ZeroCopyMemoryAgent][%s]ZeroCopyMemoryAgent "
     934              :         "is not init.", __func__), HCCL_E_INTERNAL);
     935              :     u32 devicePhyId;
     936            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId));
     937              : 
     938              :     u64 addr;
     939            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, addr));
     940              : 
     941              :     size_t size;
     942            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, size));
     943              : 
     944              :     size_t offset;
     945            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, offset));
     946              : 
     947              :     size_t shareableHandle;
     948            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, shareableHandle));
     949              : 
     950              :     size_t flags;
     951            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, flags));
     952              : 
     953            0 :     LocalIpc2RemoteAddr mapAddr;
     954            0 :     void *remoteAddr = reinterpret_cast<void *>(addr);
     955            0 :     CHK_PRT_RET((addressMgr_->GetLocalIpc2RemoteAddr(devicePhyId, remoteAddr, mapAddr) != HCCL_SUCCESS),
     956              :         HCCL_ERROR("[ZeroCopyMemoryAgent][ParseActivateCommMemory] address may not be reserved in device[%u]", devicePhyId), HCCL_E_PARA);
     957              :     
     958            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][ParseActivateCommMemory] prepare import from dev[%u] shareableHandle[%llu]", devicePhyId, shareableHandle);
     959            0 :     u64 actualAddr = mapAddr.localIpcAddr + (addr - mapAddr.remoteAddr);
     960            0 :     void* devPtr = reinterpret_cast<void*>(actualAddr);
     961            0 :     CHK_PRT_RET(actualAddr + size > mapAddr.localIpcAddr + mapAddr.length,
     962              :         HCCL_ERROR("[ZeroCopyMemoryAgent][ParseActivateCommMemory] remote addr[0x%lx] size[%llu] exceed memory range", addr, size), HCCL_E_PARA);
     963            0 :     CHK_PRT_RET(addressMgr_->IsOverlapWithActivateAddr(devPtr, size),
     964              :         HCCL_ERROR("[ZeroCopyMemoryAgent][ParseActivateCommMemory] remote addr[0x%lx] size[%llu] devPtr[%p] is overlap",
     965              :         addr, size, devPtr), HCCL_E_PARA);
     966              : 
     967            0 :     aclError ret = ACL_SUCCESS;
     968            0 :     void* pHandle = nullptr;
     969            0 :     CHK_RET(addressMgr_->ActivateCommMemoryAddr(devPtr, size));
     970            0 :     ret = aclrtMemImportFromShareableHandle(shareableHandle, deviceLogicId_, &pHandle);
     971            0 :     CHK_PRT_RET(ret != ACL_SUCCESS, HCCL_ERROR("[ZeroCopyMemoryAgent][ParseActivateCommMemory] import shareableHandle[%llu] dev[%d] failed, ret[%d]",
     972              :         shareableHandle, deviceLogicId_, ret), HCCL_E_RUNTIME);
     973              : 
     974            0 :     ret = aclrtMapMem(devPtr, size, offset, pHandle, flags);
     975            0 :     CHK_PRT_RET(ret != ACL_SUCCESS, HCCL_ERROR("[ZeroCopyMemoryAgent][ParseActivateCommMemory] map dev[%p] size[%llu] offset[%llu] handle[%p]",
     976              :         " flag[%llu] failed, ret[%d]", devPtr, size, offset, pHandle, flags, ret), HCCL_E_RUNTIME);
     977              : 
     978            0 :     CHK_RET(addressMgr_->AddRemoteImportAddr(devPtr, pHandle));
     979              : 
     980            0 :     CHK_RET(SendAckAfterParse(RequestType::ACTIVATE_COMM_MEMORY, RequestType::ACTIVATE_COMM_MEMORY_ACK, devicePhyId));
     981              : 
     982            0 :     return HCCL_SUCCESS;
     983              : }
     984              : 
     985            0 : HcclResult ZeroCopyMemoryAgent::ParseDeactivateCommMemory(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize)
     986              : {
     987            0 :     CHK_PRT_RET(!ZeroCopyMemoryAgent::IsAddressMgrInited(), HCCL_ERROR("[ZeroCopyMemoryAgent][%s]ZeroCopyMemoryAgent "
     988              :         "is not init.", __func__), HCCL_E_INTERNAL);
     989              :     u32 devicePhyId;
     990            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId));
     991              : 
     992              :     u64 addr;
     993            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, addr));
     994              : 
     995            0 :     LocalIpc2RemoteAddr mapAddr;
     996            0 :     void *remoteAddr = reinterpret_cast<void *>(addr);
     997            0 :     CHK_PRT_RET((addressMgr_->GetLocalIpc2RemoteAddr(devicePhyId, remoteAddr, mapAddr) != HCCL_SUCCESS),
     998              :         HCCL_ERROR("[ZeroCopyMemoryAgent][ParseDeactivateCommMemory] address [%p] not be set in device[%u]",
     999              :         remoteAddr, devicePhyId), HCCL_E_PARA);
    1000              : 
    1001            0 :     u64 actualAddr = mapAddr.localIpcAddr + (addr - mapAddr.remoteAddr);
    1002            0 :     void* devPtr = reinterpret_cast<void*>(actualAddr);
    1003            0 :     CHK_RET(addressMgr_->DeactivateCommMemoryAddr(devPtr));
    1004              : 
    1005            0 :     void *handle = nullptr;
    1006            0 :     CHK_RET(addressMgr_->GetRemoteImportAddr(devPtr, handle));
    1007              : 
    1008            0 :     aclError ret = ACL_SUCCESS;
    1009            0 :     ret = aclrtUnmapMem(devPtr);
    1010            0 :     CHK_PRT_RET(ret != ACL_SUCCESS, HCCL_ERROR("[ZeroCopyMemoryAgent][ParseDeactivateCommMemory] aclrtUnmapMem dev[%p] failed, ret[%d]",
    1011              :         devPtr, ret), HCCL_E_RUNTIME);
    1012            0 :     ret = aclrtFreePhysical(handle);
    1013            0 :     CHK_PRT_RET(ret != ACL_SUCCESS, HCCL_ERROR("[ZeroCopyMemoryAgent][ParseDeactivateCommMemory] aclrtFreePhysical handle[%p] failed, ret[%d]",
    1014              :         handle, ret), HCCL_E_RUNTIME);
    1015              : 
    1016            0 :     CHK_RET(addressMgr_->DelRemoteImportAddr(devPtr));
    1017              : 
    1018            0 :     CHK_RET(SendAckAfterParse(RequestType::DEACTIVATE_COMM_MEMORY, RequestType::DEACTIVATE_COMM_MEMORY_ACK, devicePhyId));
    1019              : 
    1020            0 :     return HCCL_SUCCESS;
    1021              : }
    1022              : 
    1023            0 : HcclResult ZeroCopyMemoryAgent::ParseBarrierClose(u8* &exchangeDataPtr, u32 &exchangeDataBlankSize)
    1024              : {
    1025              :     u32 devicePhyId;
    1026            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, devicePhyId));
    1027            0 :     HCCL_INFO("[ZeroCopyMemoryAgent][ParseBarrierClose] recv dev[%u] barrier close", devicePhyId);
    1028              : 
    1029            0 :     receivedBarrierClose_.insert(devicePhyId);
    1030            0 :     CHK_RET(SendAckAfterParse(RequestType::BARRIER_CLOSE, RequestType::BARRIER_CLOSE_ACK, devicePhyId));
    1031            0 :     return HCCL_SUCCESS;
    1032              : }
    1033              : 
    1034            0 : HcclResult ZeroCopyMemoryAgent::ParseReceivedRequest(std::vector<u8>& receivedData, u32 remoteRank)
    1035              : {
    1036            0 :     u8* exchangeDataPtr = receivedData.data();
    1037            0 :     u32 exchangeDataBlankSize = IPC_MEMORY_EXCHANGE_LENGTH;
    1038              : 
    1039              :     RequestType requestType;
    1040            0 :     CHK_RET(ParseData(exchangeDataPtr, exchangeDataBlankSize, requestType));
    1041              : 
    1042            0 :     HcclResult ret = HCCL_SUCCESS;
    1043            0 :     switch (requestType) {
    1044            0 :         case RequestType::SET_MEMORY_RANGE:
    1045            0 :             ret = ParseSetMemoryRange(exchangeDataPtr, exchangeDataBlankSize);
    1046            0 :             break;
    1047            0 :         case RequestType::UNSET_MEMORY_RANGE:
    1048            0 :             ret = ParseUnsetMemoryRange(exchangeDataPtr, exchangeDataBlankSize);
    1049            0 :             break;
    1050            0 :         case RequestType::ACTIVATE_COMM_MEMORY:
    1051            0 :             ret = ParseActivateCommMemory(exchangeDataPtr, exchangeDataBlankSize);
    1052            0 :             break;
    1053            0 :         case RequestType::DEACTIVATE_COMM_MEMORY:
    1054            0 :             ret = ParseDeactivateCommMemory(exchangeDataPtr, exchangeDataBlankSize);
    1055            0 :             break;
    1056            0 :         case RequestType::SET_REMOTE_BARE_TGID:
    1057            0 :             ret = ParseBareTgid(exchangeDataPtr, exchangeDataBlankSize);
    1058            0 :             break;
    1059            0 :         case RequestType::BARRIER_CLOSE:
    1060            0 :             ret = ParseBarrierClose(exchangeDataPtr, exchangeDataBlankSize);
    1061            0 :             break;
    1062            0 :         case RequestType::SET_REMOTE_BARE_TGID_ACK:
    1063            0 :             ret = ParseBareTgidAck(exchangeDataPtr, exchangeDataBlankSize);
    1064            0 :             ParseRemoteAck(requestType, remoteRank);
    1065            0 :             break;
    1066            0 :         case RequestType::SET_MEMORY_RANGE_ACK:
    1067              :         case RequestType::UNSET_MEMORY_RANGE_ACK:
    1068              :         case RequestType::ACTIVATE_COMM_MEMORY_ACK:
    1069              :         case RequestType::DEACTIVATE_COMM_MEMORY_ACK:
    1070            0 :             ParseRemoteAck(requestType, remoteRank);
    1071            0 :             break;
    1072            0 :         case RequestType::BARRIER_CLOSE_ACK:
    1073            0 :             ret = ParseBarrierCloseAck(exchangeDataPtr, exchangeDataBlankSize);
    1074            0 :             ParseRemoteAck(requestType, remoteRank);
    1075            0 :             break;
    1076            0 :         default:
    1077            0 :             HCCL_ERROR("[Parse][ReceivedRequest] invalid RequestType[%d]", requestType);
    1078            0 :             ret = HCCL_E_INTERNAL;
    1079            0 :             break;
    1080              :     }
    1081            0 :     return ret;
    1082              : }
    1083              : 
    1084            0 : std::string ZeroCopyMemoryAgent::DumpFinishInfo(RequestType requestType)
    1085              : {
    1086            0 :     auto &finishedRanks = reqMsgFinishedRanks_[static_cast<int>(requestType)];
    1087              : 
    1088            0 :     std::string msg = "Expect [";
    1089            0 :     for (auto &info : rankInfoList_) {
    1090            0 :         msg += std::to_string(info.userRank) + " ";
    1091              :     }
    1092              : 
    1093            0 :     msg += "] Actual [";
    1094            0 :     for (auto &rank : finishedRanks) {
    1095            0 :         msg += std::to_string(rank) + " ";
    1096              :     }
    1097              : 
    1098            0 :     msg += "]";
    1099            0 :     finishedRanks.clear();
    1100              : 
    1101            0 :     return msg;
    1102            0 : }
    1103              : 
    1104            0 : bool ZeroCopyMemoryAgent::IsPaused() const
    1105              : {
    1106            0 :     return !threadRun_ || isPaused_;
    1107              : }
    1108              : 
    1109            0 : bool ZeroCopyMemoryAgent::IsResumed() const
    1110              : {
    1111            0 :     return !threadRun_ || !isPaused_;
    1112              : }
    1113              : 
    1114            0 : void ZeroCopyMemoryAgent::CheckSnapshotStatus()
    1115              : {
    1116            0 :     auto snapshotStatus = SnapshotControl::GetInstance(deviceLogicId_).GetStatus();
    1117            0 :     if (isPaused_ && snapshotStatus == SnapshotStatus::POST_SNAPSHOT) {
    1118            0 :         isPaused_ = false;
    1119            0 :         HCCL_RUN_INFO("[ZeroCopyMemoryAgent][CheckSnapshotStatus] detect snapshot post-processing, "
    1120              :             "zero-copy memory agent is resumed, deviceLogicId[%d].", deviceLogicId_);
    1121            0 :     } else if (!isPaused_ && snapshotStatus == SnapshotStatus::PRE_SNAPSHOT) {
    1122            0 :         isPaused_ = true;
    1123            0 :         HCCL_RUN_INFO("[ZeroCopyMemoryAgent][CheckSnapshotStatus] detect snapshot pre-processing, "
    1124              :             "zero-copy memory agent is paused, deviceLogicId[%d].", deviceLogicId_);
    1125              :     }
    1126            0 : }
    1127              : 
    1128              : }  // namespace hccl
        

Generated by: LCOV version 2.0-1