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 "alltoallv_staged_mesh.h"
12 : #include "log.h"
13 : #include "alg_template_register.h"
14 :
15 : namespace hccl {
16 : using namespace std;
17 :
18 1 : AlltoAllVStagedMesh::AlltoAllVStagedMesh(const HcclDispatcher dispatcher)
19 1 : : AlltoAllVStagedBase(dispatcher)
20 : {
21 1 : }
22 :
23 2 : AlltoAllVStagedMesh::~AlltoAllVStagedMesh() {}
24 :
25 : // 图模式Prepare入口
26 1 : HcclResult AlltoAllVStagedMesh::Prepare(DeviceMem &sendMem, DeviceMem &recvMem, StageAlltoAllVAddrInfo &sendAddrInfo,
27 : StageAlltoAllVAddrInfo &recvAddrInfo, bool isAlltoAllZCopyMode, u32 userRank, Stream &mainStream,
28 : std::vector<Stream> &subStreams,
29 : std::vector<std::shared_ptr<LocalNotify>> &meshSignalMainToSub,
30 : std::vector<std::shared_ptr<LocalNotify>> &meshSignalSubToMain)
31 : {
32 1 : userRank_ = userRank;
33 1 : subStreamsPtr_ = &subStreams;
34 1 : meshSignalMainToSubPtr_ = &meshSignalMainToSub;
35 1 : meshSignalSubToMainPtr_ = &meshSignalSubToMain;
36 :
37 1 : HCCL_DEBUG("[AlltoAllVStagedMesh][Prepare] and isAlltoAllZCopyMode_[%d]", isAlltoAllZCopyMode_);
38 1 : return AlltoAllVStagedBase::Prepare(sendMem, recvMem, sendAddrInfo, recvAddrInfo,
39 1 : isAlltoAllZCopyMode, mainStream);
40 : }
41 :
42 1 : HcclResult AlltoAllVStagedMesh::Prepare(DeviceMem &sendMem, DeviceMem &recvMem, DeviceMem &scratchInputMem,
43 : DeviceMem &scratchOutputMem, StageAlltoAllVAddrInfo &sendAddrInfo, StageAlltoAllVAddrInfo &recvAddrInfo,
44 : bool isAlltoAllZCopyMode, u32 userRank, Stream &mainStream, std::vector<Stream> &subStreams,
45 : std::vector<std::shared_ptr<LocalNotify>> &meshSignalMainToSub,
46 : std::vector<std::shared_ptr<LocalNotify>> &meshSignalSubToMain)
47 : {
48 : (void)userRank_;
49 : (void)subStreamsPtr_;
50 : (void)meshSignalMainToSubPtr_;
51 : (void)meshSignalSubToMainPtr_;
52 : (void)scratchInputMem;
53 : (void)scratchOutputMem;
54 1 : CHK_RET(AlltoAllVStagedBase::Prepare(sendMem, recvMem, sendAddrInfo, recvAddrInfo,
55 : isAlltoAllZCopyMode, mainStream));
56 1 : HCCL_ERROR("AlltoAllv Staged Mesh is not supported in Op base mode!");
57 1 : return HCCL_E_NOT_SUPPORT;
58 : }
59 :
60 0 : HcclResult AlltoAllVStagedMesh::RunAsync(const u32 rank, const u32 rankSize, const std::vector<LINK> &links)
61 : {
62 0 : HCCL_INFO("[AlltoAllVStagedMesh][RunAsync]: rank[%u] transportSize[%llu]", rank, links.size());
63 0 : CHK_SMART_PTR_NULL(dispatcher_);
64 0 : CHK_SMART_PTR_NULL(subStreamsPtr_);
65 :
66 0 : CHK_PRT_RET(rankSize == 0, HCCL_ERROR("[AlltoAllVStagedMesh][Prepare] invilad rankSize[%u]", rankSize),
67 : HCCL_E_PARA);
68 :
69 0 : CHK_PRT_RET(rankSize != links.size(),
70 : HCCL_ERROR("[AlltoAllVStagedMesh][RunAsync]: rankSize[%u] and transport size[%llu] do not match", rankSize,
71 : links.size()),
72 : HCCL_E_PARA);
73 :
74 0 : CHK_PRT_RET(rankSize > 1 && subStreamsPtr_->size() < rankSize - 2, // 从流个数是rankSize - 2,ranksize为1时不需要校验
75 : HCCL_ERROR("[AlltoAllVStagedMesh][RunAsync]: rankSize[%u] and stream size[%llu] do not match", rankSize,
76 : subStreamsPtr_->size()),
77 : HCCL_E_PARA);
78 :
79 0 : bool sizeEqual = (sendAddrInfo_.size() == recvAddrInfo_.size() && sendAddrInfo_.size() == rankSize);
80 0 : CHK_PRT_RET(!sizeEqual,
81 : HCCL_ERROR("[AlltoAllVStagedMesh][RunAsync] invilad params: "\
82 : "sendAddrInfo size[%u] recvAddrInfo size[%u] rankSize[%u]",
83 : sendAddrInfo_.size(), recvAddrInfo_.size(), rankSize),
84 : HCCL_E_PARA);
85 :
86 0 : CHK_RET(LocalCopy(rank));
87 0 : CHK_RET(RunZCopyMode(rank, rankSize, links));
88 :
89 0 : return HCCL_SUCCESS;
90 : }
91 :
92 0 : void AlltoAllVStagedMesh::BuildSendRecvMemoryInfo(vector<TxMemoryInfo> &txMems, vector<RxMemoryInfo> &rxMems,
93 : u32 destRank)
94 : {
95 0 : u32 index = 0;
96 0 : for (auto &addrInfo : sendAddrInfo_[destRank]) {
97 0 : txMems[index].dstMemType = UserMemType::OUTPUT_MEM;
98 0 : txMems[index].dstOffset = addrInfo.remoteOffset;
99 0 : txMems[index].src = static_cast<u8 *>(sendMem_.ptr()) + addrInfo.localOffset;
100 0 : txMems[index].len = addrInfo.localLength;
101 0 : index++;
102 : }
103 :
104 0 : index = 0;
105 0 : for (auto &addrInfo : recvAddrInfo_[destRank]) {
106 0 : rxMems[index].srcMemType = UserMemType::INPUT_MEM;
107 0 : rxMems[index].srcOffset = addrInfo.remoteOffset;
108 0 : rxMems[index].dst = static_cast<u8 *>(recvMem_.ptr()) + addrInfo.localOffset;
109 0 : rxMems[index].len = addrInfo.localLength;
110 0 : index++;
111 : }
112 0 : return;
113 : }
114 :
115 0 : HcclResult AlltoAllVStagedMesh::RunZCopyMode(const u32 rank, const u32 rankSize, const std::vector<LINK> &links)
116 : {
117 : // 从stream wait, 主stream record
118 0 : for (u32 i = 0; i < rankSize - 2; i++) { // 从stream 个数 = ranksize -2
119 0 : CHK_RET(LocalNotify::Wait((*subStreamsPtr_)[i], dispatcher_, (*meshSignalMainToSubPtr_)[i],
120 : INVALID_VALUE_STAGE));
121 0 : CHK_RET(LocalNotify::Post(*mainStreamPtr_, dispatcher_, (*meshSignalMainToSubPtr_)[i],
122 : INVALID_VALUE_STAGE));
123 : }
124 :
125 0 : for (u32 i = 1; i < rankSize; i++) {
126 0 : u32 destRank = (rank + i) % rankSize;
127 0 : shared_ptr<Transport> destTransport = links[destRank];
128 0 : Stream ¤tStream = (i == 1) ? *mainStreamPtr_ : (*subStreamsPtr_)[i - 2];
129 :
130 0 : u32 sendDataNum = sendAddrInfo_[destRank].size();
131 0 : vector<TxMemoryInfo> txMems(sendDataNum);
132 0 : u32 recvDataNum = recvAddrInfo_[destRank].size();
133 0 : vector<RxMemoryInfo> rxMems(recvDataNum);
134 :
135 0 : BuildSendRecvMemoryInfo(txMems, rxMems, destRank);
136 0 : CHK_RET(LoadTask(destTransport, currentStream, txMems, rxMems));
137 0 : }
138 : // 主stream wait, 从stream record
139 0 : for (u32 i = 0; i < rankSize - 2; i++) { // 从stream 个数 = ranksize -2
140 0 : CHK_RET(LocalNotify::Wait(*mainStreamPtr_, dispatcher_, (*meshSignalSubToMainPtr_)[i],
141 : INVALID_VALUE_STAGE));
142 0 : CHK_RET(LocalNotify::Post((*subStreamsPtr_)[i], dispatcher_, (*meshSignalSubToMainPtr_)[i],
143 : INVALID_VALUE_STAGE));
144 : }
145 :
146 0 : CHK_RET(AlgTemplateBase::ExecEmptyTask(sendMem_, recvMem_, *mainStreamPtr_, dispatcher_));
147 0 : return HCCL_SUCCESS;
148 : }
149 :
150 0 : HcclResult AlltoAllVStagedMesh::LoadTask(shared_ptr<Transport> destTransport, Stream ¤tStream,
151 : vector<TxMemoryInfo> &txMems, vector<RxMemoryInfo> &rxMems) const
152 : {
153 0 : CHK_RET(destTransport->TxAck(currentStream)); // record send done
154 0 : CHK_RET(destTransport->RxAck(currentStream)); // wait send done
155 0 : CHK_RET(destTransport->TxAsync(txMems, currentStream)); // record Send Ready
156 0 : CHK_RET(destTransport->RxAsync(rxMems, currentStream)); // wait Send Ready + SDMA get
157 0 : CHK_RET(destTransport->TxAck(currentStream)); // record send done
158 0 : CHK_RET(destTransport->RxAck(currentStream)); // wait send done
159 0 : CHK_RET(destTransport->TxDataSignal(currentStream)); // record Send Ready
160 0 : CHK_RET(destTransport->RxDataSignal(currentStream)); // wait send Ready
161 0 : return HCCL_SUCCESS;
162 : }
163 : REGISTER_TEMPLATE(TemplateType::TEMPLATE_ALL_2_ALL_V_STAGED_MESH, AlltoAllVStagedMesh);
164 : } // namespace hccl
|