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