LCOV - code coverage report
Current view: top level - legacy/ascend910/platform/resource/transport/host - transport_p2p.cc (source / functions) Coverage Total Hit
Test: coverage.info Lines: 20.2 % 1013 205
Test Date: 2026-08-18 17:47:01 Functions: 28.9 % 76 22

            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 "transport_p2p.h"
      12              : #include <securec.h>
      13              : #include <sys/socket.h>
      14              : #include <sys/types.h>
      15              : #include <arpa/inet.h>
      16              : #include <unistd.h>
      17              : 
      18              : #include "mem_name_repository_pub.h"
      19              : #include "adapter_rts.h"
      20              : #include "mem_host_pub.h"
      21              : 
      22              : namespace hccl {
      23              : std::array<DeviceMem, MAX_MODULE_DEVICE_NUM> TransportP2p::notifyValueMem_;
      24              : std::array<std::mutex, MAX_MODULE_DEVICE_NUM> TransportP2p::notifyValueMutex_;
      25              : std::array<Referenced, MAX_MODULE_DEVICE_NUM> TransportP2p::instanceRef_;
      26            4 : TransportP2p::TransportP2p(
      27              :     DispatcherPub* dispatcher, const std::unique_ptr<NotifyPool>& notifyPool, MachinePara& machinePara,
      28            4 :     std::chrono::milliseconds timeout)
      29              :     : TransportBase(dispatcher, notifyPool, machinePara, timeout),
      30            4 :       remoteInputPtr_(nullptr),
      31            4 :       remoteOutputPtr_(nullptr),
      32            4 :       remoteOutputOffsetValue_(0),
      33            4 :       remoteInputOffsetValue_(0),
      34            4 :       remoteOutputMemName_(),
      35            8 :       remoteInputMemName_()
      36              : {
      37            4 :     if (machinePara_.deviceLogicId >= 0 && (static_cast<u32>(machinePara_.deviceLogicId) < MAX_MODULE_DEVICE_NUM)) {
      38            4 :         instanceRef_[machinePara_.deviceLogicId].Ref();
      39              :     }
      40            4 :     userLocalNotify_.resize(notifyNum_);
      41            4 :     userRemoteNotify_.resize(notifyNum_);
      42            4 :     userRemoteNotifyAddr_.resize(notifyNum_);
      43            4 :     userRemoteNotifyOffset_.resize(notifyNum_);
      44            4 :     remoteIpcMemPtrVector_.resize(machinePara.mem.size());
      45            4 :     remoteIpcMemOffsetValueVector_.resize(machinePara.mem.size());
      46            4 :     remoteIpcMemSizeVector_.resize(machinePara.mem.size());
      47            4 :     remoteIpcMemNameVector_.resize(machinePara.mem.size());
      48            4 : }
      49              : 
      50            6 : TransportP2p::~TransportP2p()
      51              : {
      52            4 :     HCCL_DEBUG("~TransportP2p Enter!");
      53              : 
      54              :     // 关闭rtIpcOpenMemory打开的对端共享内存和内存名称映射
      55            4 :     if (!isMemInclude_) {
      56              :         MemNameRepository::GetInstance(machinePara_.deviceLogicId)
      57            4 :             ->CloseIpcMem(static_cast<const u8*>(remoteOutputMemName_.ipcName));
      58            4 :         HCCL_DEBUG("remoteOutputMemName_.ipcName[%d]", remoteOutputMemName_.ipcName);
      59              :         MemNameRepository::GetInstance(machinePara_.deviceLogicId)
      60            4 :             ->CloseIpcMem(static_cast<const u8*>(remoteInputMemName_.ipcName));
      61            4 :         HCCL_DEBUG("remoteInputMemName_.ipcName[%d]", remoteInputMemName_.ipcName);
      62              :     }
      63            4 :     for (u32 i = 0; i < machinePara_.mem.size(); i++) {
      64              :         MemNameRepository::GetInstance(machinePara_.deviceLogicId)
      65            0 :             ->CloseIpcMem(static_cast<const u8*>(remoteIpcMemNameVector_[i].ipcName));
      66            0 :         HCCL_DEBUG("remoteIpcMemNameVector_[%u].ipcName[%s]", i, remoteIpcMemNameVector_[i].ipcName);
      67              :     }
      68              : 
      69              :     // 关闭rtIpcSetMemoryName 设置的内存名
      70            4 :     if (!isMemInclude_) {
      71              :         MemNameRepository::GetInstance(machinePara_.deviceLogicId)
      72            4 :             ->DestroyIpcMem(machinePara_.outputMem.ptr(), machinePara_.outputMem.size(), isSioToHccs_);
      73            4 :         HCCL_DEBUG(
      74              :             "machinePara_.outputMem addr:[%p], size:[%llu]", machinePara_.outputMem.ptr(),
      75              :             machinePara_.outputMem.size());
      76              :         MemNameRepository::GetInstance(machinePara_.deviceLogicId)
      77            4 :             ->DestroyIpcMem(machinePara_.inputMem.ptr(), machinePara_.inputMem.size(), isSioToHccs_);
      78            4 :         HCCL_DEBUG(
      79              :             "machinePara_.inputMem addr:[%p], size:[%llu]", machinePara_.inputMem.ptr(), machinePara_.inputMem.size());
      80              :     }
      81            4 :     for (u32 i = 0; i < machinePara_.mem.size(); i++) {
      82              :         MemNameRepository::GetInstance(machinePara_.deviceLogicId)
      83            0 :             ->DestroyIpcMem(machinePara_.mem[i].ptr(), machinePara_.mem[i].size(), isSioToHccs_);
      84            0 :         HCCL_DEBUG(
      85              :             "machinePara_.mem[%u] addr:[%p], size:[%llu]", machinePara_.mem[i].ptr(), machinePara_.mem[i].size());
      86              :     }
      87              : 
      88            4 :     SignalDestroy();
      89              : 
      90            4 :     if (machinePara_.deviceLogicId >= 0 && (static_cast<u32>(machinePara_.deviceLogicId) < MAX_MODULE_DEVICE_NUM)) {
      91            4 :         if (instanceRef_[machinePara_.deviceLogicId].Unref() == 0) {
      92            4 :             std::unique_lock<std::mutex> lock(notifyValueMutex_[machinePara_.deviceLogicId]);
      93            4 :             notifyValueMem_[machinePara_.deviceLogicId].free();
      94            4 :         }
      95              :     }
      96            4 :     HCCL_DEBUG("~TransportP2p Success!");
      97            6 : }
      98              : 
      99            2 : HcclResult TransportP2p::Init()
     100              : {
     101            2 :     HCCL_INFO(
     102              :         "machineType=[%d], serverId=[%s], localDeviceId=[%d], remoteDeviceId=[%d], "
     103              :         "localRank=[%u], localUserRank=[%u], remoteRank=[%u], remoteUserRank=[%u], "
     104              :         "deviceType=[%d], input_ptr=[%p], output_ptr=[%p], linkAttribute=[0x%x], linkMode=[%d], "
     105              :         "notifyNum[%u], isIndOp[%d], custom exchange data size [%llu], specifyLink[%d].",
     106              :         machinePara_.machineType, machinePara_.serverId.c_str(), machinePara_.localDeviceId,
     107              :         machinePara_.remoteDeviceId, machinePara_.localUserrank, machinePara_.localWorldRank,
     108              :         machinePara_.remoteUserrank, machinePara_.remoteWorldRank, machinePara_.deviceType, machinePara_.inputMem.ptr(),
     109              :         machinePara_.outputMem.ptr(), machinePara_.linkAttribute, machinePara_.linkMode, machinePara_.notifyNum,
     110              :         machinePara_.isIndOp, machinePara_.exchangeInfo.size(), machinePara_.specifyLink);
     111            2 :     HcclUs startut = TIME_NOW();
     112              : 
     113              :     /* make input memory shared interprocess and assigned a name */
     114            2 :     if (!machinePara_.isNewOneSide) {
     115            0 :         CHK_SMART_PTR_NULL(machinePara_.inputMem);
     116            0 :         CHK_SMART_PTR_NULL(machinePara_.outputMem);
     117              :     }
     118              : 
     119            2 :     CHK_PTR_NULL(dispatcher_);
     120            2 :     CHK_SMART_PTR_NULL(notifyPool_);
     121            2 :     CHK_RET(CheckDeviceId());
     122            2 :     CHK_RET(CheckExchangeData());
     123            2 :     SetMemIncludeFlag();
     124              :     // 上层初始化时保证 machinePara_.sockets 非空
     125            2 :     if (machinePara_.sockets.size() == 0) {
     126            0 :         HCCL_ERROR("machinePara sockets is empty.");
     127            0 :         return HCCL_E_INTERNAL;
     128              :     }
     129            2 :     defaultSocket_ = machinePara_.sockets[0];
     130            2 :     CHK_PTR_NULL(defaultSocket_);
     131              : 
     132            2 :     CHK_RET(CheckLinkMode());
     133              : 
     134              :     /* 本端与远端交换tgid 信息 */
     135            2 :     CHK_RET(ExchangeTgidMesg()); // tgid 无法合并交换,因为依赖对端的tgid判定是同一个进程还是跨进程
     136              : 
     137            2 :     CHK_RET(SetLinkType()); // 需要在交换sdid之后调用,确定是否超节点内节点间HCCS场景
     138              : 
     139            2 :     CHK_RET(FillExchangeDataTotalSize());
     140              : 
     141            2 :     CHK_RET(ConstructExchangeForSend());
     142              : 
     143            2 :     HcclResult ret = defaultSocket_->Send(exchangeDataForSend_.data(), exchangeDataTotalSize_);
     144            2 :     CHK_PRT_RET(
     145              :         ret != HCCL_SUCCESS,
     146              :         HCCL_ERROR(
     147              :             "[TransportP2p][Init] failed to send exchangeData exchangeDataTotalSize[%llu], custom exchange data "
     148              :             "size [%llu].",
     149              :             exchangeDataTotalSize_, machinePara_.exchangeInfo.size()),
     150              :         ret);
     151              : 
     152            2 :     exchangeDataForRecv_.resize(exchangeDataTotalSize_);
     153            2 :     ret = defaultSocket_->Recv(exchangeDataForRecv_.data(), exchangeDataTotalSize_);
     154            2 :     CHK_PRT_RET(
     155              :         ret != HCCL_SUCCESS,
     156              :         HCCL_ERROR(
     157              :             "[TransportP2p][Init] failed to recv exchangeData exchangeDataTotalSize[%llu], custom exchange data "
     158              :             "size [%llu].",
     159              :             exchangeDataTotalSize_, machinePara_.exchangeInfo.size()),
     160              :         ret);
     161              : 
     162            2 :     HCCL_DEBUG("[TransportP2p][Init] Socket Data Received");
     163              : 
     164            2 :     CHK_RET(ParseReceivedExchangeData());
     165              : 
     166            2 :     SetTransportRelationship();
     167            2 :     SetUseSdmaToSignalRecord();
     168            2 :     CHK_RET(CreateNotifyValueBuffer());
     169              : 
     170            2 :     HcclUs endut = TIME_NOW();
     171            2 :     HCCL_INFO("Time:%lld us", DURATION_US(endut - startut));
     172              : 
     173            2 :     HCCL_USER_CRITICAL_LOG(
     174              :         "create hccl transport:communicator[%s], local rank[%u], remote rank[%u], "
     175              :         "transporttype[%s]",
     176              :         machinePara_.tag.c_str(), machinePara_.localUserrank, machinePara_.remoteUserrank,
     177              :         GetLinkTypeEnumStr(GetLinkType()).c_str());
     178              : 
     179            2 :     return HCCL_SUCCESS;
     180              : }
     181              : 
     182            4 : void TransportP2p::SetUseSdmaToSignalRecord()
     183              : {
     184              :     // AICPU展开时,在节点间使用SDMA进行notify record操作,STARS可检出节点间链路异常,触发HCCL重执行
     185              :     useSdmaToSignalRecord_
     186            8 :         = ((transportAttr_.relationship & HCCL_TRANSPORT_RELATIONSHIP_SAME_SERVER) == 0)
     187            4 :           && ((transportAttr_.linkType == LinkType::LINK_HCCS_SW) || (transportAttr_.linkType == LinkType::LINK_HCCS));
     188            4 : }
     189              : 
     190            2 : HcclResult TransportP2p::ParseSpecifyLink(LinkTypeInServer& linkType)
     191              : {
     192            2 :     if (machinePara_.specifyLink == LinkTypeInServer::RESERVED_LINK_TYPE || machinePara_.specifyLink == linkType) {
     193            2 :         return HCCL_SUCCESS; // 未指定切换链路,保持默认
     194            0 :     } else if (machinePara_.specifyLink == LinkTypeInServer::HCCS_SW_TYPE && linkType == LinkTypeInServer::SIO_TYPE) {
     195              :         // 切换链路基于ipc实现, 多线程场景暂不支持
     196            0 :         s32 sendPid = 0;
     197            0 :         CHK_RET(SalGetBareTgid(&sendPid));
     198            0 :         CHK_PRT_RET(
     199              :             sendPid == recvPid_, HCCL_WARNING("%s specifyLink is not support in multi-thread", __func__), HCCL_SUCCESS);
     200              : 
     201              :         // A3 DIE间通信场景, 将链路从SIO切换到HCCS
     202            0 :         linkType = LinkTypeInServer::HCCS_SW_TYPE;
     203            0 :         isSioToHccs_ = true;
     204            0 :         HCCL_INFO("%s specifyLink change to HCCS_SW_TYPE", __func__);
     205            0 :     } else {
     206            0 :         HCCL_ERROR("%s fail, linkType:%d, specifyLink:%d is not support", __func__, linkType, machinePara_.specifyLink);
     207            0 :         return HCCL_E_NOT_SUPPORT;
     208              :     }
     209            0 :     return HCCL_SUCCESS;
     210              : }
     211              : 
     212            2 : HcclResult TransportP2p::SetLinkType()
     213              : {
     214              :     // 计算linkType
     215            2 :     LinkTypeInServer linkType = LinkTypeInServer::HCCS_TYPE;
     216            2 :     if (recvSdid_ != INVALID_INT) { // 超节点内节点间走p2p通信时,链路类型为LINK_HCCS_SW
     217            0 :         linkType = LinkTypeInServer::HCCS_SW_TYPE;
     218              :     } else {
     219            2 :         CHK_RET(hrtGetPairDeviceLinkType(
     220              :             static_cast<u32>(machinePara_.localDeviceId), static_cast<u32>(machinePara_.remoteDeviceId), linkType));
     221              :     }
     222              : 
     223            2 :     CHK_RET(ParseSpecifyLink(linkType));
     224              : 
     225            2 :     switch (linkType) {
     226            2 :         case LinkTypeInServer::HCCS_TYPE:
     227            2 :             transportAttr_.linkType = hccl::LinkType::LINK_HCCS;
     228            2 :             break;
     229            0 :         case LinkTypeInServer::HCCS_SW_TYPE:
     230            0 :             transportAttr_.linkType = hccl::LinkType::LINK_HCCS_SW;
     231            0 :             break;
     232            0 :         case LinkTypeInServer::SIO_TYPE:
     233            0 :             transportAttr_.linkType = hccl::LinkType::LINK_SIO;
     234            0 :             break;
     235            0 :         default:
     236            0 :             transportAttr_.linkType = hccl::LinkType::LINK_PCIE;
     237            0 :             break;
     238              :     }
     239              : 
     240            2 :     HCCL_DEBUG("[TransportP2p] transportattr linktype: 0x%x", transportAttr_.linkType);
     241            2 :     return HCCL_SUCCESS;
     242              : }
     243              : 
     244            2 : HcclResult TransportP2p::CreateNotifyValueBuffer()
     245              : {
     246            2 :     if (!useSdmaToSignalRecord_) {
     247            2 :         return HCCL_SUCCESS;
     248              :     }
     249              : 
     250            0 :     u32 notifySize = 0;
     251            0 :     CHK_RET(hrtGetNotifySize(notifySize));
     252            0 :     std::unique_lock<std::mutex> lock(notifyValueMutex_[machinePara_.deviceLogicId]);
     253            0 :     if (notifyValueMem_[machinePara_.deviceLogicId].ptr() == nullptr) {
     254            0 :         u64 notifyVaule = 1; // notify值写1表示record
     255            0 :         CHK_RET(DeviceMem::alloc(notifyValueMem_[machinePara_.deviceLogicId], notifyValueSize_));
     256            0 :         HCCL_DEBUG(
     257              :             "create notify value buffer[%p], size[%u]", notifyValueMem_[machinePara_.deviceLogicId].ptr(), notifySize);
     258              : 
     259            0 :         CHK_RET(hrtMemSyncCopy(
     260              :             notifyValueMem_[machinePara_.deviceLogicId].ptr(), notifyValueMem_[machinePara_.deviceLogicId].size(),
     261              :             &notifyVaule, notifySize, HcclRtMemcpyKind::HCCL_RT_MEMCPY_KIND_HOST_TO_DEVICE));
     262              :     }
     263            0 :     transportAttr_.signalRecordBuff.address = reinterpret_cast<u64>(notifyValueMem_[machinePara_.deviceLogicId].ptr());
     264            0 :     transportAttr_.signalRecordBuff.length = notifySize;
     265              : 
     266            0 :     HCCL_DEBUG(
     267              :         "[TransportP2p] transportattr signalRecordBuff.address[%p], signalRecordBuff.length[%llu]",
     268              :         transportAttr_.signalRecordBuff.address, transportAttr_.signalRecordBuff.length);
     269            0 :     return HCCL_SUCCESS;
     270            0 : }
     271              : 
     272            2 : void TransportP2p::SetTransportRelationship()
     273              : {
     274            2 :     if (transportAttr_.linkType == hccl::LinkType::LINK_SIO) {
     275              :         // 芯片内
     276            0 :         transportAttr_.relationship |= HCCL_TRANSPORT_RELATIONSHIP_SAME_CHIP;
     277            0 :         transportAttr_.relationship |= HCCL_TRANSPORT_RELATIONSHIP_SAME_SERVER;
     278            0 :         transportAttr_.relationship |= HCCL_TRANSPORT_RELATIONSHIP_SAME_SUPERPOD;
     279            2 :     } else if (recvSdid_ == INVALID_INT) {
     280              :         // 节点内
     281            2 :         transportAttr_.relationship |= HCCL_TRANSPORT_RELATIONSHIP_SAME_SERVER;
     282            2 :         transportAttr_.relationship |= HCCL_TRANSPORT_RELATIONSHIP_SAME_SUPERPOD;
     283              :     } else {
     284              :         // 节点间
     285            0 :         transportAttr_.relationship |= HCCL_TRANSPORT_RELATIONSHIP_SAME_SUPERPOD;
     286              :     }
     287              : 
     288            2 :     HCCL_DEBUG("[TransportP2p] transportattr relationship: 0x%x", transportAttr_.relationship);
     289            2 :     return;
     290              : }
     291              : 
     292            2 : HcclResult TransportP2p::FillExchangeDataTotalSize()
     293              : {
     294            2 :     exchangeDataTotalSize_ = 0;
     295            2 :     s32 sendPid = 0;
     296            2 :     CHK_RET(SalGetBareTgid(&sendPid));
     297            2 :     u64 ipcMemDataSize = 0;
     298            2 :     if (sendPid != recvPid_ || recvSdid_ != INVALID_INT) {
     299              :         // 输入输出内存
     300            0 :         HCCL_DEBUG("[TransportP2p][FillExchangeDataTotalSize] Inter Proc");
     301            0 :         ipcMemDataSize = HCCL_IPC_MEM_NAME_LEN + sizeof(u64) + sizeof(u64); // size + offset
     302            0 :         if (!isMemInclude_) {
     303            0 :             exchangeInfoSize_.ipcMenSize = ipcMemDataSize * (2 + machinePara_.mem.size());
     304              :         } else {
     305              :             // in和out包含在整块CCLbuf的时候,不需要传ipcName,但是size和offset不能少
     306            0 :             exchangeInfoSize_.ipcMenSize = ipcMemDataSize * machinePara_.mem.size() + 2 * (sizeof(u64) + sizeof(u64));
     307              :         }
     308              :     } else {
     309            2 :         HCCL_DEBUG("[TransportP2p][FillExchangeDataTotalSize] intra Proc");
     310            2 :         ipcMemDataSize = sizeof(u64) + sizeof(u64); // addr + length
     311              :         exchangeInfoSize_.ipcMenSize
     312            2 :             = ipcMemDataSize * (2 + machinePara_.mem.size()); // 2: input  & output + mem.size()
     313              :     }
     314              : 
     315            2 :     if (!machinePara_.isNewOneSide) {
     316              :         // notify 信息
     317            0 :         if (machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
     318            0 :             || machinePara_.machineType == MachineType::MACHINE_CLIENT_TYPE) {
     319            0 :             exchangeInfoSize_.notifySize = NOTIFY_INFO_LENGTH;
     320              :         }
     321            0 :         if (machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
     322            0 :             || machinePara_.machineType == MachineType::MACHINE_SERVER_TYPE) {
     323            0 :             exchangeInfoSize_.notifySize += NOTIFY_INFO_LENGTH;
     324              :         }
     325              :         // 3.新增notify资源
     326            0 :         exchangeInfoSize_.notifySize += NOTIFY_INFO_LENGTH * notifyNum_;
     327              :     }
     328              : 
     329              :     // 自定义信息
     330            2 :     exchangeInfoSize_.exDataSize = machinePara_.exchangeInfo.size();
     331              : 
     332              :     // 独立算子内存
     333            2 :     if (machinePara_.isIndOp) {
     334              :         // userDeviceMem数量\userDeviceMem\userHostMem数量\userHostMem
     335            0 :         const int kMemCountItems = 2;
     336              :         exchangeInfoSize_.indOpMemSize
     337            0 :             = ipcMemDataSize * (machinePara_.userDeviceMem.size() + machinePara_.userHostMem.size());
     338            0 :         exchangeInfoSize_.indOpMemSize += sizeof(u64) * kMemCountItems;
     339              :     }
     340              : 
     341            2 :     exchangeDataTotalSize_ = exchangeInfoSize_.ipcMenSize + exchangeInfoSize_.notifySize + exchangeInfoSize_.exDataSize
     342            2 :                              + exchangeInfoSize_.indOpMemSize + sizeof(ExchangeInfoSize);
     343            2 :     HCCL_INFO(
     344              :         "[TransportP2p][FillExchangeDataTotalSize] exchangeDataTotalSize[%llu] memSize[%d]", exchangeDataTotalSize_,
     345              :         machinePara_.mem.size());
     346            2 :     return HCCL_SUCCESS;
     347              : }
     348              : 
     349            2 : HcclResult TransportP2p::ConstructExchangeForSend()
     350              : {
     351            2 :     exchangeDataForSend_.resize(exchangeDataTotalSize_);
     352            2 :     u8* exchangeDataPtr = exchangeDataForSend_.data();
     353            2 :     u64 exchangeDataBlankSize = exchangeDataTotalSize_;
     354            2 :     CHK_RET(ConstructDataLenForSend(exchangeDataPtr, exchangeDataBlankSize));
     355            2 :     u64 blankSizeRecord = exchangeDataBlankSize;
     356              : 
     357            2 :     s32 sendPid = 0;
     358            2 :     CHK_RET(SalGetBareTgid(&sendPid));
     359            2 :     HCCL_DEBUG("%s sendPid %d, recvPid %d, recvSdid %d", __func__, sendPid, recvPid_, recvSdid_);
     360            2 :     if (sendPid != recvPid_ || recvSdid_ != INVALID_INT) { // 跨进程方式交换
     361              :         // 构造IPC内存地址交换数据结构
     362            0 :         for (auto ipcMem : machinePara_.mem) {
     363            0 :             CHK_RET(ConstructIpcMemInfoForSend(ipcMem.ptr(), ipcMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     364            0 :         }
     365            0 :         if (!isMemInclude_) {
     366            0 :             CHK_RET(ConstructIpcMemInfoForSend(
     367              :                 machinePara_.outputMem.ptr(), machinePara_.outputMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     368            0 :             CHK_RET(ConstructIpcMemInfoForSend(
     369              :                 machinePara_.inputMem.ptr(), machinePara_.inputMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     370              :         } else {
     371            0 :             CHK_RET(ConstructMemIncludeInfoForSend(exchangeDataPtr, exchangeDataBlankSize));
     372              :         }
     373            0 :     } else {
     374              :         // 构造进程内内存地址交换数据结构
     375            2 :         CHK_RET(ConstructIntraProcMemInfoForSend(
     376              :             machinePara_.outputMem.ptr(), machinePara_.outputMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     377            2 :         CHK_RET(ConstructIntraProcMemInfoForSend(
     378              :             machinePara_.inputMem.ptr(), machinePara_.inputMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     379            2 :         for (auto ipcMem : machinePara_.mem) {
     380            0 :             CHK_RET(
     381              :                 ConstructIntraProcMemInfoForSend(ipcMem.ptr(), ipcMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     382            0 :         }
     383              :     }
     384            2 :     CHK_RET(SumCheckSizeAndConsisten(
     385              :         ExInfoType::EX_IPCMEN_SIZE, exchangeInfoSize_.ipcMenSize, blankSizeRecord, exchangeDataBlankSize));
     386              : 
     387            2 :     CHK_RET(ConstructNotifyInfoForSend(exchangeDataPtr, exchangeDataBlankSize));
     388            2 :     CHK_RET(ConstructNotifyVectorInfoForSend(exchangeDataPtr, exchangeDataBlankSize)); // 新增notify资源的创建
     389            2 :     CHK_RET(SumCheckSizeAndConsisten(
     390              :         ExInfoType::EX_NOTIFY_SIZE, exchangeInfoSize_.notifySize, blankSizeRecord, exchangeDataBlankSize));
     391              : 
     392            2 :     CHK_RET(ConstructExchangeDataForSend(exchangeDataPtr, exchangeDataBlankSize));
     393            2 :     CHK_RET(SumCheckSizeAndConsisten(
     394              :         ExInfoType::EX_EXDATA_SIZE, exchangeInfoSize_.exDataSize, blankSizeRecord, exchangeDataBlankSize));
     395              : 
     396              :     // 独立算子内存资源,无需检查大小
     397            2 :     if (machinePara_.isIndOp) {
     398            0 :         if (sendPid != recvPid_ || recvSdid_ != INVALID_INT) { // 跨进程方式交换
     399            0 :             CHK_RET(ConstructNumInfoForSend(machinePara_.userDeviceMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     400            0 :             for (auto ipcMem : machinePara_.userDeviceMem) {
     401            0 :                 CHK_RET(
     402              :                     ConstructIpcMemInfoForSend(ipcMem.ptr(), ipcMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     403            0 :             }
     404            0 :             CHK_RET(ConstructNumInfoForSend(machinePara_.userHostMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     405            0 :             for (auto ipcMem : machinePara_.userHostMem) {
     406            0 :                 CHK_RET(
     407              :                     ConstructIpcMemInfoForSend(ipcMem.ptr(), ipcMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     408            0 :             }
     409            0 :         } else {
     410            0 :             CHK_RET(ConstructNumInfoForSend(machinePara_.userDeviceMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     411            0 :             for (auto ipcMem : machinePara_.userDeviceMem) {
     412            0 :                 CHK_RET(ConstructIntraProcMemInfoForSend(
     413              :                     ipcMem.ptr(), ipcMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     414            0 :             }
     415            0 :             CHK_RET(ConstructNumInfoForSend(machinePara_.userHostMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     416            0 :             for (auto ipcMem : machinePara_.userHostMem) {
     417            0 :                 CHK_RET(ConstructIntraProcMemInfoForSend(
     418              :                     ipcMem.ptr(), ipcMem.size(), exchangeDataPtr, exchangeDataBlankSize));
     419            0 :             }
     420              :         }
     421              :     }
     422            2 :     if (exchangeDataBlankSize != 0) {
     423            0 :         HCCL_ERROR(
     424              :             "[TransportP2p][ConstructExchangeForSend] failed to construct exchange Data "
     425              :             "exchangeDataBlankSize[%llu]",
     426              :             exchangeDataBlankSize);
     427            0 :         return HCCL_E_INTERNAL;
     428              :     }
     429            2 :     return HCCL_SUCCESS; // this function should not be called in normal process
     430              : }
     431              : 
     432              : // exchangeDataPtr对指针进行了引用,因为需要改变exchangeDataPtr的值
     433              : HcclResult
     434            0 : TransportP2p::ConstructIpcMemInfoForSend(void* ptr, u64 size, u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     435              : {
     436              :     HcclResult ret;
     437              :     u64 memOffset;
     438            0 :     SecIpcName_t memName;
     439              : 
     440            0 :     if (!machinePara_.isNewOneSide) {
     441              :         ret = MemNameRepository::GetInstance(machinePara_.deviceLogicId)
     442            0 :                   ->SetIpcMem(
     443            0 :                       ptr, size, memName.ipcName, HCCL_IPC_MEM_NAME_LEN, memOffset, recvPid_, recvSdid_, isSioToHccs_);
     444            0 :         CHK_PRT_RET(
     445              :             ret != HCCL_SUCCESS,
     446              :             HCCL_ERROR(
     447              :                 "[Send][IpcMemMesg]errNo[0x%016llx], In send ipc mesg, get para mem name failed. "
     448              :                 "mem addr[%p] local rank[%u]",
     449              :                 HCCL_ERROR_CODE(ret), machinePara_.outputMem.ptr(), machinePara_.localUserrank),
     450              :             ret);
     451              :     }
     452              : 
     453              :     // 设置ipc mem属性,指定通信链路从sio切换至hccs
     454            0 :     if (isSioToHccs_) {
     455            0 :         u32 ipcAttr = 1; // 0: SIO(默认), 1: HCCS
     456            0 :         CHK_RET(hrtIpcSetMemoryAttr(memName.ipcName, ACL_RT_IPC_MEM_ATTR_ACCESS_LINK, ipcAttr));
     457              :     }
     458              : 
     459            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, memName.ipcName, HCCL_IPC_MEM_NAME_LEN));
     460            0 :     exchangeDataPtr += HCCL_IPC_MEM_NAME_LEN;
     461            0 :     exchangeDataBlankSize -= HCCL_IPC_MEM_NAME_LEN;
     462            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &size, sizeof(u64)));
     463            0 :     exchangeDataPtr += sizeof(u64);
     464            0 :     exchangeDataBlankSize -= sizeof(u64);
     465            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &memOffset, sizeof(u64)));
     466            0 :     exchangeDataPtr += sizeof(u64);
     467            0 :     exchangeDataBlankSize -= sizeof(u64);
     468              : 
     469            0 :     return HCCL_SUCCESS;
     470            0 : }
     471              : 
     472              : HcclResult
     473            4 : TransportP2p::ConstructIntraProcMemInfoForSend(void* ptr, u64 size, u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     474              : {
     475            4 :     if (!machinePara_.isNewOneSide) {
     476            0 :         CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &ptr, sizeof(u64)));
     477              :     }
     478            4 :     exchangeDataPtr += sizeof(u64);
     479            4 :     exchangeDataBlankSize -= sizeof(u64);
     480            4 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &size, sizeof(u64)));
     481            4 :     exchangeDataPtr += sizeof(u64);
     482            4 :     exchangeDataBlankSize -= sizeof(u64);
     483              : 
     484            4 :     return HCCL_SUCCESS;
     485              : }
     486              : 
     487            0 : HcclResult TransportP2p::ConstructNumInfoForSend(u64 num, u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     488              : {
     489            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &num, sizeof(u64)));
     490            0 :     exchangeDataPtr += sizeof(u64);
     491            0 :     exchangeDataBlankSize -= sizeof(u64);
     492            0 :     return HCCL_SUCCESS;
     493              : }
     494              : 
     495            0 : HcclResult TransportP2p::ParseMemNumInfo(u64& memNum, u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     496              : {
     497            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(&memNum, sizeof(u64), exchangeDataPtr, sizeof(u64)));
     498            0 :     exchangeDataPtr += sizeof(u64);
     499            0 :     exchangeDataBlankSize -= sizeof(u64);
     500            0 :     return HCCL_SUCCESS;
     501              : }
     502              : 
     503            0 : HcclResult TransportP2p::ParseIpcMemInfo(
     504              :     void** memPtr, u64& size, u8* memName, u64& offset, u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     505              : {
     506            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(memName, HCCL_IPC_MEM_NAME_LEN, exchangeDataPtr, HCCL_IPC_MEM_NAME_LEN));
     507            0 :     exchangeDataPtr += HCCL_IPC_MEM_NAME_LEN;
     508            0 :     exchangeDataBlankSize -= HCCL_IPC_MEM_NAME_LEN;
     509              : 
     510            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(&size, sizeof(u64), exchangeDataPtr, sizeof(u64)));
     511            0 :     exchangeDataPtr += sizeof(u64);
     512            0 :     exchangeDataBlankSize -= sizeof(u64);
     513              : 
     514            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(&offset, sizeof(u64), exchangeDataPtr, sizeof(u64)));
     515            0 :     exchangeDataPtr += sizeof(u64);
     516            0 :     exchangeDataBlankSize -= sizeof(u64);
     517              : 
     518            0 :     if (!machinePara_.isNewOneSide) {
     519              :         /* 根据名字,获取对端IPC 内存 */
     520            0 :         HcclResult ret = WaitPeerMemConfig(memPtr, const_cast<u8*>(memName), size, offset);
     521            0 :         CHK_PRT_RET(
     522              :             ret != HCCL_SUCCESS,
     523              :             HCCL_ERROR(
     524              :                 "[Recv][IpcMemMesg]errNo[0x%016llx]In recv ipc mem mesg, wait peer mem config "
     525              :                 "failed. local rank[%u]",
     526              :                 HCCL_ERROR_CODE(ret), machinePara_.localUserrank),
     527              :             ret);
     528              : 
     529            0 :         CHK_PTR_NULL(*memPtr);
     530              :     }
     531              : 
     532            0 :     HCCL_DEBUG(
     533              :         "localUserrank[%u] receive from remoteUserrank[%u]", machinePara_.localUserrank, machinePara_.remoteUserrank);
     534              : 
     535            0 :     return HCCL_SUCCESS;
     536              : }
     537              : 
     538            4 : HcclResult TransportP2p::ParseIntraProcMemInfo(u64* addr, u64* size, u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     539              : {
     540            4 :     if (!machinePara_.isNewOneSide) {
     541            0 :         CHK_SAFETY_FUNC_RET(memcpy_s(addr, sizeof(u64), exchangeDataPtr, sizeof(u64)));
     542            0 :         CHK_PTR_NULL(reinterpret_cast<void*>(*addr));
     543              :     }
     544              : 
     545            4 :     exchangeDataPtr += sizeof(u64);
     546            4 :     exchangeDataBlankSize -= sizeof(u64);
     547            4 :     CHK_SAFETY_FUNC_RET(memcpy_s(size, sizeof(u64), exchangeDataPtr, sizeof(u64)));
     548            4 :     exchangeDataPtr += sizeof(u64);
     549            4 :     exchangeDataBlankSize -= sizeof(u64);
     550            4 :     return HCCL_SUCCESS;
     551              : }
     552              : 
     553            0 : HcclResult TransportP2p::ParseNotifyInfo(u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     554              : {
     555            0 :     s32 sendPid = 0;
     556            0 :     CHK_RET(SalGetBareTgid(&sendPid)); // 当前进程id
     557            0 :     HCCL_INFO("LinkRecvNotifyMesg, sendPid[%d], recvPid[%d]", sendPid, recvPid_);
     558              : 
     559            0 :     if (machinePara_.isAicpuModeEn) {
     560            0 :         if ((machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
     561            0 :              || machinePara_.machineType == MachineType::MACHINE_SERVER_TYPE)
     562            0 :             && machinePara_.isAicpuModeEn == true) {
     563            0 :             std::vector<u8> data(NOTIFY_INFO_LENGTH, 0);
     564            0 :             CHK_SAFETY_FUNC_RET(memcpy_s(&data[0], data.size(), exchangeDataPtr, NOTIFY_INFO_LENGTH));
     565            0 :             exchangeDataPtr += NOTIFY_INFO_LENGTH;
     566            0 :             exchangeDataBlankSize -= NOTIFY_INFO_LENGTH;
     567            0 :             CHK_RET(OpenRemoteNotify(data, remoteSendReadyDeviceNotify_));
     568            0 :         }
     569              : 
     570            0 :         if ((machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
     571            0 :              || machinePara_.machineType == MachineType::MACHINE_CLIENT_TYPE)
     572            0 :             && machinePara_.isAicpuModeEn == true) {
     573            0 :             std::vector<u8> data(NOTIFY_INFO_LENGTH, 0);
     574            0 :             CHK_SAFETY_FUNC_RET(memcpy_s(&data[0], data.size(), exchangeDataPtr, NOTIFY_INFO_LENGTH));
     575            0 :             exchangeDataPtr += NOTIFY_INFO_LENGTH;
     576            0 :             exchangeDataBlankSize -= NOTIFY_INFO_LENGTH;
     577            0 :             CHK_RET(OpenRemoteNotify(data, remoteSendDoneDeviceNotify_));
     578            0 :         }
     579              :     } else {
     580            0 :         if (machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
     581            0 :             || machinePara_.machineType == MachineType::MACHINE_SERVER_TYPE) {
     582            0 :             std::vector<u8> data(NOTIFY_INFO_LENGTH, 0);
     583            0 :             CHK_SAFETY_FUNC_RET(memcpy_s(&data[0], data.size(), exchangeDataPtr, NOTIFY_INFO_LENGTH));
     584            0 :             exchangeDataPtr += NOTIFY_INFO_LENGTH;
     585            0 :             exchangeDataBlankSize -= NOTIFY_INFO_LENGTH;
     586            0 :             CHK_RET(OpenRemoteNotify(data, remoteSendReadyNotify_));
     587              : 
     588              :             HcclSignalInfo notifyInfo;
     589            0 :             CHK_RET(remoteSendReadyNotify_->GetNotifyData(notifyInfo));
     590            0 :             CHK_RET(remoteSendReadyNotify_->GetNotifyOffset(remoteSendReadyOffset_));
     591              : 
     592            0 :             remoteSendReadyAddress_ = notifyInfo.addr;
     593            0 :         }
     594              : 
     595            0 :         if (machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
     596            0 :             || machinePara_.machineType == MachineType::MACHINE_CLIENT_TYPE) {
     597            0 :             std::vector<u8> data(NOTIFY_INFO_LENGTH, 0);
     598            0 :             CHK_SAFETY_FUNC_RET(memcpy_s(&data[0], data.size(), exchangeDataPtr, NOTIFY_INFO_LENGTH));
     599            0 :             exchangeDataPtr += NOTIFY_INFO_LENGTH;
     600            0 :             exchangeDataBlankSize -= NOTIFY_INFO_LENGTH;
     601            0 :             CHK_RET(OpenRemoteNotify(data, remoteSendDoneNotify_));
     602              :             HcclSignalInfo notifyInfo;
     603            0 :             CHK_RET(remoteSendDoneNotify_->GetNotifyData(notifyInfo));
     604            0 :             CHK_RET(remoteSendDoneNotify_->GetNotifyOffset(remoteSendDoneOffset_));
     605              : 
     606            0 :             remoteSendDoneAddress_ = notifyInfo.addr;
     607            0 :         }
     608              :     }
     609            0 :     return HCCL_SUCCESS;
     610              : }
     611              : 
     612            2 : HcclResult TransportP2p::ParseNotifyInfoEx(u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     613              : {
     614            2 :     if (machinePara_.isNewOneSide) {
     615            2 :         return HCCL_SUCCESS;
     616              :     }
     617            0 :     return ParseNotifyInfo(exchangeDataPtr, exchangeDataBlankSize);
     618              : }
     619              : 
     620            2 : HcclResult TransportP2p::ParseNotifyVectorInfo(u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     621              : {
     622            2 :     if (machinePara_.isNewOneSide) {
     623            2 :         return HCCL_SUCCESS;
     624              :     }
     625            0 :     for (u32 i = 0; i < notifyNum_; i++) {
     626            0 :         std::vector<u8> data(NOTIFY_INFO_LENGTH, 0);
     627            0 :         CHK_SAFETY_FUNC_RET(memcpy_s(&data[0], data.size(), exchangeDataPtr, NOTIFY_INFO_LENGTH));
     628            0 :         exchangeDataPtr += NOTIFY_INFO_LENGTH;
     629            0 :         exchangeDataBlankSize -= NOTIFY_INFO_LENGTH;
     630            0 :         CHK_RET(OpenRemoteNotify(data, userRemoteNotify_[i]));
     631              : 
     632            0 :         if (!machinePara_.isAicpuModeEn) {
     633              :             HcclSignalInfo notifyInfo;
     634            0 :             CHK_RET(userRemoteNotify_[i]->GetNotifyData(notifyInfo));
     635            0 :             CHK_RET(userRemoteNotify_[i]->GetNotifyOffset(userRemoteNotifyOffset_[i]));
     636            0 :             userRemoteNotifyAddr_[i] = notifyInfo.addr;
     637              :         }
     638            0 :     }
     639            0 :     return HCCL_SUCCESS;
     640              : }
     641              : 
     642              : HcclResult
     643            2 : TransportP2p::ParseCheckDataLen(ExchangeInfoSize& remoteInfoSize, u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     644              : {
     645            2 :     CHK_SAFETY_FUNC_RET(memcpy_s(&remoteInfoSize, sizeof(remoteInfoSize), exchangeDataPtr, sizeof(ExchangeInfoSize)));
     646            2 :     exchangeDataPtr += sizeof(ExchangeInfoSize);
     647            2 :     exchangeDataBlankSize -= sizeof(ExchangeInfoSize);
     648            2 :     return HCCL_SUCCESS;
     649              : }
     650              : 
     651            2 : HcclResult TransportP2p::ConstructNotifyInfoForSend(u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     652              : {
     653            2 :     if (machinePara_.isNewOneSide) {
     654            2 :         return HCCL_SUCCESS;
     655              :     }
     656            0 :     if (machinePara_.isAicpuModeEn) {
     657            0 :         if ((machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
     658            0 :              || machinePara_.machineType == MachineType::MACHINE_CLIENT_TYPE)) {
     659            0 :             RemoteRankInfo info(machinePara_.remoteDeviceId, machinePara_.remoteWorldRank, recvPid_, recvSdid_);
     660            0 :             CHK_RET(
     661              :                 notifyPool_->Alloc(machinePara_.tag, info, localSendReadyDeviceNotify_, NotifyLoadType::DEVICE_NOTIFY));
     662            0 :             std::vector<u8> data(NOTIFY_INFO_LENGTH, 0);
     663            0 :             CHK_RET(localSendReadyDeviceNotify_->Serialize(data));
     664            0 :             CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &data[0], data.size()));
     665            0 :             exchangeDataPtr += data.size();
     666            0 :             exchangeDataBlankSize -= data.size();
     667            0 :         }
     668              : 
     669            0 :         if ((machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
     670            0 :              || machinePara_.machineType == MachineType::MACHINE_SERVER_TYPE)) {
     671            0 :             RemoteRankInfo info(machinePara_.remoteDeviceId, machinePara_.remoteWorldRank, recvPid_, recvSdid_);
     672            0 :             CHK_RET(
     673              :                 notifyPool_->Alloc(machinePara_.tag, info, localSendDoneDeviceNotify_, NotifyLoadType::DEVICE_NOTIFY));
     674            0 :             std::vector<u8> data(NOTIFY_INFO_LENGTH, 0);
     675            0 :             CHK_RET(localSendDoneDeviceNotify_->Serialize(data));
     676            0 :             CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &data[0], data.size()));
     677            0 :             exchangeDataPtr += data.size();
     678            0 :             exchangeDataBlankSize -= data.size();
     679            0 :         }
     680              :     } else {
     681            0 :         if (machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
     682            0 :             || machinePara_.machineType == MachineType::MACHINE_CLIENT_TYPE) {
     683            0 :             RemoteRankInfo info(machinePara_.remoteDeviceId, machinePara_.remoteWorldRank, recvPid_, recvSdid_);
     684            0 :             CHK_RET(notifyPool_->Alloc(machinePara_.tag, info, localSendReadyNotify_));
     685            0 :             std::vector<u8> data(NOTIFY_INFO_LENGTH, 0);
     686            0 :             CHK_RET(localSendReadyNotify_->Serialize(data));
     687            0 :             CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &data[0], data.size()));
     688            0 :             exchangeDataPtr += data.size();
     689            0 :             exchangeDataBlankSize -= data.size();
     690            0 :         }
     691              : 
     692            0 :         if (machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
     693            0 :             || machinePara_.machineType == MachineType::MACHINE_SERVER_TYPE) {
     694            0 :             RemoteRankInfo info(machinePara_.remoteDeviceId, machinePara_.remoteWorldRank, recvPid_, recvSdid_);
     695            0 :             CHK_RET(notifyPool_->Alloc(machinePara_.tag, info, localSendDoneNotify_));
     696            0 :             std::vector<u8> data(NOTIFY_INFO_LENGTH, 0);
     697            0 :             CHK_RET(localSendDoneNotify_->Serialize(data));
     698            0 :             CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &data[0], data.size()));
     699            0 :             exchangeDataPtr += data.size();
     700            0 :             exchangeDataBlankSize -= data.size();
     701            0 :         }
     702              :     }
     703            0 :     return HCCL_SUCCESS;
     704              : }
     705              : 
     706            2 : HcclResult TransportP2p::ConstructNotifyVectorInfoForSend(u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     707              : {
     708            2 :     if (machinePara_.isNewOneSide) {
     709            2 :         return HCCL_SUCCESS;
     710              :     }
     711            0 :     NotifyLoadType notifyLoadType
     712            0 :         = machinePara_.isAicpuModeEn ? NotifyLoadType::DEVICE_NOTIFY : NotifyLoadType::HOST_NOTIFY;
     713            0 :     for (u32 i = 0; i < notifyNum_; i++) {
     714            0 :         RemoteRankInfo info(machinePara_.remoteDeviceId, machinePara_.remoteWorldRank, recvPid_, recvSdid_);
     715            0 :         CHK_RET(notifyPool_->Alloc(machinePara_.tag, info, userLocalNotify_[i], notifyLoadType));
     716            0 :         std::vector<u8> data(NOTIFY_INFO_LENGTH, 0);
     717            0 :         CHK_RET(userLocalNotify_[i]->Serialize(data));
     718            0 :         CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &data[0], data.size()));
     719            0 :         exchangeDataPtr += data.size();
     720            0 :         exchangeDataBlankSize -= data.size();
     721            0 :     }
     722            0 :     return HCCL_SUCCESS;
     723              : }
     724              : 
     725            2 : HcclResult TransportP2p::ConstructDataLenForSend(u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
     726              : {
     727            2 :     CHK_SAFETY_FUNC_RET(
     728              :         memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &exchangeInfoSize_, sizeof(exchangeInfoSize_)));
     729            2 :     exchangeDataPtr += sizeof(exchangeInfoSize_);
     730            2 :     exchangeDataBlankSize -= sizeof(exchangeInfoSize_);
     731            2 :     return HCCL_SUCCESS;
     732              : }
     733              : 
     734            2 : HcclResult TransportP2p::ParseReceivedExchangeData()
     735              : {
     736            2 :     s32 sendPid = 0;
     737            2 :     CHK_RET(SalGetBareTgid(&sendPid)); // 当前进程id
     738            2 :     HCCL_INFO("ParseReceivedExchangeData, sendPid[%d], recvPid[%d]", sendPid, recvPid_);
     739            2 :     u8* exchangeDataPtr = exchangeDataForRecv_.data();
     740            2 :     u64 exchangeDataBlankSize = exchangeDataTotalSize_;
     741              :     ExchangeInfoSize remoteInfoSize;
     742            2 :     CHK_RET(ParseCheckDataLen(remoteInfoSize, exchangeDataPtr, exchangeDataBlankSize));
     743            2 :     if (!exchangeInfoSize_.compare(remoteInfoSize)) {
     744            0 :         HCCL_ERROR(
     745              :             "remoteExchangeDataSize check fail, localIpcMenSize[%u] localNotifySize[%u] localExDataSize[%u]"
     746              :             "remoteIpcMenSize[%u] remoteNotifySize[%u] remoteExDataSize[%u]",
     747              :             exchangeInfoSize_.ipcMenSize, exchangeInfoSize_.notifySize, exchangeInfoSize_.exDataSize,
     748              :             remoteInfoSize.ipcMenSize, remoteInfoSize.notifySize, remoteInfoSize.exDataSize);
     749            0 :         return HCCL_E_INTERNAL;
     750              :     }
     751              : 
     752            2 :     if (sendPid != recvPid_ || recvSdid_ != INVALID_INT) {
     753            0 :         for (u32 i = 0; i < remoteIpcMemPtrVector_.size(); ++i) {
     754            0 :             CHK_RET(ParseIpcMemInfo(
     755              :                 &remoteIpcMemPtrVector_[i], remoteIpcMemSizeVector_[i], remoteIpcMemNameVector_[i].ipcName,
     756              :                 remoteIpcMemOffsetValueVector_[i], exchangeDataPtr, exchangeDataBlankSize));
     757            0 :             HCCL_INFO(
     758              :                 "[TransportP2p][ParseReceivedExchangeData]index[%d]: remoteIpcMemPtr:[%p], "
     759              :                 "remoteIpcMemSize:[%llu]",
     760              :                 i, remoteIpcMemPtrVector_[i], remoteIpcMemSizeVector_[i]);
     761              :         }
     762            0 :         if (!isMemInclude_) {
     763            0 :             CHK_RET(ParseIpcMemInfo(
     764              :                 &remoteOutputPtr_, remoteOutputSize_, remoteOutputMemName_.ipcName, remoteOutputOffsetValue_,
     765              :                 exchangeDataPtr, exchangeDataBlankSize));
     766            0 :             CHK_RET(ParseIpcMemInfo(
     767              :                 &remoteInputPtr_, remoteInputSize_, remoteInputMemName_.ipcName, remoteInputOffsetValue_,
     768              :                 exchangeDataPtr, exchangeDataBlankSize));
     769              :         } else {
     770            0 :             CHK_RET(ParseMemIncludeInfo(&remoteOutputPtr_, remoteOutputSize_, exchangeDataPtr, exchangeDataBlankSize));
     771            0 :             CHK_RET(ParseMemIncludeInfo(&remoteInputPtr_, remoteInputSize_, exchangeDataPtr, exchangeDataBlankSize));
     772              :         }
     773            0 :     } else {
     774              :         u64 memAddr;
     775            2 :         CHK_RET(ParseIntraProcMemInfo(&memAddr, &remoteOutputSize_, exchangeDataPtr, exchangeDataBlankSize));
     776            2 :         remoteOutputPtr_ = reinterpret_cast<void*>(memAddr);
     777            2 :         CHK_RET(ParseIntraProcMemInfo(&memAddr, &remoteInputSize_, exchangeDataPtr, exchangeDataBlankSize));
     778            2 :         remoteInputPtr_ = reinterpret_cast<void*>(memAddr);
     779            2 :         for (u32 i = 0; i < remoteIpcMemPtrVector_.size(); ++i) {
     780            0 :             CHK_RET(
     781              :                 ParseIntraProcMemInfo(&memAddr, &remoteIpcMemSizeVector_[i], exchangeDataPtr, exchangeDataBlankSize));
     782            0 :             remoteIpcMemPtrVector_[i] = reinterpret_cast<void*>(memAddr);
     783            0 :             HCCL_INFO(
     784              :                 "[TransportP2p][ParseReceivedExchangeData]index[%d]: remoteIpcMemPtr:[%p], "
     785              :                 "remoteIpcMemSize:[%llu]",
     786              :                 i, remoteIpcMemPtrVector_[i], remoteIpcMemSizeVector_[i]);
     787              :         }
     788              :     }
     789              :     // 将本端和远端的Mem都打印。
     790            2 :     HCCL_INFO(
     791              :         "[TransportP2p][ParseReceivedExchangeData]remoteOutputPtr_[%p], remoteOutputSize_[%llu], "
     792              :         "remoteInputPtr_[%p], remoteInputSize_[%llu]",
     793              :         remoteOutputPtr_, remoteOutputSize_, remoteInputPtr_, remoteInputSize_);
     794              : 
     795            2 :     CHK_RET(ParseNotifyInfoEx(exchangeDataPtr, exchangeDataBlankSize));
     796            2 :     CHK_RET(ParseNotifyVectorInfo(exchangeDataPtr, exchangeDataBlankSize));
     797            2 :     CHK_RET(ParseExchangeData(exchangeDataPtr, exchangeDataBlankSize));
     798              : 
     799            2 :     if (machinePara_.isIndOp) {
     800            0 :         if (sendPid != recvPid_ || recvSdid_ != INVALID_INT) {
     801              :             u64 deviceMemNum;
     802            0 :             CHK_RET(ParseMemNumInfo(deviceMemNum, exchangeDataPtr, exchangeDataBlankSize));
     803            0 :             remoteIndOpDeviceMemPtrVector_.resize(deviceMemNum);
     804            0 :             remoteIndOpDeviceMemSizeVector_.resize(deviceMemNum);
     805            0 :             remoteIndOpDeviceMemOffsetValueVector_.resize(deviceMemNum);
     806            0 :             remoteIndOpDeviceMemNameVector_.resize(deviceMemNum);
     807            0 :             for (u64 i = 0; i < deviceMemNum; ++i) {
     808            0 :                 CHK_RET(ParseIpcMemInfo(
     809              :                     &remoteIndOpDeviceMemPtrVector_[i], remoteIndOpDeviceMemSizeVector_[i],
     810              :                     remoteIndOpDeviceMemNameVector_[i].ipcName, remoteIndOpDeviceMemOffsetValueVector_[i],
     811              :                     exchangeDataPtr, exchangeDataBlankSize));
     812            0 :                 HCCL_INFO(
     813              :                     "[TransportP2p][ParseReceivedExchangeData]independent operator device mem index[%d]: "
     814              :                     "remoteIndOpDeviceMemPtr:[%p], remoteIndOpDeviceMemSize:[%llu]",
     815              :                     i, remoteIndOpDeviceMemPtrVector_[i], remoteIndOpDeviceMemSizeVector_[i]);
     816              :             }
     817              :             u64 hostMemNum;
     818            0 :             CHK_RET(ParseMemNumInfo(hostMemNum, exchangeDataPtr, exchangeDataBlankSize));
     819            0 :             remoteIndOpHostMemPtrVector_.resize(hostMemNum);
     820            0 :             remoteIndOpHostMemSizeVector_.resize(hostMemNum);
     821            0 :             remoteIndOpHostMemOffsetValueVector_.resize(hostMemNum);
     822            0 :             remoteIndOpHostMemNameVector_.resize(hostMemNum);
     823            0 :             for (u64 i = 0; i < hostMemNum; ++i) {
     824            0 :                 CHK_RET(ParseIpcMemInfo(
     825              :                     &remoteIndOpHostMemPtrVector_[i], remoteIndOpHostMemSizeVector_[i],
     826              :                     remoteIndOpHostMemNameVector_[i].ipcName, remoteIndOpHostMemOffsetValueVector_[i], exchangeDataPtr,
     827              :                     exchangeDataBlankSize));
     828            0 :                 HCCL_INFO(
     829              :                     "[TransportP2p][ParseReceivedExchangeData]independent operator host mem index[%d]: "
     830              :                     "remoteIndOpHostMemPtr:[%p], remoteIndOpHostMemSize:[%llu]",
     831              :                     i, remoteIndOpHostMemPtrVector_[i], remoteIndOpHostMemSizeVector_[i]);
     832              :             }
     833            0 :         } else {
     834              :             u64 deviceMemNum;
     835              :             u64 memAddr;
     836            0 :             CHK_RET(ParseMemNumInfo(deviceMemNum, exchangeDataPtr, exchangeDataBlankSize));
     837            0 :             remoteIndOpDeviceMemPtrVector_.resize(deviceMemNum);
     838            0 :             remoteIndOpDeviceMemSizeVector_.resize(deviceMemNum);
     839            0 :             for (u32 i = 0; i < deviceMemNum; ++i) {
     840            0 :                 CHK_RET(ParseIntraProcMemInfo(
     841              :                     &memAddr, &remoteIndOpDeviceMemSizeVector_[i], exchangeDataPtr, exchangeDataBlankSize));
     842            0 :                 remoteIndOpDeviceMemPtrVector_[i] = reinterpret_cast<void*>(memAddr);
     843            0 :                 HCCL_INFO(
     844              :                     "[TransportP2p][ParseReceivedExchangeData]independent operator device mem index[%d]: "
     845              :                     "remoteIndOpDeviceMemPtr:[%p], remoteIndOpDeviceMemSize:[%llu]",
     846              :                     i, remoteIndOpDeviceMemPtrVector_[i], remoteIndOpDeviceMemSizeVector_[i]);
     847              :             }
     848              :             u64 hostMemNum;
     849            0 :             CHK_RET(ParseMemNumInfo(hostMemNum, exchangeDataPtr, exchangeDataBlankSize));
     850            0 :             remoteIndOpHostMemPtrVector_.resize(hostMemNum);
     851            0 :             remoteIndOpHostMemSizeVector_.resize(hostMemNum);
     852            0 :             for (u32 i = 0; i < hostMemNum; ++i) {
     853            0 :                 CHK_RET(ParseIntraProcMemInfo(
     854              :                     &memAddr, &remoteIndOpHostMemSizeVector_[i], exchangeDataPtr, exchangeDataBlankSize));
     855            0 :                 remoteIndOpHostMemPtrVector_[i] = reinterpret_cast<void*>(memAddr);
     856            0 :                 HCCL_INFO(
     857              :                     "[TransportP2p][ParseReceivedExchangeData]independent operator host mem index[%d]: "
     858              :                     "remoteIndOpHostMemPtr:[%p], remoteIndOpHostMemSize:[%llu]",
     859              :                     i, remoteIndOpHostMemPtrVector_[i], remoteIndOpHostMemSizeVector_[i]);
     860              :             }
     861              :         }
     862              :     }
     863              : 
     864            2 :     if (exchangeDataBlankSize != 0) {
     865            0 :         HCCL_ERROR(
     866              :             "[TransportP2p][ParseReceivedExchangeData] failed to Parse exchange Data "
     867              :             "exchangeDataBlankSize[%llu]",
     868              :             exchangeDataBlankSize);
     869            0 :         return HCCL_E_INTERNAL;
     870              :     }
     871            2 :     return HCCL_SUCCESS; // this function should not be called in normal process
     872              : }
     873              : 
     874            0 : HcclResult TransportP2p::SignalRecord(
     875              :     std::shared_ptr<RemoteNotify>& remoteSignal, u64 remoteSignalAddr, u64 remoteSignalOffset, Stream& stream)
     876              : {
     877            0 :     return dispatcher_->SignalRecord(
     878              :         remoteSignal->ptr(), stream, machinePara_.remoteWorldRank, remoteSignalOffset, INVALID_VALUE_STAGE, false,
     879            0 :         remoteSignalAddr);
     880              : }
     881              : 
     882            0 : HcclResult TransportP2p::TxDataSignal(Stream& stream)
     883              : {
     884              :     HcclResult ret;
     885              :     /* 发起send_ready_event事件 */
     886            0 :     ret = SignalRecord(remoteSendReadyNotify_, remoteSendReadyAddress_, remoteSendReadyOffset_, stream);
     887            0 :     CHK_PRT_RET(
     888              :         ret != HCCL_SUCCESS,
     889              :         HCCL_ERROR(
     890              :             "[TransportP2p][TxDataSignal]errNo[0x%016llx]In tx data signal, signal record failed.",
     891              :             HCCL_ERROR_CODE(ret)),
     892              :         ret);
     893            0 :     return HCCL_SUCCESS;
     894              : }
     895              : 
     896            0 : HcclResult TransportP2p::RxDataSignal(Stream& stream)
     897              : {
     898              :     /* 等待send_ready_event事件 */
     899            0 :     CHK_RET(dispatcher_->SignalWait(
     900              :         localSendReadyNotify_->ptr(), stream, machinePara_.localUserrank, machinePara_.remoteWorldRank,
     901              :         INVALID_VALUE_STAGE, false, localSendReadyNotify_->notifyId_));
     902            0 :     return HCCL_SUCCESS;
     903              : }
     904              : 
     905            0 : HcclResult TransportP2p::TxAck(Stream& stream)
     906              : {
     907              :     /* 发起send_done_signal事件 */
     908            0 :     CHK_RET(SignalRecord(remoteSendDoneNotify_, remoteSendDoneAddress_, remoteSendDoneOffset_, stream));
     909            0 :     return HCCL_SUCCESS;
     910              : }
     911              : 
     912            0 : HcclResult TransportP2p::RxAck(Stream& stream)
     913              : {
     914              :     /* 等待send_done_signal事件 */
     915            0 :     CHK_RET(dispatcher_->SignalWait(
     916              :         localSendDoneNotify_->ptr(), stream, machinePara_.localUserrank, machinePara_.remoteWorldRank,
     917              :         INVALID_VALUE_STAGE, false, localSendDoneNotify_->notifyId_));
     918            0 :     return HCCL_SUCCESS;
     919              : }
     920              : 
     921            0 : HcclResult TransportP2p::TxPrepare(Stream& stream)
     922              : {
     923            0 :     CHK_RET(TxAck(stream));
     924              : 
     925            0 :     return HCCL_SUCCESS;
     926              : }
     927              : 
     928            0 : HcclResult TransportP2p::RxPrepare(Stream& stream)
     929              : {
     930            0 :     CHK_RET(RxAck(stream));
     931              : 
     932            0 :     return HCCL_SUCCESS;
     933              : }
     934              : 
     935            0 : HcclResult TransportP2p::TxDone(Stream& stream)
     936              : {
     937            0 :     HcclResult ret = RxDataSignal(stream);
     938            0 :     CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_ERROR("[TransportP2p][TxDone]RxDataSignal failed"), ret);
     939            0 :     return HCCL_SUCCESS;
     940              : }
     941              : 
     942            0 : HcclResult TransportP2p::RxDone(Stream& stream)
     943              : {
     944            0 :     HcclResult ret = TxDataSignal(stream);
     945            0 :     CHK_PRT_RET(ret != HCCL_SUCCESS, HCCL_ERROR("[TransportP2p][RxDone]TxDataSignal failed"), ret);
     946            0 :     return HCCL_SUCCESS;
     947              : }
     948              : 
     949            0 : HcclResult TransportP2p::Post(u32 notifyIdx, Stream& stream)
     950              : {
     951              :     // 校验notifyIdx有效性
     952            0 :     bool bRet = (notifyIdx >= notifyNum_);
     953            0 :     CHK_PRT_RET(
     954              :         bRet,
     955              :         HCCL_ERROR(
     956              :             "[TransportP2p][Post]notifyNum[%u], notifyIdx[%u] out of range[0, %u]", notifyNum_, notifyIdx,
     957              :             notifyNum_ - 1),
     958              :         HCCL_E_INTERNAL);
     959              : 
     960              :     // 发起send_done_signal事件
     961            0 :     CHK_RET(SignalRecord(
     962              :         userRemoteNotify_[notifyIdx], userRemoteNotifyAddr_[notifyIdx], userRemoteNotifyOffset_[notifyIdx], stream));
     963            0 :     return HCCL_SUCCESS;
     964              : }
     965              : 
     966            0 : HcclResult TransportP2p::Wait(u32 notifyIdx, Stream& stream, const u32 timeOut)
     967              : {
     968              :     // 校验notifyIdx有效性
     969            0 :     bool bRet = (notifyIdx >= notifyNum_);
     970            0 :     CHK_PRT_RET(
     971              :         bRet,
     972              :         HCCL_ERROR(
     973              :             "[TransportP2p][Wait]notifyNum[%u], notifyIdx[%u] out of range[0, %u]", notifyNum_, notifyIdx,
     974              :             notifyNum_ - 1),
     975              :         HCCL_E_INTERNAL);
     976              : 
     977              :     // 等待send_done_signal事件
     978            0 :     CHK_RET(dispatcher_->SignalWait(
     979              :         userLocalNotify_[notifyIdx]->ptr(), stream, machinePara_.localUserrank, machinePara_.remoteWorldRank,
     980              :         INVALID_VALUE_STAGE, false, userLocalNotify_[notifyIdx]->notifyId_, timeOut));
     981            0 :     return HCCL_SUCCESS;
     982              : }
     983              : 
     984            0 : HcclResult TransportP2p::ExchangeMemAndNotifyWithoutIpc()
     985              : {
     986              :     HcclResult ret;
     987              :     /* 发送 output 内存 */
     988            0 :     ret = SendMemMesgWithoutIpc(machinePara_.outputMem.ptr(), machinePara_.outputMem.size());
     989            0 :     CHK_PRT_RET(
     990              :         ret != HCCL_SUCCESS,
     991              :         HCCL_ERROR(
     992              :             "[ExchangeI][pcMesg]In exchange ipc mesg, send ipc mem output mesg fail. ret[%d], "
     993              :             "ptr[%p], size[%llu]",
     994              :             ret, machinePara_.outputMem.ptr(), machinePara_.outputMem.size()),
     995              :         ret);
     996              : 
     997              :     /* 发送 input 内存 */
     998            0 :     ret = SendMemMesgWithoutIpc(machinePara_.inputMem.ptr(), machinePara_.inputMem.size());
     999            0 :     CHK_PRT_RET(
    1000              :         ret != HCCL_SUCCESS,
    1001              :         HCCL_ERROR(
    1002              :             "[ExchangeI][pcMesg]In exchange ipc mesg, send ipc mem input mesg fail. ret[%d], "
    1003              :             "ptr[%p], size[%llu]",
    1004              :             ret, machinePara_.inputMem.ptr(), machinePara_.inputMem.size()),
    1005              :         ret);
    1006              : 
    1007              :     /* 发送 notify 信息 */
    1008            0 :     CHK_RET(LinkSendNotifyMesg());
    1009              : 
    1010              :     /* 接收 output 内存 */
    1011              :     u64 memAddr;
    1012            0 :     ret = RecvMemMesgWithoutIpc(memAddr, remoteOutputMemName_.ipcName, remoteOutputOffsetValue_);
    1013            0 :     remoteOutputPtr_ = reinterpret_cast<void*>(memAddr);
    1014            0 :     CHK_PRT_RET(
    1015              :         ret != HCCL_SUCCESS,
    1016              :         HCCL_ERROR(
    1017              :             "[Exchange][IpcMesg]In exchange ipc mesg, receive ipc output mem mesg fail. ret[%d], "
    1018              :             "ptr[%p], memptr[%p], offset[%llu]",
    1019              :             ret, remoteOutputPtr_, remoteOutputMemName_.ipcName, remoteOutputOffsetValue_),
    1020              :         ret);
    1021              : 
    1022              :     /* 接收 input 内存 */
    1023            0 :     ret = RecvMemMesgWithoutIpc(memAddr, remoteInputMemName_.ipcName, remoteInputOffsetValue_);
    1024            0 :     remoteInputPtr_ = reinterpret_cast<void*>(memAddr);
    1025            0 :     CHK_PRT_RET(
    1026              :         ret != HCCL_SUCCESS,
    1027              :         HCCL_ERROR(
    1028              :             "[Exchange][IpcMesg]In exchange ipc mesg, receive ipc input mem mesg fail. ret[%d], "
    1029              :             "ptr[%p], memptr[%p], offset[%llu]",
    1030              :             ret, remoteInputPtr_, remoteInputMemName_.ipcName, remoteInputOffsetValue_),
    1031              :         ret);
    1032              : 
    1033              :     /* 接收 notify 信息 */
    1034            0 :     CHK_RET(LinkRecvNotifyMesg());
    1035            0 :     return HCCL_SUCCESS;
    1036              : }
    1037              : 
    1038            0 : HcclResult TransportP2p::ExchangeMemAndNotifyWithIpc()
    1039              : {
    1040              :     HcclResult ret;
    1041              : 
    1042              :     /* 发送IPC output 内存 */
    1043            0 :     ret = SendIpcMemMesg(machinePara_.outputMem.ptr(), machinePara_.outputMem.size());
    1044            0 :     CHK_PRT_RET(
    1045              :         ret != HCCL_SUCCESS,
    1046              :         HCCL_ERROR(
    1047              :             "[ExchangeI][pcMesg]In exchange ipc mesg, send ipc mem output mesg fail. ret[%d], "
    1048              :             "ptr[%p], size[%llu]",
    1049              :             ret, machinePara_.outputMem.ptr(), machinePara_.outputMem.size()),
    1050              :         ret);
    1051              : 
    1052              :     /* 发送IPC input 内存 */
    1053            0 :     ret = SendIpcMemMesg(machinePara_.inputMem.ptr(), machinePara_.inputMem.size());
    1054            0 :     CHK_PRT_RET(
    1055              :         ret != HCCL_SUCCESS,
    1056              :         HCCL_ERROR(
    1057              :             "[ExchangeI][pcMesg]In exchange ipc mesg, send ipc mem input mesg fail. ret[%d], "
    1058              :             "ptr[%p], size[%llu]",
    1059              :             ret, machinePara_.inputMem.ptr(), machinePara_.inputMem.size()),
    1060              :         ret);
    1061              : 
    1062              :     /* 发送IPC notify 信息 */
    1063            0 :     CHK_RET(LinkSendNotifyMesg());
    1064              : 
    1065              :     /* 接收IPC output 内存 */
    1066            0 :     ret = RecvIpcMemMesg(&remoteOutputPtr_, remoteOutputMemName_.ipcName, remoteOutputOffsetValue_);
    1067            0 :     CHK_PRT_RET(
    1068              :         ret != HCCL_SUCCESS,
    1069              :         HCCL_ERROR(
    1070              :             "[Exchange][IpcMesg]In exchange ipc mesg, receive ipc output mem mesg fail. ret[%d], "
    1071              :             "ptr[%p], memptr[%p], offset[%llu]",
    1072              :             ret, remoteOutputPtr_, remoteOutputMemName_.ipcName, remoteOutputOffsetValue_),
    1073              :         ret);
    1074              : 
    1075              :     /* 接收IPC input 内存 */
    1076            0 :     ret = RecvIpcMemMesg(&remoteInputPtr_, remoteInputMemName_.ipcName, remoteInputOffsetValue_);
    1077            0 :     CHK_PRT_RET(
    1078              :         ret != HCCL_SUCCESS,
    1079              :         HCCL_ERROR(
    1080              :             "[Exchange][IpcMesg]In exchange ipc mesg, receive ipc input mem mesg fail. ret[%d], "
    1081              :             "ptr[%p], memptr[%p], offset[%llu]",
    1082              :             ret, remoteInputPtr_, remoteInputMemName_.ipcName, remoteInputOffsetValue_),
    1083              :         ret);
    1084              : 
    1085              :     /* 发送IPC notify 信息 */
    1086            0 :     CHK_RET(LinkRecvNotifyMesg());
    1087            0 :     return HCCL_SUCCESS;
    1088              : }
    1089              : 
    1090            0 : HcclResult TransportP2p::ExchangeMemAndNotifyMesg()
    1091              : {
    1092            0 :     s32 sendPid = 0;
    1093            0 :     CHK_RET(SalGetBareTgid(&sendPid)); // 当前进程id
    1094            0 :     HCCL_INFO("ExchangeMemAndNotifyMesg, sendPid[%d], recvPid[%d]", sendPid, recvPid_);
    1095            0 :     if (sendPid != recvPid_) {
    1096            0 :         CHK_RET(ExchangeMemAndNotifyWithIpc()); // 跨进程时处于安全考虑,交换的是IPC Memory Name
    1097              :     } else {
    1098            0 :         CHK_RET(ExchangeMemAndNotifyWithoutIpc()); // 不跨进程时,仍然使用vnic来交换,直接交换VA,不需要转成Name
    1099              :     }
    1100            0 :     return HCCL_SUCCESS;
    1101              : }
    1102              : 
    1103            0 : HcclResult TransportP2p::SendMemMesgWithoutIpc(void* ptr, u64 size) const
    1104              : {
    1105              :     HcclResult ret;
    1106              :     /* send memaddr to remote rank */
    1107            0 :     std::stringstream ss;
    1108            0 :     ss << ptr;
    1109            0 :     std::string memAddr = ss.str();
    1110            0 :     ret = defaultSocket_->Send(memAddr);
    1111            0 :     CHK_PRT_RET(
    1112              :         ret != HCCL_SUCCESS,
    1113              :         HCCL_ERROR(
    1114              :             "[Send]errNo[0x%016llx], In send ipc mesg, send name failed.remote "
    1115              :             "userrank[%u] local rank[%u]",
    1116              :             HCCL_ERROR_CODE(ret), machinePara_.remoteUserrank, machinePara_.localUserrank),
    1117              :         ret);
    1118              : 
    1119              :     /* send memsize to remote rank */
    1120            0 :     std::string memSize = std::to_string(size);
    1121            0 :     ret = defaultSocket_->Send(memSize);
    1122            0 :     CHK_PRT_RET(
    1123              :         ret != HCCL_SUCCESS,
    1124              :         HCCL_ERROR(
    1125              :             "[Send]errNo[0x%016llx]In send ipc mesg, send size failed. remote rank[%u] "
    1126              :             "size[%s] local rank[%u]",
    1127              :             HCCL_ERROR_CODE(ret), machinePara_.remoteUserrank, memSize.c_str(), machinePara_.localUserrank),
    1128              :         ret);
    1129            0 :     return HCCL_SUCCESS;
    1130            0 : }
    1131              : 
    1132            0 : HcclResult TransportP2p::RecvMemMesgWithoutIpc(u64& addr, [[maybe_unused]] u8* memName, u64& offset)
    1133              : {
    1134              :     HcclResult ret;
    1135            0 :     std::string memAddr;
    1136              : 
    1137              :     /* 获取对端地址 */
    1138            0 :     ret = defaultSocket_->Recv(memAddr);
    1139            0 :     CHK_PRT_RET(
    1140              :         ret != HCCL_SUCCESS,
    1141              :         HCCL_ERROR(
    1142              :             "[Recv]errNo[0x%016llx]In recv ipc mem mesg, receive mem name failed."
    1143              :             "remote userrank[%u] local rank[%u]",
    1144              :             HCCL_ERROR_CODE(ret), machinePara_.remoteUserrank, machinePara_.localUserrank),
    1145              :         ret);
    1146              : 
    1147            0 :     CHK_RET(SalStrToULonglong(memAddr, HCCL_BASE_HEX, addr));
    1148              :     /* 获取对端内存的大小 */
    1149            0 :     std::string remoteMemSize;
    1150            0 :     u64 size = 0;
    1151            0 :     ret = defaultSocket_->Recv(remoteMemSize);
    1152            0 :     CHK_PRT_RET(
    1153              :         ret != HCCL_SUCCESS,
    1154              :         HCCL_ERROR(
    1155              :             "[Recv]errNo[0x%016llx]In recv ipc mem mesg, receive offset name failed."
    1156              :             "remote userrank[%u] local rank[%u], remoteMemSize[%s]",
    1157              :             HCCL_ERROR_CODE(ret), machinePara_.remoteUserrank, machinePara_.localUserrank, remoteMemSize.c_str()),
    1158              :         ret);
    1159              : 
    1160            0 :     CHK_RET(SalStrToULonglong(remoteMemSize, HCCL_BASE_DECIMAL, size));
    1161              :     /* 获取对端内存的偏移值 */
    1162            0 :     offset = 0;
    1163            0 :     return ret;
    1164            0 : }
    1165              : 
    1166            0 : HcclResult TransportP2p::SendIpcMemMesg(void* ptr, u64 size) const
    1167              : {
    1168              :     HcclResult ret;
    1169              :     /* make memory shared interprocess and assigned a name */
    1170              :     u64 offset;
    1171            0 :     SecIpcName_t memName;
    1172            0 :     ret = MemNameRepository::GetInstance(machinePara_.deviceLogicId)
    1173            0 :               ->SetIpcMem(ptr, size, memName.ipcName, HCCL_IPC_MEM_NAME_LEN, offset, recvPid_, recvSdid_, isSioToHccs_);
    1174            0 :     CHK_PRT_RET(
    1175              :         ret != HCCL_SUCCESS,
    1176              :         HCCL_ERROR(
    1177              :             "[Send][IpcMemMesg]errNo[0x%016llx], In send ipc mesg, get para mem name failed. "
    1178              :             "mem addr[%p] local rank[%u]",
    1179              :             HCCL_ERROR_CODE(ret), ptr, machinePara_.localUserrank),
    1180              :         ret);
    1181              : 
    1182            0 :     std::string memOffset = std::to_string(offset);
    1183              :     /* send memName to remote rank */
    1184            0 :     ret = defaultSocket_->Send(memName.ipcName, HCCL_IPC_MEM_NAME_LEN);
    1185            0 :     CHK_PRT_RET(
    1186              :         ret != HCCL_SUCCESS,
    1187              :         HCCL_ERROR(
    1188              :             "[Send][IpcMemMesg]errNo[0x%016llx], In send ipc mesg, send name failed.remote "
    1189              :             "userrank[%u] local rank[%u]",
    1190              :             HCCL_ERROR_CODE(ret), machinePara_.remoteUserrank, machinePara_.localUserrank),
    1191              :         ret);
    1192            0 :     HCCL_INFO(
    1193              :         "localUserrank=%u, ptr=%p, remoteUserrank=%u, mem_offset=%s", machinePara_.localUserrank, ptr,
    1194              :         machinePara_.remoteUserrank, memOffset.c_str());
    1195              : 
    1196              :     /* send memsize to remote rank */
    1197            0 :     std::string memSize = std::to_string(size);
    1198            0 :     ret = defaultSocket_->Send(memSize);
    1199            0 :     CHK_PRT_RET(
    1200              :         ret != HCCL_SUCCESS,
    1201              :         HCCL_ERROR(
    1202              :             "[Send][IpcMemMesg]errNo[0x%016llx]In send ipc mesg, send size failed. remote rank[%u] "
    1203              :             "size[%s] local rank[%u]",
    1204              :             HCCL_ERROR_CODE(ret), machinePara_.remoteUserrank, memSize.c_str(), machinePara_.localUserrank),
    1205              :         ret);
    1206              : 
    1207              :     /* send memOffset to remote rank */
    1208            0 :     ret = defaultSocket_->Send(memOffset);
    1209            0 :     CHK_PRT_RET(
    1210              :         ret != HCCL_SUCCESS,
    1211              :         HCCL_ERROR(
    1212              :             "[Send][IpcMemMesg]errNo[0x%016llx]In send ipc mesg, send offset failed. remote rank[%u] "
    1213              :             "offset[%s] local rank[%u]",
    1214              :             HCCL_ERROR_CODE(ret), machinePara_.remoteUserrank, memOffset.c_str(), machinePara_.localUserrank),
    1215              :         ret);
    1216              : 
    1217            0 :     HCCL_DEBUG(
    1218              :         "localUserrank=%u, ptr=%p, remoteUserrank=%u, offset=%s", machinePara_.localUserrank, ptr,
    1219              :         machinePara_.remoteUserrank, memOffset.c_str());
    1220            0 :     return HCCL_SUCCESS;
    1221            0 : }
    1222              : 
    1223            0 : HcclResult TransportP2p::RecvIpcMemMesg(void** memPtr, u8* memName, u64& offset)
    1224              : {
    1225              :     HcclResult ret;
    1226              :     /* 获取对端内存名字 */
    1227            0 :     ret = defaultSocket_->Recv(memName, HCCL_IPC_MEM_NAME_LEN);
    1228            0 :     CHK_PRT_RET(
    1229              :         ret != HCCL_SUCCESS,
    1230              :         HCCL_ERROR(
    1231              :             "[Recv][IpcMemMesg]errNo[0x%016llx]In recv ipc mem mesg, receive mem name failed."
    1232              :             "remote userrank[%u] local rank[%u]",
    1233              :             HCCL_ERROR_CODE(ret), machinePara_.remoteUserrank, machinePara_.localUserrank),
    1234              :         ret);
    1235              :     /* 获取对端内存的大小 */
    1236            0 :     std::string remoteMemSize;
    1237            0 :     u64 size = 0;
    1238            0 :     ret = defaultSocket_->Recv(remoteMemSize);
    1239            0 :     CHK_PRT_RET(
    1240              :         ret != HCCL_SUCCESS,
    1241              :         HCCL_ERROR(
    1242              :             "[Recv][IpcMemMesg]errNo[0x%016llx]In recv ipc mem mesg, receive offset name failed."
    1243              :             "remote userrank[%u] local rank[%u], remoteMemSize[%s]",
    1244              :             HCCL_ERROR_CODE(ret), machinePara_.remoteUserrank, machinePara_.localUserrank, remoteMemSize.c_str()),
    1245              :         ret);
    1246              : 
    1247            0 :     CHK_RET(SalStrToULonglong(remoteMemSize, HCCL_BASE_DECIMAL, size));
    1248              : 
    1249              :     /* 获取对端内存的偏移值 */
    1250            0 :     std::string remoteOffsetName;
    1251            0 :     ret = defaultSocket_->Recv(remoteOffsetName);
    1252            0 :     CHK_PRT_RET(
    1253              :         ret != HCCL_SUCCESS,
    1254              :         HCCL_ERROR(
    1255              :             "[Recv][IpcMemMesg]errNo[0x%016llx]In recv ipc mem mesg, receive offset name failed."
    1256              :             "remote userrank[%u] local rank[%u], remoteOffsetName[%s]",
    1257              :             HCCL_ERROR_CODE(ret), machinePara_.remoteUserrank, machinePara_.localUserrank, remoteOffsetName.c_str()),
    1258              :         ret);
    1259              : 
    1260            0 :     CHK_RET(SalStrToULonglong(remoteOffsetName, HCCL_BASE_DECIMAL, offset));
    1261              : 
    1262              :     /* 根据名字,获取对端IPC 内存 */
    1263            0 :     ret = WaitPeerMemConfig(memPtr, const_cast<u8*>(memName), size, offset);
    1264            0 :     CHK_PRT_RET(
    1265              :         ret != HCCL_SUCCESS,
    1266              :         HCCL_ERROR(
    1267              :             "[Recv][IpcMemMesg]errNo[0x%016llx]In recv ipc mem mesg, wait peer mem config "
    1268              :             "failed. local rank[%u]",
    1269              :             HCCL_ERROR_CODE(ret), machinePara_.localUserrank),
    1270              :         ret);
    1271              : 
    1272            0 :     HCCL_DEBUG(
    1273              :         "localUserrank[%u] receive from remoteUserrank[%u]", machinePara_.localUserrank, machinePara_.remoteUserrank);
    1274              : 
    1275            0 :     return HCCL_SUCCESS;
    1276            0 : }
    1277              : 
    1278            0 : HcclResult TransportP2p::TxAsync(UserMemType dstMemType, u64 dstOffset, const void* src, u64 len, Stream& stream)
    1279              : {
    1280              :     HcclResult ret;
    1281              :     /* 源端发起数据传输 */
    1282            0 :     if (((machinePara_.linkAttribute & 0x2) == 0) && (src != nullptr)) { // 不支持目的端发起
    1283            0 :         void* dstMemPtr = nullptr;
    1284            0 :         CHK_RET(GetRemoteMem(dstMemType, &dstMemPtr));
    1285              : 
    1286            0 :         DeviceMem dstDevMem(static_cast<s8*>(dstMemPtr) + dstOffset, len);
    1287            0 :         DeviceMem srcDevMem(const_cast<void*>(src), len);
    1288              :         /* 增加hccl 数据传输时数据地址和size记录 */
    1289            0 :         HCCL_INFO(
    1290              :             "HCCL_KEY_INFO: srcAddr=[%p],srcSize=[%llu],dstAddr=[%p],dstSize=[%llu]", srcDevMem.ptr(), srcDevMem.size(),
    1291              :             dstDevMem.ptr(), dstDevMem.size());
    1292            0 :         CHK_RET(HcclD2DMemcpyAsync(
    1293              :             dispatcher_, dstDevMem, srcDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1294            0 :     }
    1295              : 
    1296              :     /* 发起send_ready_signal事件 */
    1297            0 :     ret = SignalRecord(remoteSendReadyNotify_, remoteSendReadyAddress_, remoteSendReadyOffset_, stream);
    1298            0 :     CHK_PRT_RET(
    1299              :         ret != HCCL_SUCCESS,
    1300              :         HCCL_ERROR("[TransportP2p][TxAsync]errNo[0x%016llx]In tx async, signal record failed.", HCCL_ERROR_CODE(ret)),
    1301              :         ret);
    1302              : 
    1303            0 :     return HCCL_SUCCESS;
    1304              : }
    1305              : 
    1306            0 : HcclResult TransportP2p::TxData(UserMemType dstMemType, u64 dstOffset, const void* src, u64 len, Stream& stream)
    1307              : {
    1308              :     /* 源端发起数据传输 */
    1309            0 :     if (((machinePara_.linkAttribute & 0x2) == 0) && (src != nullptr)) { // 不支持目的端发起
    1310            0 :         void* dstMemPtr = nullptr;
    1311            0 :         CHK_RET(GetRemoteMem(dstMemType, &dstMemPtr));
    1312              : 
    1313            0 :         DeviceMem dstDevMem(static_cast<s8*>(dstMemPtr) + dstOffset, len);
    1314            0 :         DeviceMem srcDevMem(const_cast<void*>(src), len);
    1315              :         /* 增加hccl 数据传输时数据地址和size记录 */
    1316            0 :         HCCL_INFO(
    1317              :             "HCCL_KEY_INFO: srcAddr=[%p],srcSize=[%llu],dstAddr=[%p],dstSize=[%llu]", srcDevMem.ptr(), srcDevMem.size(),
    1318              :             dstDevMem.ptr(), dstDevMem.size());
    1319            0 :         CHK_RET(HcclD2DMemcpyAsync(
    1320              :             dispatcher_, dstDevMem, srcDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1321            0 :     }
    1322              : 
    1323            0 :     return HCCL_SUCCESS;
    1324              : }
    1325              : 
    1326            0 : HcclResult TransportP2p::RxData(UserMemType srcMemType, u64 srcOffset, void* dst, u64 len, Stream& stream)
    1327              : {
    1328              :     /* 目的端发起数据传输 */
    1329            0 :     if ((machinePara_.linkAttribute & 0x2) && (dst != nullptr)) { // 支持目的端发起
    1330            0 :         void* srcMemPtr = nullptr;
    1331            0 :         CHK_RET(GetRemoteMem(srcMemType, &srcMemPtr));
    1332              : 
    1333            0 :         DeviceMem srcDevMem(static_cast<s8*>(srcMemPtr) + srcOffset, len);
    1334            0 :         DeviceMem dstDevMem(static_cast<s8*>(dst), len);
    1335            0 :         CHK_RET(HcclD2DMemcpyAsync(
    1336              :             dispatcher_, dstDevMem, srcDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1337            0 :     }
    1338              : 
    1339            0 :     return HCCL_SUCCESS;
    1340              : }
    1341              : 
    1342            0 : HcclResult TransportP2p::TxAsync(std::vector<TxMemoryInfo>& txMems, Stream& stream)
    1343              : {
    1344              :     HcclResult ret;
    1345              :     /* 源端发起数据传输 */
    1346            0 :     if ((machinePara_.linkAttribute & 0x2) == 0) { // 不支持目的端发起
    1347            0 :         for (auto& mem : txMems) {
    1348            0 :             CHK_PTR_NULL(mem.src);
    1349            0 :             void* dstMemPtr = nullptr;
    1350            0 :             CHK_RET(GetRemoteMem(mem.dstMemType, &dstMemPtr));
    1351              : 
    1352            0 :             DeviceMem dstDevMem(static_cast<s8*>(dstMemPtr) + mem.dstOffset, mem.len);
    1353            0 :             DeviceMem srcDevMem(const_cast<void*>(mem.src), mem.len);
    1354              :             /* 增加hccl 数据传输时数据地址和size记录 */
    1355            0 :             HCCL_INFO(
    1356              :                 "HCCL_KEY_INFO: srcAddr=[%p],srcSize=[%llu],dstAddr=[%p],dstSize=[%llu]", srcDevMem.ptr(),
    1357              :                 srcDevMem.size(), dstDevMem.ptr(), dstDevMem.size());
    1358            0 :             CHK_RET(HcclD2DMemcpyAsync(
    1359              :                 dispatcher_, dstDevMem, srcDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1360            0 :         }
    1361              :     }
    1362              : 
    1363              :     /* 发起send_ready_signal事件 */
    1364            0 :     ret = SignalRecord(remoteSendReadyNotify_, remoteSendReadyAddress_, remoteSendReadyOffset_, stream);
    1365            0 :     CHK_PRT_RET(
    1366              :         ret != HCCL_SUCCESS,
    1367              :         HCCL_ERROR("[TransportP2p][TxAsync]errNo[0x%016llx]In tx async, signal record failed.", HCCL_ERROR_CODE(ret)),
    1368              :         ret);
    1369              : 
    1370            0 :     return HCCL_SUCCESS;
    1371              : }
    1372              : 
    1373            0 : HcclResult TransportP2p::RxAsync(UserMemType srcMemType, u64 srcOffset, void* dst, u64 len, Stream& stream)
    1374              : {
    1375              :     /* 等待send_ready_signal事件 */
    1376            0 :     CHK_RET(dispatcher_->SignalWait(
    1377              :         localSendReadyNotify_->ptr(), stream, machinePara_.localUserrank, machinePara_.remoteWorldRank,
    1378              :         INVALID_VALUE_STAGE, false, localSendReadyNotify_->notifyId_));
    1379              : 
    1380              :     /* 目的端发起数据传输 */
    1381            0 :     if ((machinePara_.linkAttribute & 0x2) && (dst != nullptr)) { // 支持目的端发起
    1382            0 :         void* srcMemPtr = nullptr;
    1383            0 :         CHK_RET(GetRemoteMem(srcMemType, &srcMemPtr));
    1384              : 
    1385            0 :         DeviceMem srcDevMem(static_cast<s8*>(srcMemPtr) + srcOffset, len);
    1386            0 :         DeviceMem dstDevMem(static_cast<s8*>(dst), len);
    1387            0 :         CHK_RET(HcclD2DMemcpyAsync(
    1388              :             dispatcher_, dstDevMem, srcDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1389            0 :     }
    1390              : 
    1391            0 :     return HCCL_SUCCESS;
    1392              : }
    1393              : 
    1394            0 : HcclResult TransportP2p::RxAsync(std::vector<RxMemoryInfo>& rxMems, Stream& stream)
    1395              : {
    1396              :     /* 等待send_ready_signal事件 */
    1397            0 :     CHK_RET(dispatcher_->SignalWait(
    1398              :         localSendReadyNotify_->ptr(), stream, machinePara_.localUserrank, machinePara_.remoteWorldRank,
    1399              :         INVALID_VALUE_STAGE, false, localSendReadyNotify_->notifyId_));
    1400              : 
    1401              :     /* 目的端发起数据传输 */
    1402            0 :     if ((machinePara_.linkAttribute & 0x2) != 0) { // 支持目的端发起
    1403            0 :         for (auto& mem : rxMems) {
    1404            0 :             CHK_PTR_NULL(mem.dst);
    1405            0 :             void* srcMemPtr = nullptr;
    1406            0 :             CHK_RET(GetRemoteMem(mem.srcMemType, &srcMemPtr));
    1407              : 
    1408            0 :             DeviceMem srcDevMem(static_cast<s8*>(srcMemPtr) + mem.srcOffset, mem.len);
    1409            0 :             DeviceMem dstDevMem(static_cast<s8*>(mem.dst), mem.len);
    1410            0 :             CHK_RET(HcclD2DMemcpyAsync(
    1411              :                 dispatcher_, dstDevMem, srcDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1412            0 :         }
    1413              :     }
    1414              : 
    1415            0 :     return HCCL_SUCCESS;
    1416              : }
    1417              : 
    1418            0 : HcclResult TransportP2p::DataReceivedAck(Stream& stream)
    1419              : {
    1420            0 :     CHK_RET(TxAck(stream));
    1421            0 :     CHK_RET(RxAck(stream));
    1422            0 :     CHK_RET(TxDataSignal(stream));
    1423            0 :     CHK_RET(RxDataSignal(stream));
    1424              : 
    1425            0 :     return HCCL_SUCCESS;
    1426              : }
    1427              : 
    1428            0 : HcclResult TransportP2p::GetLocalNotify(std::vector<HcclSignalInfo>& localNotify)
    1429              : {
    1430            0 :     if (machinePara_.isNewOneSide) {
    1431            0 :         return HCCL_SUCCESS;
    1432              :     }
    1433              :     HcclSignalInfo notifyInfo;
    1434              : 
    1435            0 :     if ((machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
    1436            0 :          || machinePara_.machineType == MachineType::MACHINE_CLIENT_TYPE)) {
    1437            0 :         if (machinePara_.isAicpuModeEn) {
    1438            0 :             CHK_SMART_PTR_NULL(localSendReadyDeviceNotify_);
    1439            0 :             CHK_RET(localSendReadyDeviceNotify_->GetNotifyData(notifyInfo));
    1440            0 :             localNotify.push_back(notifyInfo);
    1441              :         } else {
    1442            0 :             CHK_SMART_PTR_NULL(localSendReadyNotify_);
    1443            0 :             CHK_RET(localSendReadyNotify_->GetNotifyData(notifyInfo));
    1444            0 :             localNotify.push_back(notifyInfo);
    1445              :         }
    1446              :     }
    1447              : 
    1448            0 :     if (machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
    1449            0 :         || machinePara_.machineType == MachineType::MACHINE_SERVER_TYPE) {
    1450            0 :         if (machinePara_.isAicpuModeEn) {
    1451            0 :             CHK_SMART_PTR_NULL(localSendDoneDeviceNotify_);
    1452            0 :             CHK_RET(localSendDoneDeviceNotify_->GetNotifyData(notifyInfo));
    1453            0 :             localNotify.push_back(notifyInfo);
    1454              :         } else {
    1455            0 :             CHK_SMART_PTR_NULL(localSendDoneNotify_);
    1456            0 :             CHK_RET(localSendDoneNotify_->GetNotifyData(notifyInfo));
    1457            0 :             localNotify.push_back(notifyInfo);
    1458              :         }
    1459              :     }
    1460              : 
    1461            0 :     bool bRet = !(notifyNum_ == userLocalNotify_.size());
    1462            0 :     CHK_PRT_RET(
    1463              :         bRet,
    1464              :         HCCL_ERROR(
    1465              :             "[TransportP2p][GetLocalNotify]size of userLocalNotify_ doesn't equal to notifyNum_[%u]", notifyNum_),
    1466              :         HCCL_E_INTERNAL);
    1467              : 
    1468              :     // 提取新增的notify资源
    1469            0 :     for (u32 i = 0; i < notifyNum_; i++) {
    1470            0 :         CHK_SMART_PTR_NULL(userLocalNotify_[i]);
    1471            0 :         CHK_RET(userLocalNotify_[i]->GetNotifyData(notifyInfo));
    1472            0 :         localNotify.push_back(notifyInfo);
    1473              :     }
    1474            0 :     return HCCL_SUCCESS;
    1475              : }
    1476              : 
    1477            0 : HcclResult TransportP2p::GetRemoteNotify(std::vector<HcclSignalInfo>& localNotify)
    1478              : {
    1479            0 :     if (machinePara_.isNewOneSide) {
    1480            0 :         return HCCL_SUCCESS;
    1481              :     }
    1482              :     HcclSignalInfo notifyInfo;
    1483            0 :     if ((machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
    1484            0 :          || machinePara_.machineType == MachineType::MACHINE_SERVER_TYPE)) {
    1485            0 :         if (machinePara_.isAicpuModeEn) {
    1486            0 :             CHK_SMART_PTR_NULL(remoteSendReadyDeviceNotify_);
    1487            0 :             CHK_RET(remoteSendReadyDeviceNotify_->GetNotifyData(notifyInfo));
    1488            0 :             localNotify.push_back(notifyInfo);
    1489              :         } else {
    1490            0 :             CHK_SMART_PTR_NULL(remoteSendReadyNotify_);
    1491            0 :             CHK_RET(remoteSendReadyNotify_->GetNotifyData(notifyInfo));
    1492            0 :             localNotify.push_back(notifyInfo);
    1493              :         }
    1494              :     }
    1495              : 
    1496            0 :     if (machinePara_.linkMode != LinkMode::LINK_SIMPLEX_MODE
    1497            0 :         || machinePara_.machineType == MachineType::MACHINE_CLIENT_TYPE) {
    1498            0 :         if (machinePara_.isAicpuModeEn) {
    1499            0 :             CHK_SMART_PTR_NULL(remoteSendDoneDeviceNotify_);
    1500            0 :             CHK_RET(remoteSendDoneDeviceNotify_->GetNotifyData(notifyInfo));
    1501            0 :             localNotify.push_back(notifyInfo);
    1502              :         } else {
    1503            0 :             CHK_SMART_PTR_NULL(remoteSendDoneNotify_);
    1504            0 :             CHK_RET(remoteSendDoneNotify_->GetNotifyData(notifyInfo));
    1505            0 :             localNotify.push_back(notifyInfo);
    1506              :         }
    1507              :     }
    1508              : 
    1509            0 :     bool bRet = !(notifyNum_ == userRemoteNotify_.size());
    1510            0 :     CHK_PRT_RET(
    1511              :         bRet,
    1512              :         HCCL_ERROR(
    1513              :             "[TransportP2p][GetRemoteNotify]size of userRemoteNotify_ doesn't equal to notifyNum_[%u]", notifyNum_),
    1514              :         HCCL_E_INTERNAL);
    1515              : 
    1516              :     // 新增notify的提取
    1517            0 :     for (u32 i = 0; i < notifyNum_; i++) {
    1518            0 :         CHK_SMART_PTR_NULL(userRemoteNotify_[i]);
    1519            0 :         CHK_RET(userRemoteNotify_[i]->GetNotifyData(notifyInfo));
    1520            0 :         localNotify.push_back(notifyInfo);
    1521              :     }
    1522            0 :     return HCCL_SUCCESS;
    1523              : }
    1524              : 
    1525            0 : HcclResult TransportP2p::GetIndOpRemoteMem(HcclMem** remoteMem, uint32_t* memNum)
    1526              : {
    1527            0 :     CHK_PRT_RET(remoteMem == nullptr, HCCL_ERROR("[%s] remoteMem is nullptr", __func__), HCCL_E_PARA);
    1528            0 :     CHK_PRT_RET(memNum == nullptr, HCCL_ERROR("[%s] memNum is nullptr", __func__), HCCL_E_PARA);
    1529              : 
    1530            0 :     *remoteMem = nullptr;
    1531            0 :     *memNum = 0;
    1532            0 :     uint32_t totalCount = remoteIndOpHostMemPtrVector_.size() + remoteIndOpDeviceMemPtrVector_.size();
    1533            0 :     if (totalCount == 0) {
    1534            0 :         HCCL_DEBUG("[%s] No remote memory regions available", __func__);
    1535            0 :         return HCCL_SUCCESS;
    1536              :     }
    1537              :     // 检查向量大小是否匹配
    1538            0 :     if (remoteIndOpHostMemPtrVector_.size() != remoteIndOpHostMemSizeVector_.size()
    1539            0 :         || remoteIndOpDeviceMemPtrVector_.size() != remoteIndOpDeviceMemSizeVector_.size()) {
    1540            0 :         HCCL_ERROR("[%s] Memory pointer and size vectors size mismatch", __func__);
    1541            0 :         return HCCL_E_INTERNAL;
    1542              :     }
    1543              :     // 外部需要手动释放内存
    1544            0 :     HcclMem* resultArray = static_cast<HcclMem*>(malloc(totalCount * sizeof(HcclMem)));
    1545            0 :     CHK_PTR_NULL(resultArray);
    1546            0 :     uint32_t index = 0;
    1547            0 :     for (size_t i = 0; i < remoteIndOpDeviceMemPtrVector_.size(); ++i) {
    1548            0 :         resultArray[index].type = HcclMemType::HCCL_MEM_TYPE_DEVICE;
    1549            0 :         resultArray[index].addr = remoteIndOpDeviceMemPtrVector_[i];
    1550            0 :         resultArray[index].size = remoteIndOpDeviceMemSizeVector_[i];
    1551            0 :         index++;
    1552              :     }
    1553            0 :     for (size_t i = 0; i < remoteIndOpHostMemPtrVector_.size(); ++i) {
    1554            0 :         resultArray[index].type = HcclMemType::HCCL_MEM_TYPE_HOST;
    1555            0 :         resultArray[index].addr = remoteIndOpHostMemPtrVector_[i];
    1556            0 :         resultArray[index].size = remoteIndOpHostMemSizeVector_[i];
    1557            0 :         index++;
    1558              :     }
    1559            0 :     *remoteMem = resultArray;
    1560            0 :     *memNum = index;
    1561              : 
    1562            0 :     HCCL_DEBUG("[%s] Successfully returned %u remote memory regions", __func__, index);
    1563              : 
    1564            0 :     return HCCL_SUCCESS;
    1565              : }
    1566              : 
    1567            0 : HcclResult TransportP2p::GetRemoteMem(UserMemType memType, void** remotePtr)
    1568              : {
    1569            0 :     switch (memType) {
    1570            0 :         case UserMemType::INPUT_MEM: {
    1571            0 :             *remotePtr = remoteInputPtr_;
    1572            0 :             break;
    1573              :         }
    1574              : 
    1575            0 :         case UserMemType::OUTPUT_MEM: {
    1576            0 :             *remotePtr = remoteOutputPtr_;
    1577            0 :             break;
    1578              :         }
    1579              : 
    1580            0 :         default: {
    1581            0 :             HCCL_ERROR("[Get][RemoteMem]not support dst_mem_type=%d", memType);
    1582            0 :             return HCCL_E_NOT_SUPPORT;
    1583              :         }
    1584              :     }
    1585              : 
    1586            0 :     return HCCL_SUCCESS;
    1587              : }
    1588              : 
    1589            0 : HcclResult TransportP2p::GetRemoteMem(std::vector<void*>* remotePtr)
    1590              : {
    1591            0 :     *remotePtr = remoteIpcMemPtrVector_;
    1592            0 :     return HCCL_SUCCESS;
    1593              : }
    1594              : 
    1595            0 : HcclResult TransportP2p::GetRemoteMemSize(UserMemType memType, u64& size)
    1596              : {
    1597            0 :     switch (memType) {
    1598            0 :         case UserMemType::INPUT_MEM: {
    1599            0 :             size = remoteInputSize_;
    1600            0 :             break;
    1601              :         }
    1602              : 
    1603            0 :         case UserMemType::OUTPUT_MEM: {
    1604            0 :             size = remoteOutputSize_;
    1605            0 :             break;
    1606              :         }
    1607              : 
    1608            0 :         default: {
    1609            0 :             HCCL_ERROR("[Get][RemoteMem]not support dst_mem_type=%d", memType);
    1610            0 :             return HCCL_E_NOT_SUPPORT;
    1611              :         }
    1612              :     }
    1613              : 
    1614            0 :     return HCCL_SUCCESS;
    1615              : }
    1616              : 
    1617            0 : HcclResult TransportP2p::WaitPeerMemConfig(void** memPtr, const u8* memName, uint64_t size, u64 offset)
    1618              : {
    1619            0 :     CHK_PTR_NULL(memPtr);
    1620            0 :     CHK_PTR_NULL(memName);
    1621              : 
    1622            0 :     bool firstOpened = false;
    1623              :     // 支持进程间、进程内都可以通过name获取对端内存
    1624              :     HcclResult ret = MemNameRepository::GetInstance(machinePara_.deviceLogicId)
    1625            0 :                          ->OpenIpcMem(memPtr, size, memName, HCCL_IPC_MEM_NAME_LEN, offset, firstOpened, isSioToHccs_);
    1626            0 :     CHK_PRT_RET(
    1627              :         ret != HCCL_SUCCESS,
    1628              :         HCCL_ERROR(
    1629              :             "[Wait][WaitPeerMemConfig]errNo[0x%016llx]In link pcie, open mem failed. "
    1630              :             "offset[%llu], size[%llu Byte], linkType[%d]",
    1631              :             HCCL_ERROR_CODE(ret), offset, size, transportAttr_.linkType),
    1632              :         ret);
    1633            0 :     return HCCL_SUCCESS;
    1634              : }
    1635              : 
    1636            0 : HcclResult TransportP2p::PostReady(Stream& stream)
    1637              : {
    1638            0 :     CHK_RET(SignalRecord(remoteSendReadyNotify_, remoteSendReadyAddress_, remoteSendReadyOffset_, stream));
    1639            0 :     return HCCL_SUCCESS;
    1640              : }
    1641              : 
    1642            0 : HcclResult TransportP2p::WaitReady(Stream& stream)
    1643              : {
    1644            0 :     CHK_RET(dispatcher_->SignalWait(
    1645              :         localSendReadyNotify_->ptr(), stream, machinePara_.localUserrank, machinePara_.remoteWorldRank,
    1646              :         INVALID_VALUE_STAGE, false, localSendReadyNotify_->notifyId_));
    1647            0 :     return HCCL_SUCCESS;
    1648              : }
    1649              : 
    1650            0 : HcclResult TransportP2p::PostFin(Stream& stream)
    1651              : {
    1652            0 :     CHK_RET(SignalRecord(remoteSendDoneNotify_, remoteSendDoneAddress_, remoteSendDoneOffset_, stream));
    1653            0 :     return HCCL_SUCCESS;
    1654              : }
    1655              : 
    1656            0 : HcclResult TransportP2p::WaitFin(Stream& stream)
    1657              : {
    1658            0 :     CHK_RET(dispatcher_->SignalWait(
    1659              :         localSendDoneNotify_->ptr(), stream, machinePara_.localUserrank, machinePara_.remoteWorldRank,
    1660              :         INVALID_VALUE_STAGE, false, localSendDoneNotify_->notifyId_));
    1661            0 :     return HCCL_SUCCESS;
    1662              : }
    1663              : 
    1664              : HcclResult
    1665            0 : TransportP2p::WriteSync(struct Transport::Buffer& remoteBuf, struct Transport::Buffer& localBuf, Stream& stream)
    1666              : {
    1667            0 :     DeviceMem remoteDevMem(const_cast<void*>(remoteBuf.addr), remoteBuf.size);
    1668            0 :     DeviceMem localDevMem(const_cast<void*>(localBuf.addr), localBuf.size);
    1669            0 :     HCCL_INFO(
    1670              :         "HCCL_KEY_INFO: localAddr=[%p],localSize=[%llu],remoteAddr=[%p],remoteSize=[%llu]", localDevMem.ptr(),
    1671              :         localDevMem.size(), remoteDevMem.ptr(), remoteDevMem.size());
    1672            0 :     CHK_RET(HcclD2DMemcpyAsync(
    1673              :         dispatcher_, remoteDevMem, localDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1674            0 :     return HCCL_SUCCESS;
    1675            0 : }
    1676              : 
    1677              : HcclResult
    1678            0 : TransportP2p::WriteAsyncEx(struct Transport::Buffer& remoteBuf, struct Transport::Buffer& localBuf, Stream& stream)
    1679              : {
    1680            0 :     bool isLocalHostAddr = false;
    1681            0 :     bool isRemoteHostAddr = false;
    1682            0 :     struct Transport::Buffer newLocalBuf {};
    1683            0 :     struct Transport::Buffer newRemoteBuf {};
    1684            0 :     CHK_RET(ReplaceMemAddr(localBuf, remoteBuf, newLocalBuf, newRemoteBuf, isLocalHostAddr, isRemoteHostAddr));
    1685            0 :     DeviceMem dstDevMem(const_cast<void*>(newRemoteBuf.addr), newRemoteBuf.size);
    1686            0 :     DeviceMem srcDevMem(const_cast<void*>(newLocalBuf.addr), newLocalBuf.size);
    1687            0 :     CHK_RET(reinterpret_cast<DispatcherPub*>(dispatcher_)
    1688              :                 ->MemcpyAsync(dstDevMem, srcDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1689            0 :     return HCCL_SUCCESS;
    1690            0 : }
    1691              : 
    1692              : HcclResult
    1693            0 : TransportP2p::WriteAsync(struct Transport::Buffer& remoteBuf, struct Transport::Buffer& localBuf, Stream& stream)
    1694              : {
    1695            0 :     if (machinePara_.isNewOneSide) {
    1696            0 :         return WriteAsyncEx(remoteBuf, localBuf, stream);
    1697              :     }
    1698              : 
    1699            0 :     DeviceMem remoteDevMem(const_cast<void*>(remoteBuf.addr), remoteBuf.size);
    1700            0 :     DeviceMem localDevMem(const_cast<void*>(localBuf.addr), localBuf.size);
    1701            0 :     HCCL_INFO(
    1702              :         "HCCL_KEY_INFO: localAddr=[%p],localSize=[%llu],remoteAddr=[%p],remoteSize=[%llu]", localDevMem.ptr(),
    1703              :         localDevMem.size(), remoteDevMem.ptr(), remoteDevMem.size());
    1704            0 :     CHK_RET(HcclD2DMemcpyAsync(
    1705              :         dispatcher_, remoteDevMem, localDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1706            0 :     return HCCL_SUCCESS;
    1707            0 : }
    1708              : 
    1709            0 : HcclResult TransportP2p::WriteReduceAsync(
    1710              :     struct Transport::Buffer& remoteBuf, struct Transport::Buffer& localBuf, const HcclDataType datatype,
    1711              :     HcclReduceOp redOp, Stream& stream)
    1712              : {
    1713            0 :     HCCL_INFO(
    1714              :         "HCCL_KEY_INFO: localAddr=[%p],localSize=[%llu],remoteAddr=[%p],remoteSize=[%llu]", localBuf.addr,
    1715              :         localBuf.size, remoteBuf.addr, remoteBuf.size);
    1716              : 
    1717            0 :     u64 reduceAttr = 0;
    1718            0 :     if (IsSpInlineReduce()) {
    1719            0 :         reduceAttr = INLINE_REDUCE_BIT;
    1720              :     }
    1721            0 :     CHK_RET(HcclReduceAsync(
    1722              :         dispatcher_, const_cast<void*>(localBuf.addr), remoteBuf.size / SIZE_TABLE[datatype], datatype, redOp, stream,
    1723              :         const_cast<void*>(remoteBuf.addr), GetRemoteRank(), GetLinkType(), reduceAttr));
    1724            0 :     return HCCL_SUCCESS;
    1725              : }
    1726              : 
    1727              : HcclResult
    1728            0 : TransportP2p::ReadSync(struct Transport::Buffer& localBuf, struct Transport::Buffer& remoteBuf, Stream& stream)
    1729              : {
    1730            0 :     DeviceMem remoteDevMem(const_cast<void*>(remoteBuf.addr), remoteBuf.size);
    1731            0 :     DeviceMem localDevMem(const_cast<void*>(localBuf.addr), localBuf.size);
    1732            0 :     HCCL_INFO(
    1733              :         "HCCL_KEY_INFO: localAddr=[%p],localSize=[%llu],remoteAddr=[%p],remoteSize=[%llu]", localDevMem.ptr(),
    1734              :         localDevMem.size(), remoteDevMem.ptr(), remoteDevMem.size());
    1735            0 :     CHK_RET(HcclD2DMemcpyAsync(
    1736              :         dispatcher_, localDevMem, remoteDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1737            0 :     return HCCL_SUCCESS;
    1738            0 : }
    1739              : 
    1740            0 : HcclResult TransportP2p::ReadReduceSync(
    1741              :     struct Transport::Buffer& localBuf, struct Transport::Buffer& remoteBuf, const HcclDataType datatype,
    1742              :     HcclReduceOp redOp, Stream& stream)
    1743              : {
    1744            0 :     HCCL_INFO(
    1745              :         "HCCL_KEY_INFO: localAddr=[%p],localSize=[%llu],remoteAddr=[%p],remoteSize=[%llu]", localBuf.addr,
    1746              :         localBuf.size, remoteBuf.addr, remoteBuf.size);
    1747              : 
    1748            0 :     u64 reduceAttr = 0;
    1749            0 :     if (IsSpInlineReduce()) {
    1750            0 :         reduceAttr = INLINE_REDUCE_BIT;
    1751              :     }
    1752            0 :     CHK_RET(HcclReduceAsync(
    1753              :         dispatcher_, const_cast<void*>(remoteBuf.addr), remoteBuf.size / SIZE_TABLE[datatype], datatype, redOp, stream,
    1754              :         const_cast<void*>(localBuf.addr), GetRemoteRank(), GetLinkType(), reduceAttr));
    1755            0 :     return HCCL_SUCCESS;
    1756              : }
    1757              : 
    1758              : HcclResult
    1759            0 : TransportP2p::ReadAsyncEx(struct Transport::Buffer& localBuf, struct Transport::Buffer& remoteBuf, Stream& stream)
    1760              : {
    1761            0 :     bool isLocalHostAddr = false;
    1762            0 :     bool isRemoteHostAddr = false;
    1763            0 :     struct Transport::Buffer newLocalBuf {};
    1764            0 :     struct Transport::Buffer newRemoteBuf {};
    1765            0 :     CHK_RET(ReplaceMemAddr(localBuf, remoteBuf, newLocalBuf, newRemoteBuf, isLocalHostAddr, isRemoteHostAddr));
    1766            0 :     DeviceMem dstDevMem(const_cast<void*>(newLocalBuf.addr), newLocalBuf.size);
    1767            0 :     DeviceMem srcDevMem(const_cast<void*>(newRemoteBuf.addr), newRemoteBuf.size);
    1768            0 :     CHK_RET(reinterpret_cast<DispatcherPub*>(dispatcher_)
    1769              :                 ->MemcpyAsync(dstDevMem, srcDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType));
    1770            0 :     return HCCL_SUCCESS;
    1771            0 : }
    1772              : 
    1773              : HcclResult
    1774            0 : TransportP2p::ReadAsync(struct Transport::Buffer& localBuf, struct Transport::Buffer& remoteBuf, Stream& stream)
    1775              : {
    1776            0 :     if (machinePara_.isNewOneSide) {
    1777            0 :         return ReadAsyncEx(localBuf, remoteBuf, stream);
    1778              :     }
    1779            0 :     DeviceMem dstDevMem(const_cast<void*>(localBuf.addr), localBuf.size);
    1780            0 :     DeviceMem srcDevMem(const_cast<void*>(remoteBuf.addr), remoteBuf.size);
    1781            0 :     return HcclD2DMemcpyAsync(
    1782            0 :         dispatcher_, dstDevMem, srcDevMem, stream, machinePara_.remoteWorldRank, transportAttr_.linkType);
    1783            0 : }
    1784              : 
    1785            6 : HcclResult TransportP2p::SumCheckSizeAndConsisten(
    1786              :     ExInfoType exInfoType, u32 rightInfoSize, u64& blankSizeRecord, u64 exchangeDataBlankSize)
    1787              : {
    1788            6 :     u32 checkInfoSize = blankSizeRecord - exchangeDataBlankSize;
    1789            6 :     if (checkInfoSize != rightInfoSize) {
    1790            0 :         HCCL_ERROR(
    1791              :             "[SumCheckSizeAndConsisten] ExInfoType[%d] check size failed, checkInfoSize[%u] rightInfoSize[%u]",
    1792              :             exInfoType, checkInfoSize, rightInfoSize);
    1793            0 :         return HCCL_E_INTERNAL;
    1794              :     }
    1795            6 :     blankSizeRecord = exchangeDataBlankSize;
    1796            6 :     return HCCL_SUCCESS;
    1797              : }
    1798              : 
    1799            0 : HcclResult TransportP2p::ConstructMemIncludeInfoForSend(u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
    1800              : {
    1801            0 :     u64 outputSize = machinePara_.outputMem.size();
    1802              :     u64 outputOffset
    1803            0 :         = reinterpret_cast<u64>(machinePara_.outputMem.ptr()) - reinterpret_cast<u64>(machinePara_.mem[0].ptr());
    1804            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &outputSize, sizeof(u64)));
    1805            0 :     exchangeDataPtr += sizeof(u64);
    1806            0 :     exchangeDataBlankSize -= sizeof(u64);
    1807            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &outputOffset, sizeof(u64)));
    1808            0 :     exchangeDataPtr += sizeof(u64);
    1809            0 :     exchangeDataBlankSize -= sizeof(u64);
    1810              : 
    1811            0 :     u64 inputSize = machinePara_.inputMem.size();
    1812              :     u64 inputOffset
    1813            0 :         = reinterpret_cast<u64>(machinePara_.inputMem.ptr()) - reinterpret_cast<u64>(machinePara_.mem[0].ptr());
    1814            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &inputSize, sizeof(u64)));
    1815            0 :     exchangeDataPtr += sizeof(u64);
    1816            0 :     exchangeDataBlankSize -= sizeof(u64);
    1817            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(exchangeDataPtr, exchangeDataBlankSize, &inputOffset, sizeof(u64)));
    1818            0 :     exchangeDataPtr += sizeof(u64);
    1819            0 :     exchangeDataBlankSize -= sizeof(u64);
    1820              : 
    1821            0 :     return HCCL_SUCCESS;
    1822              : }
    1823              : 
    1824            0 : HcclResult TransportP2p::ParseMemIncludeInfo(void** memPtr, u64& size, u8*& exchangeDataPtr, u64& exchangeDataBlankSize)
    1825              : {
    1826            0 :     u64 memOffset = 0;
    1827            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(&size, sizeof(u64), exchangeDataPtr, sizeof(u64)));
    1828            0 :     exchangeDataPtr += sizeof(u64);
    1829            0 :     exchangeDataBlankSize -= sizeof(u64);
    1830            0 :     CHK_SAFETY_FUNC_RET(memcpy_s(&memOffset, sizeof(u64), exchangeDataPtr, sizeof(u64)));
    1831            0 :     exchangeDataPtr += sizeof(u64);
    1832            0 :     exchangeDataBlankSize -= sizeof(u64);
    1833            0 :     if (!machinePara_.isNewOneSide) {
    1834            0 :         *memPtr = reinterpret_cast<void*>(reinterpret_cast<u64>(remoteIpcMemPtrVector_[0]) + memOffset);
    1835              :     }
    1836            0 :     return HCCL_SUCCESS;
    1837              : }
    1838              : 
    1839            2 : void TransportP2p::SetMemIncludeFlag()
    1840              : {
    1841            2 :     if (machinePara_.mem.empty()) {
    1842            2 :         return;
    1843              :     }
    1844              :     // 当前只取mem[0]  ->expMem
    1845            0 :     u64 memPtr = reinterpret_cast<u64>(machinePara_.mem[0].ptr());
    1846            0 :     u64 memEndPtr = memPtr + machinePara_.mem[0].size();
    1847            0 :     u64 inputMemPtr = reinterpret_cast<u64>(machinePara_.inputMem.ptr());
    1848            0 :     u64 inputMemEndPtr = inputMemPtr + machinePara_.inputMem.size();
    1849            0 :     u64 outputMemPtr = reinterpret_cast<u64>(machinePara_.outputMem.ptr());
    1850            0 :     u64 outputMemEndPtr = outputMemPtr + machinePara_.outputMem.size();
    1851            0 :     HCCL_DEBUG(
    1852              :         "[SetMemIncludeFlag] memPtr[%u] memEndPtr[%u], inputMemPtr[%u] inputMemEndPtr[%u],",
    1853              :         "outputMemPtr[%u] outputMemEndPtr[%u]", memPtr, memEndPtr, inputMemPtr, inputMemEndPtr, outputMemPtr,
    1854              :         outputMemEndPtr);
    1855            0 :     if ((memPtr <= inputMemPtr && inputMemEndPtr <= memEndPtr)
    1856            0 :         && (memPtr <= outputMemPtr && outputMemEndPtr <= memEndPtr)) {
    1857            0 :         isMemInclude_ = true;
    1858              :     }
    1859            0 :     return;
    1860              : }
    1861              : 
    1862            0 : HcclResult TransportP2p::ReplaceMemAddr(
    1863              :     Transport::Buffer& localMem, Transport::Buffer& remoteMem, Transport::Buffer& newLocalMem,
    1864              :     Transport::Buffer& newRemoteMem, bool& isLocalHostAddr, bool& isRemoteHostAddr)
    1865              : {
    1866            0 :     HCCL_DEBUG(
    1867              :         "[TransportP2p][ReplaceMemAddr]old localAddr=[%p],localSize=[%llu],remoteAddr=[%p],remoteSize=[%llu]",
    1868              :         localMem.addr, localMem.size, remoteMem.addr, remoteMem.size);
    1869              : 
    1870            0 :     isLocalHostAddr = false;
    1871            0 :     isRemoteHostAddr = false;
    1872            0 :     void* localAddr = const_cast<void*>(localMem.addr);
    1873            0 :     u64 localSize = localMem.size;
    1874            0 :     auto localKey = BufferKey<uintptr_t, u64>(reinterpret_cast<uintptr_t>(localAddr), localSize);
    1875            0 :     auto localBufferPair = localHcclMemExMgr_.Find(localKey);
    1876            0 :     if (localBufferPair.first) {
    1877            0 :         std::shared_ptr<HcclMemEx>& localBufMemPtr = localBufferPair.second;
    1878            0 :         u64 localDataOffSet = static_cast<u8*>(localAddr) - static_cast<u8*>(localBufMemPtr->addr);
    1879            0 :         newLocalMem.addr = static_cast<void*>(static_cast<u8*>(localBufMemPtr->devAddr) + localDataOffSet);
    1880            0 :         newLocalMem.size = localMem.size;
    1881            0 :         if (localBufMemPtr->type == HcclMemType::HCCL_MEM_TYPE_HOST) {
    1882            0 :             isLocalHostAddr = true;
    1883              :         }
    1884              :     } else {
    1885            0 :         HCCL_DEBUG("[TransportP2p][ReplaceMemAddr] Can't find localBufferPair by key {%p, %llu}", localAddr, localSize);
    1886            0 :         newLocalMem.addr = localAddr;
    1887            0 :         newLocalMem.size = localSize;
    1888              :     }
    1889              : 
    1890            0 :     void* remoteAddr = const_cast<void*>(remoteMem.addr);
    1891            0 :     auto remoteKey = BufferKey<uintptr_t, u64>(reinterpret_cast<uintptr_t>(remoteAddr), remoteMem.size);
    1892            0 :     auto remoteBufferPair = remoteHcclMemExMgr_.Find(remoteKey);
    1893            0 :     if (remoteBufferPair.first) {
    1894            0 :         std::shared_ptr<HcclMemEx>& remoteBufMemPtr = remoteBufferPair.second;
    1895            0 :         u64 remoteDataOffSet = static_cast<u8*>(remoteAddr) - static_cast<u8*>(remoteBufMemPtr->addr);
    1896            0 :         newRemoteMem.addr = static_cast<void*>(static_cast<u8*>(remoteBufMemPtr->devAddr) + remoteDataOffSet);
    1897            0 :         if (remoteBufMemPtr->type == HcclMemType::HCCL_MEM_TYPE_HOST) {
    1898            0 :             isRemoteHostAddr = true;
    1899              :         }
    1900              :     } else {
    1901            0 :         HCCL_DEBUG(
    1902              :             "[TransportP2p][ReplaceMemAddr] Can't find remoteBuffer by key {%p, %llu}", remoteAddr, remoteMem.size);
    1903            0 :         newRemoteMem.addr = remoteMem.addr;
    1904              :     }
    1905            0 :     newRemoteMem.size = remoteMem.size;
    1906              : 
    1907            0 :     HCCL_DEBUG(
    1908              :         "[TransportP2p][ReplaceMemAddr]old localAddr=[%p],localSize=[%llu],remoteAddr=[%p],remoteSize=[%llu], "
    1909              :         "isLocalHostAddr[%u] isRemoteHostAddr[%u]",
    1910              :         newLocalMem.addr, newLocalMem.size, newRemoteMem.addr, newRemoteMem.size,
    1911              :         static_cast<uint32_t>(isLocalHostAddr), static_cast<uint32_t>(isRemoteHostAddr));
    1912            0 :     return HCCL_SUCCESS;
    1913            0 : }
    1914              : 
    1915            0 : HcclResult TransportP2p::InitHcclMemExMgrWithMem(HcclMemEx* bufMem, u32 bufSize, HcclMemExMgr& hcommMemExMgr)
    1916              : {
    1917            0 :     for (u32 i = 0; i < bufSize; i++) {
    1918            0 :         HcclMemEx& bufMemTmp = bufMem[i];
    1919              : 
    1920            0 :         std::shared_ptr<HcclMemEx> hcclMemEx = nullptr;
    1921            0 :         hcclMemEx = std::make_shared<HcclMemEx>();
    1922            0 :         CHK_PTR_NULL(hcclMemEx);
    1923              : 
    1924            0 :         HcclMemEx* hcclMemExPtr = reinterpret_cast<HcclMemEx*>(hcclMemEx.get());
    1925            0 :         *hcclMemExPtr = bufMemTmp;
    1926              : 
    1927            0 :         hccl::BufferKey<uintptr_t, u64> tempKey(reinterpret_cast<uintptr_t>(bufMemTmp.addr), bufMemTmp.size);
    1928            0 :         auto resultPair = hcommMemExMgr.Add(tempKey, hcclMemEx);
    1929            0 :         if (!resultPair.second) {
    1930            0 :             HCCL_ERROR(
    1931              :                 "[TransportP2p][InitHcclMemExMgrWithMem]add addr:%p, size[%lu], type[%u], devAddr[%p] fail",
    1932              :                 bufMemTmp.addr, bufMemTmp.size, static_cast<u32>(bufMemTmp.type), bufMemTmp.devAddr);
    1933            0 :             return HCCL_E_INTERNAL;
    1934              :         } else {
    1935            0 :             HCCL_INFO(
    1936              :                 "[TransportP2p][InitHcclMemExMgrWithMem]add addr:%p, size[%lu], type[%u], devAddr[%p] done",
    1937              :                 bufMemTmp.addr, bufMemTmp.size, static_cast<u32>(bufMemTmp.type), bufMemTmp.devAddr);
    1938              :         }
    1939            0 :     }
    1940              : 
    1941            0 :     HCCL_INFO("[TransportP2p][InitHcclMemExMgrWithMem] done");
    1942            0 :     return HCCL_SUCCESS;
    1943              : }
    1944              : 
    1945            0 : HcclResult TransportP2p::InitHcclMemExMgr(MachinePara& machinePara)
    1946              : {
    1947            0 :     HCCL_INFO("[TransportP2p][InitHcclMemExMgr] start");
    1948            0 :     CHK_RET(InitHcclMemExMgrWithMem(machinePara.localBufMem, machinePara.localBufSize, localHcclMemExMgr_));
    1949            0 :     machinePara.localBufMem = nullptr;
    1950            0 :     machinePara.localBufSize = 0;
    1951            0 :     HCCL_INFO("[TransportP2p][InitHcclMemExMgr] local done");
    1952            0 :     CHK_RET(InitHcclMemExMgrWithMem(machinePara.remoteBufMem, machinePara.remoteBufSize, remoteHcclMemExMgr_));
    1953            0 :     machinePara.remoteBufMem = nullptr;
    1954            0 :     machinePara.remoteBufSize = 0;
    1955            0 :     HCCL_INFO("[TransportP2p][InitHcclMemExMgr] remote done");
    1956            0 :     return HCCL_SUCCESS;
    1957              : }
    1958              : } // namespace hccl
        

Generated by: LCOV version 2.0-1