LCOV - code coverage report
Current view: top level - base_comm/resources/endpoint_pairs/channels/host - host_cpu_urma_channel.cc (source / functions) Coverage Total Hit
Test: coverage.info Lines: 57.2 % 236 135
Test Date: 2026-07-28 12:11:00 Functions: 62.5 % 24 15

            Line data    Source code
       1              : /**
       2              : * Copyright (c) 2026 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 "host_cpu_urma_channel.h"
      12              : #include "endpoint.h"
      13              : #include "orion_adpt_utils.h"
      14              : #include "hcomm_adapter_urma.h"
      15              : 
      16              : // Orion
      17              : #include "topo_common_types.h"
      18              : #include "virtual_topo.h"
      19              : #include "env_config/env_config.h"
      20              : 
      21              : namespace hcomm {
      22              : constexpr uint16_t DEFAULT_LISTENING_PORT = 60001;
      23              : 
      24           17 : HostCpuUrmaChannel::HostCpuUrmaChannel(EndpointHandle endpointHandle, const HcommChannelDesc &channelDesc):
      25           17 :     endpointHandle_(endpointHandle), channelDesc_(channelDesc) {}
      26              : 
      27           34 : HostCpuUrmaChannel::~HostCpuUrmaChannel()
      28              : {
      29           17 :     if (channelDesc_.socket == nullptr && socket_ != nullptr) {
      30            1 :         SocketMgr::GetInstance(devicePhyId_).PutSocket(socketConfig_, socket_);
      31            1 :         socket_ = nullptr;
      32              :     }
      33           34 : }
      34              : 
      35           10 : HcclResult HostCpuUrmaChannel::ParseInputParam()
      36              : {
      37              :     // 1. 从 endpointHandle_,获得 localEp_ 和 rdmaHandle_
      38           10 :     Endpoint* localEpPtr = reinterpret_cast<Endpoint*>(endpointHandle_);
      39           10 :     CHK_PTR_NULL(localEpPtr);
      40            9 :     localEp_ = localEpPtr->GetEndpointDesc();
      41            9 :     rdmaHandle_ = localEpPtr->GetRdmaHandle();
      42              : 
      43            9 :     HCCL_INFO("[HostCpuUrmaChannel][%s] localProtocol[%d]", __func__, localEp_.protocol);
      44              : 
      45              :     // 2. 从 channelDesc_,获得 remoteEp_, socket_ 和 notifyNum
      46            9 :     remoteEp_ = channelDesc_.remoteEndpoint;
      47            9 :     socket_ = reinterpret_cast<Hccl::Socket*>(channelDesc_.socket);
      48            9 :     commonRes_.bufferVec.clear();
      49              : 
      50            9 :     if (channelDesc_.exchangeAllMems) {
      51              :         // 3. Get memHandles from endpoint
      52            1 :         HCCL_INFO("[HostCpuUrmaChannel][%s] exchangeAllMems == True. Get memHandles from endpoint.", __func__);
      53            1 :         std::shared_ptr<Hccl::LocalUbRmaBuffer> *memHandles = nullptr;
      54            1 :         uint32_t memHandleNum = 0;
      55            1 :         CHK_RET(static_cast<HcclResult>(HcommMemGetAllMemHandles(
      56              :             endpointHandle_, reinterpret_cast<void**>(&memHandles), &memHandleNum)));
      57            1 :         HCCL_INFO("[HostCpuUrmaChannel][%s] Got memHandleNum[%u].", __func__, memHandleNum);
      58            2 :         for (uint32_t i = 0; i < memHandleNum; ++i) {
      59            1 :             std::shared_ptr<Hccl::LocalUbRmaBuffer> &localUbRmaBuffer = memHandles[i];
      60            1 :             HCCL_INFO("[HostCpuUrmaChannel][%s] Got memHandle No.%u: addr[0x%llx], size[0x%llx], memType[%d], memInfo[%s].",
      61              :                 __func__, i, static_cast<unsigned long long>(localUbRmaBuffer->GetAddr()),
      62              :                 static_cast<unsigned long long>(localUbRmaBuffer->GetSize()),
      63              :                 static_cast<int>(localUbRmaBuffer->GetBuf()->GetMemType()),
      64              :                 localUbRmaBuffer->GetBuf()->GetMemInfo().c_str());
      65            1 :             commonRes_.bufferVec.push_back(localUbRmaBuffer.get());
      66              :         }
      67              :     } else {
      68              :         // 3. 从 channelDesc 的 memHandle,获得 bufs_
      69            8 :         HCCL_WARNING("[HostCpuUrmaChannel][%s] exchangeAllMems is false.", __func__);
      70              :     }
      71              : 
      72            9 :     return HCCL_SUCCESS;
      73              : }
      74              : 
      75            1 : HcclResult HostCpuUrmaChannel::StartListen()
      76              : {
      77            1 :     uint16_t port = channelDesc_.port;
      78            1 :     HCCL_INFO("[HostCpuUrmaChannel::%s] Start. EndpointHandle[0x%llx], port[%u]", __func__, reinterpret_cast<uint64_t>(endpointHandle_), port);
      79            1 :     if (port == 0) {
      80            0 :         port = DEFAULT_LISTENING_PORT;
      81            0 :         HCCL_INFO("[HostCpuUrmaChannel::%s] channelDesc port is 0, use default port [%u]", __func__, port);
      82              :     }
      83            1 :     CHK_RET(static_cast<HcclResult>(HcommEndpointStartListen(endpointHandle_, port, nullptr)));
      84            0 :     HCCL_INFO("[HostCpuUrmaChannel::%s] SUCCESS. port[%u].", __func__, port);
      85            0 :     return HCCL_SUCCESS;
      86              : }
      87              : 
      88            8 : HcclResult HostCpuUrmaChannel::BuildSocket()
      89              : {
      90            8 :     if (socket_ != nullptr) {
      91            7 :         return HCCL_SUCCESS;
      92              :     }
      93            1 :     HCCL_INFO("[HostCpuUrmaChannel::%s] socket ptr is NULL, rebuild Socket", __func__);
      94              : 
      95            1 :     Hccl::LinkData linkData = BuildDefaultLinkData();
      96            1 :     CHK_RET(EndpointDescPairToLinkData(localEp_, remoteEp_, linkData));
      97            1 :     HCCL_INFO("[HostCpuUrmaChannel::%s] built linkData: %s", __func__, linkData.Describe().c_str());
      98            1 :     uint16_t port = channelDesc_.port;
      99            1 :     if (port == 0) {
     100            0 :         port = DEFAULT_LISTENING_PORT;
     101            0 :         HCCL_INFO("[HostCpuUrmaChannel::%s] channelDesc port is 0, use default port [%u]", __func__, port);
     102              :     }
     103              :     
     104            1 :     std::string socketTag = (channelDesc_.channelName != nullptr)
     105            3 :         ? std::string(channelDesc_.channelName) : "AUTOMATIC_SOCKET_TAG";
     106            1 :     Hccl::SocketConfig socketConfig = (channelDesc_.role != HCOMM_SOCKET_ROLE_RESERVED)
     107            1 :         ? Hccl::SocketConfig(linkData, port, socketTag, channelDesc_.role == HCOMM_SOCKET_ROLE_SERVER)
     108            1 :         : Hccl::SocketConfig(linkData, socketTag, true);
     109            1 :     CHK_RET(SocketMgr::GetInstance(devicePhyId_).GetSocket(socketConfig, socket_));
     110              : 
     111            1 :     HCCL_INFO("[HostCpuUrmaChannel::%s] SUCCESS. port[%u].", __func__, port);
     112            1 :     return HCCL_SUCCESS;
     113            1 : }
     114              : 
     115            8 : HcclResult HostCpuUrmaChannel::BuildConnection()
     116              : {
     117            8 :     UbConnBuildContext ctx;
     118            8 :     CHK_RET(PrepareUbConnBuildContext(localEp_, remoteEp_, channelDesc_.qos, ctx));
     119              : 
     120            8 :     Hccl::OpMode opMode = Hccl::OpMode::OPBASE;
     121            8 :     std::unique_ptr<Hccl::HostUbConnection> ubConn = nullptr;
     122            8 :     switch (ctx.protocol) {
     123            0 :         case Hccl::LinkProtocol::UB_TP:
     124            0 :             EXCEPTION_CATCH(
     125              :                 ubConn = std::make_unique<Hccl::HostUbTpConnection>(rdmaHandle_, ctx.locAddr, ctx.rmtAddr, opMode,
     126              :                     Hccl::HrtUbJfcMode::NORMAL, ctx.qosPre),
     127              :                 return HCCL_E_PTR
     128              :             );
     129            0 :             break;
     130            8 :         case Hccl::LinkProtocol::UB_CTP:
     131            8 :             EXCEPTION_CATCH(
     132              :                 ubConn = std::make_unique<Hccl::HostUbCtpConnection>(rdmaHandle_, ctx.locAddr, ctx.rmtAddr, opMode,
     133              :                     Hccl::HrtUbJfcMode::NORMAL, ctx.qosPre),
     134              :                 return HCCL_E_PTR
     135              :             );
     136            8 :             break;
     137            0 :         default:
     138            0 :             HCCL_ERROR("%s No LinkProtocol protocol[%s] to match", __func__, ctx.protocol.Describe().c_str());
     139            0 :             break;
     140              :     }
     141            8 :     CHK_SMART_PTR_NULL(ubConn);
     142              : 
     143            8 :     commonRes_.connVec.clear();
     144            8 :     commonRes_.connVec.emplace_back(ubConn.get());
     145            8 :     connections_.clear();
     146            8 :     connections_.push_back(std::move(ubConn));
     147              : 
     148            8 :     return HCCL_SUCCESS;
     149            8 : }
     150              : 
     151            8 : HcclResult HostCpuUrmaChannel::BuildUbMemTransport()
     152              : {
     153            8 :     Hccl::BaseMemTransport::LocCntNotifyRes locCntNotifyRes{};
     154            8 :     const Hccl::Socket &socket = *socket_;
     155            8 :     bool isRecvFirst = socket.GetRole() == Hccl::SocketRole::CLIENT ? true : false;
     156              : 
     157            8 :     Hccl::LinkData linkData = BuildDefaultLinkData();
     158            8 :     CHK_RET(EndpointDescPairToLinkData(localEp_, remoteEp_, linkData));
     159              : 
     160              :     // make_unique / make_shared / release 包一层抛异常的宏
     161            8 :     EXCEPTION_CATCH(
     162              :         memTransport_ = std::make_unique<Hccl::UbMemTransport>(
     163              :             commonRes_, attr_, linkData, socket, rdmaHandle_, locCntNotifyRes, isRecvFirst
     164              :         ),
     165              :         return HCCL_E_PTR
     166              :     );
     167            8 :     return HCCL_SUCCESS;
     168            8 : }
     169              : 
     170           10 : HcclResult HostCpuUrmaChannel::Init()
     171              : {
     172              :     s32 devLogicId;
     173           10 :     CHK_RET(hrtGetDevice(&devLogicId));
     174           10 :     CHK_RET(hrtGetDevicePhyIdByIndex(static_cast<u32>(devLogicId), devicePhyId_));
     175           10 :     CHK_RET(ParseInputParam());
     176            9 :     if (channelDesc_.role != HCOMM_SOCKET_ROLE_CLIENT) {
     177            1 :         CHK_RET(StartListen());
     178              :     }
     179            8 :     CHK_RET(BuildSocket());
     180            8 :     CHK_RET(BuildConnection());
     181            8 :     CHK_RET(BuildUbMemTransport());
     182              :     // urma函数初始化
     183            8 :     CHK_RET(DlUrmaFunction::GetInstance().DlUrmaFunctionInit());
     184              :     // 获取urma read/write 单个wr的最大传输数据大小
     185            8 :     CHK_RET(HccpRaGetDevBaseAttr(rdmaHandle_, &devBaseAttr_));
     186              : 
     187            8 :     return HCCL_SUCCESS;
     188              : }
     189              : 
     190            0 : HcclResult HostCpuUrmaChannel::GetNotifyNum(uint32_t *notifyNum) const
     191              : {
     192            0 :     HCCL_INFO("[HostCpuUrmaChannel::%s] not supported yet.", __func__);
     193            0 :     return HCCL_SUCCESS;
     194              : }
     195              : 
     196            1 : HcclResult HostCpuUrmaChannel::GetRemoteMems(uint32_t *memNum, CommMem **remoteMem, char ***memInfos)
     197              : {
     198            1 :     return memTransport_->GetRemoteMems(memNum, remoteMem, memInfos);
     199              : }
     200              : 
     201            0 : ChannelStatus HostCpuUrmaChannel::GetStatus()
     202              : {
     203            0 :     memTransport_->SetIsHost();
     204            0 :     ChannelStatus out = Channel::TransportStatusToChannelStatus(memTransport_->GetStatus());
     205            0 :     return out;
     206              : }
     207              : 
     208            0 : HcclResult hcomm::HostCpuUrmaChannel::NotifyRecord(const uint32_t remoteNotifyIdx)
     209              : {
     210            0 :     HCCL_INFO("[HostCpuUrmaChannel::%s] not supported yet.", __func__);
     211            0 :     return HCCL_E_NOT_SUPPORT;
     212              : }
     213              : 
     214            0 : HcclResult hcomm::HostCpuUrmaChannel::NotifyWait(const uint32_t localNotifyIdx, const uint32_t timeout)
     215              : {
     216            0 :     HCCL_INFO("[HostCpuUrmaChannel::%s] not supported yet.", __func__);
     217            0 :     return HCCL_E_NOT_SUPPORT;
     218              : }
     219              : 
     220            0 : HcclResult hcomm::HostCpuUrmaChannel::WriteWithNotify(void *dst, const void *src, const uint64_t len, uint32_t remoteNotifyIdx)
     221              : {
     222            0 :     HCCL_INFO("[HostCpuUrmaChannel::%s] not supported yet.", __func__);
     223            0 :     return HCCL_E_NOT_SUPPORT;
     224              : }
     225              : 
     226            4 : HcclResult HostCpuUrmaChannel::GetLocSeg(const void *addr, const size_t size, u64 *seg)
     227              : {
     228            4 :     if (commonRes_.bufferVec.empty()) {
     229            1 :         HCCL_ERROR("[HostCpuUrmaChannel::%s] commonRes_.bufferVec is empty.", __func__);
     230            1 :         return HCCL_E_INTERNAL;
     231              :     }
     232              : 
     233            3 :     bool isAddrInRange = false;
     234            4 :     for (auto &it : commonRes_.bufferVec) {
     235            4 :         CHK_PTR_NULL(it);
     236            3 :         Hccl::Buffer iterBuf(it->GetAddr(), it->GetSize());
     237            3 :         if (iterBuf.Contains(reinterpret_cast<uintptr_t>(addr), size)) {
     238            2 :             auto localUbRmaBuffer = dynamic_cast<Hccl::LocalUbRmaBuffer *>(it);
     239            2 :             CHK_PTR_NULL(localUbRmaBuffer);
     240            1 :             *seg = localUbRmaBuffer->GetTargetSeg();
     241            1 :             isAddrInRange = true;
     242            1 :             break;
     243              :         }
     244            3 :     }
     245              : 
     246            2 :     if (!isAddrInRange) {
     247            1 :         HCCL_ERROR("GetLocSeg addr[%p] size[%llu] is not in commonRes_.bufferVec", addr, size);
     248            1 :         return HCCL_E_INTERNAL;
     249              :     }
     250            1 :     return HCCL_SUCCESS;
     251              : }
     252              : 
     253            4 : HcclResult HostCpuUrmaChannel::GetSplitNum(uint64_t len, uint64_t maxJettyWrDataLen, uint64_t &splitNum)
     254              : {
     255            4 :     if (len == 0 || maxJettyWrDataLen == 0) {
     256            1 :         HCCL_ERROR("[HostCpuUrmaChannel::%s] invalid len[%llu] or maxJettyWrDataLen[%llu].", __func__, len, maxJettyWrDataLen);
     257            1 :         return HCCL_E_PARA;
     258              :     }
     259            3 :     if ((len % maxJettyWrDataLen) == 0) {
     260            2 :         splitNum = len / maxJettyWrDataLen;
     261              :     } else {
     262            1 :         splitNum = (len / maxJettyWrDataLen) + 1;
     263              :     }
     264            3 :     return HCCL_SUCCESS;
     265              : }
     266              : 
     267            0 : HcclResult HostCpuUrmaChannel::GetLocalAndRemoteSeg(urma_opcode_t opcode, void *dst, const void *src, uint64_t len, u64 &localSeg, u64 &remoteSeg)
     268              : {
     269            0 :     if (opcode == URMA_OPC_WRITE) {
     270            0 :         CHK_RET(GetLocSeg(src, len, &localSeg));
     271            0 :         CHK_RET(memTransport_->GetRemoteSeg(dst, len, &remoteSeg));
     272            0 :     } else if (opcode == URMA_OPC_READ) {
     273            0 :         CHK_RET(GetLocSeg(dst, len, &localSeg));
     274            0 :         CHK_RET(memTransport_->GetRemoteSeg(src, len, &remoteSeg));
     275              :     } 
     276            0 :     return HCCL_SUCCESS;
     277              : }
     278              : 
     279              : constexpr u32 RELAX_ORDER       = 1; // Relax Order
     280              : constexpr u32 STRONG_ORDER      = 2; // Strong Order
     281            0 : HcclResult HostCpuUrmaChannel::UrmaPostJettySendWr(urma_opcode_t opcode, void *dst, const void *src, uint64_t len)
     282              : {
     283              :     // 构造urma的wr
     284            0 :     urma_jfs_wr_t urmaWriteWr{};
     285            0 :     urmaWriteWr.opcode = opcode;
     286            0 :     urmaWriteWr.flag.bs.place_order = (fenceFlag_ == true ? STRONG_ORDER : RELAX_ORDER);
     287            0 :     urmaWriteWr.flag.bs.comp_order = 1;     // comp_order要一直保持为1,
     288            0 :     urmaWriteWr.flag.bs.fence = (fenceFlag_ == true ? 1 : 0);
     289            0 :     urmaWriteWr.flag.bs.complete_enable = 0;
     290            0 :     urmaWriteWr.flag.bs.inline_flag = 0;
     291            0 :     urmaWriteWr.tjetty = reinterpret_cast<urma_target_jetty_t*>(connections_[0]->GetTJettyVa());
     292            0 :     urmaWriteWr.user_ctx = 0; // 跟ibvs中的wr_id对应
     293            0 :     urmaWriteWr.next = nullptr;
     294              : 
     295              :     //  获取切片数量
     296            0 :     uint64_t splitNum = 0;
     297            0 :     uint64_t maxJettyWrDataLen = (opcode == URMA_OPC_WRITE) ? devBaseAttr_.maxWriteSize : devBaseAttr_.maxReadSize;
     298            0 :     CHK_RET(GetSplitNum(len, maxJettyWrDataLen, splitNum));
     299              : 
     300              :     u64 localSeg;
     301              :     u64 remoteSeg;
     302            0 :     CHK_RET(GetLocalAndRemoteSeg(opcode, dst, src, len, localSeg, remoteSeg));
     303              : 
     304            0 :     uint64_t offset = 0;
     305            0 :     for (uint64_t i = 0; i < splitNum; i++) {
     306            0 :         urma_jfs_wr_t *badWr = nullptr;
     307            0 :         uint64_t chunkLen = std::min(len - offset, maxJettyWrDataLen);
     308              :         // 源地址 数据长度 tseg
     309            0 :         urma_sge_t srclist = {0};
     310            0 :         urmaWriteWr.rw.src.sge = &srclist;
     311            0 :         urmaWriteWr.rw.src.sge->addr = reinterpret_cast<uint64_t>(static_cast<char *>(const_cast<void *>(src)) + offset);
     312            0 :         urmaWriteWr.rw.src.sge->len = chunkLen;
     313            0 :         urmaWriteWr.rw.src.sge->tseg = (opcode == URMA_OPC_WRITE) ? reinterpret_cast<urma_target_seg_t*>(localSeg) : reinterpret_cast<urma_target_seg_t*>(remoteSeg);
     314            0 :         urmaWriteWr.rw.src.num_sge = 1;
     315              : 
     316              :         // 目的地址 数据长度 tseg
     317            0 :         urma_sge_t dstlist = {0};
     318            0 :         urmaWriteWr.rw.dst.sge = &dstlist;
     319            0 :         urmaWriteWr.rw.dst.sge->addr = reinterpret_cast<uint64_t>(static_cast<const char *>(dst) + offset); // 远端地址
     320            0 :         urmaWriteWr.rw.dst.sge->len = chunkLen;
     321            0 :         urmaWriteWr.rw.dst.sge->tseg = (opcode == URMA_OPC_WRITE) ? reinterpret_cast<urma_target_seg_t*>(remoteSeg) : reinterpret_cast<urma_target_seg_t*>(localSeg);
     322            0 :         urmaWriteWr.rw.dst.num_sge = 1;
     323              : 
     324              :         // 只有最后一个wr上报cqe
     325            0 :         if (i == splitNum - 1) {
     326            0 :             urmaWriteWr.flag.bs.complete_enable = 1;
     327            0 :             urmaWriteWr.flag.bs.place_order = STRONG_ORDER; // 最后一个wr设置为strong order
     328              :         }
     329            0 :         CHK_RET(HrtUrmaPostJettySendWr(reinterpret_cast<urma_jetty_t*>(connections_[0]->GetJettyVa()), &urmaWriteWr, &badWr));
     330            0 :         offset += chunkLen;
     331              :     }
     332            0 :     fenceFlag_ = false;
     333            0 :     wqeNum_++;
     334            0 :     HCCL_INFO("UrmaPostJettySendWr opencode[%u] fenceFlag_[%u] wqeNum_[%u] splitNum[%llu] SUCCESS.", opcode, fenceFlag_, wqeNum_, splitNum);
     335            0 :     return HCCL_SUCCESS;
     336              : }
     337              : 
     338            0 : HcclResult hcomm::HostCpuUrmaChannel::Write(void *dst, const void *src, uint64_t len)
     339              : {
     340            0 :     CHK_RET(UrmaPostJettySendWr(URMA_OPC_WRITE, dst, src, len));
     341            0 :     return HCCL_SUCCESS;
     342              : }
     343              : 
     344            0 : HcclResult hcomm::HostCpuUrmaChannel::Read(void *dst, const void *src, uint64_t len)
     345              : {
     346            0 :     CHK_RET(UrmaPostJettySendWr(URMA_OPC_READ, dst, src, len));
     347            0 :     return HCCL_SUCCESS;
     348              : }
     349              : 
     350            2 : HcclResult hcomm::HostCpuUrmaChannel::ChannelFence()
     351              : {
     352            2 :     std::lock_guard<std::mutex> lock(fenceMutex_);
     353            2 :     HCCL_INFO("[HostCpuUrmaChannel::%s] start, wqeNum_ = %u va[%llu]", __func__, wqeNum_, connections_[0]->GetCqVa());
     354            2 :     CHK_PRT_RET(wqeNum_ == 0, HCCL_INFO("[HostCpuUrmaChannel::%s] no need to fence since no wqeNum[%u].", __func__), HCCL_SUCCESS);
     355            1 :     std::vector<urma_cr_t> wc(wqeNum_);
     356              : 
     357              :     auto timeout = std::chrono::milliseconds(
     358            1 :         static_cast<uint64_t>(Hccl::EnvConfig::GetInstance().GetRtsConfig().GetExecTimeOut()) * 1000ULL); // 乘1000转为毫秒
     359            1 :     auto startTime = std::chrono::steady_clock::now();
     360              :     while (true) {
     361            1 :         auto actualNum = HrtUrmaPollJfc(reinterpret_cast<urma_jfc_t*>(connections_[0]->GetCqVa()), wqeNum_, wc.data());
     362            1 :         if (actualNum < 0) {
     363            1 :             HCCL_ERROR("[HostCpuUrmaChannel::%s] urma_poll_jfc failed. actualNum=%d", __func__, actualNum);
     364            1 :             return HCCL_E_NETWORK;
     365              :         }
     366              : 
     367            0 :         uint32_t actualNum32 = static_cast<uint32_t>(actualNum);
     368            0 :         if (actualNum32 > wqeNum_) {
     369            0 :             HCCL_ERROR("[HostCpuUrmaChannel::%s] urma_poll_jfc polled more completions (%u) than expected (%u).",
     370              :                 __func__, actualNum32, wqeNum_);
     371            0 :             return HCCL_E_INTERNAL;
     372            0 :         } else if (actualNum32 > 0) {
     373            0 :             for (uint32_t i = 0; i < actualNum32; i++) {
     374            0 :                 if (wc[i].status != URMA_CR_SUCCESS) {
     375            0 :                     HCCL_ERROR("[HostCpuUrmaChannel::%s] urma_poll_jfc error. wc[%u] status:%d", __func__, i, wc[i].status);
     376            0 :                     return HCCL_E_NETWORK;
     377              :                 }
     378              :             }
     379            0 :             wqeNum_ -= actualNum32; // 减去已完成的数量,继续等待剩余的完成
     380            0 :             if (wqeNum_ == 0) {
     381            0 :                 break; // 所有的wqe都已完成,退出循环
     382              :             }
     383              :         }
     384              : 
     385            0 :         if ((std::chrono::steady_clock::now() - startTime) >= timeout) {
     386            0 :             HCCL_ERROR("[HostCpuUrmaChannel::%s] call urma_poll_jfc timeout.", __func__);
     387            0 :             return HCCL_E_TIMEOUT;
     388              :         }
     389            0 :     }
     390              : 
     391            0 :     wqeNum_ = 0; // 所有wqe都已完成,重置计算器
     392            0 :     fenceFlag_ = true;
     393            0 :     return HCCL_SUCCESS;
     394            2 : }
     395              : 
     396            1 : HcclResult hcomm::HostCpuUrmaChannel::Clean()
     397              : {
     398            1 :     HCCL_INFO("[HostCpuUrmaChannel::%s] not supported yet.", __func__);
     399            1 :     return HCCL_E_NOT_SUPPORT;
     400              : }
     401              : 
     402            1 : HcclResult hcomm::HostCpuUrmaChannel::Resume()
     403              : {
     404            1 :     HCCL_INFO("[HostCpuUrmaChannel::%s] not supported yet.", __func__);
     405            1 :     return HCCL_E_NOT_SUPPORT;
     406              : }
     407              : 
     408              : } // namespace hcomm
        

Generated by: LCOV version 2.0-1