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