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 : ¬ifyVaule, 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
|