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 "log.h"
12 :
13 : #include "alg_data_trans_wrapper.h"
14 : #include "ins_alg_template/ins_temp_all_reduce_mesh_1D_one_shot.h"
15 :
16 : namespace Hccl {
17 0 : InsTempAllReduceMesh1DOneShot::InsTempAllReduceMesh1DOneShot(
18 : const RankId virtualRank, const u32 tempRankSize, const std::vector<std::vector<RankId>>& tempVTopo,
19 0 : const std::map<RankId, u32>& tempVirtRankMap)
20 0 : : InsAlgTemplateBase(virtualRank, tempRankSize, tempVTopo, tempVirtRankMap)
21 0 : {}
22 :
23 0 : InsTempAllReduceMesh1DOneShot::~InsTempAllReduceMesh1DOneShot() {}
24 :
25 0 : u32 InsTempAllReduceMesh1DOneShot::CalcScratchMultiple(BufferType input, BufferType output) const
26 : {
27 : (void)input;
28 : (void)output;
29 : // one shot 场景,scratch Buffer 需要是 usrIn的rankSize倍
30 0 : return tempRankSize_;
31 : }
32 :
33 0 : HcclResult InsTempAllReduceMesh1DOneShot::CalcRes(AlgTempResReq& tempResReq)
34 : {
35 0 : tempResReq.queNum = tempVTopo_[0].size();
36 0 : tempResReq.streamNum = tempResReq.queNum;
37 0 : tempResReq.queNotifys = CreateMasterSlaveQueNotifiesRequest(tempResReq.queNum);
38 :
39 0 : QId centerQ = 0;
40 0 : tempResReq.localWaitGroupCntNotify.emplace_back(centerQ, 0);
41 0 : tempResReq.localBcastPostCntNotify.emplace_back(centerQ, 0);
42 :
43 0 : CHK_PRT_RET(
44 : CalcResLinksMesh(myRank_, tempRankSize_, tempVTopo_, linkNumBtwPeers_, tempResReq) != HcclResult::HCCL_SUCCESS,
45 : HCCL_ERROR("[CollAlgFactory] [InsTempAllReduceMesh1DOneShot] Rank [%d], resLinks calculation error!", myRank_),
46 : HcclResult::HCCL_E_INTERNAL);
47 :
48 0 : return HcclResult::HCCL_SUCCESS;
49 : }
50 :
51 0 : HcclResult InsTempAllReduceMesh1DOneShot::CalcSlice(const u64 dataSize, RankSliceInfo& sliceInfoVec)
52 : {
53 0 : std::vector<SliceInfo> tmp(tempVTopo_.size());
54 0 : sliceInfoVec.resize(tempRankSize_, tmp);
55 0 : AllignInfo allignInfo = {false, 0, dataType_}; // 参数填充,CalcRsAgSliceInfoMesh实际没用到
56 :
57 0 : CHK_RET(CalcRsAgSliceInfoMesh(myRank_, tempRankSize_, allignInfo, dataSize, sliceInfoVec));
58 :
59 0 : return HcclResult::HCCL_SUCCESS;
60 0 : }
61 :
62 0 : HcclResult InsTempAllReduceMesh1DOneShot::GenExtIns(
63 : const TempFuncs& tempFuncs, const TemplateDataParams& tempAlgParams, const ResLinks& tempLinks,
64 : std::vector<InsQuePtr>& tempInsQues)
65 : {
66 0 : HCCL_INFO("[InsTempAllReduceMesh1DOneShot][Run] AllReduceMesh1DOneShot begin: rank[%d] start", myRank_);
67 :
68 0 : opMode_ = tempFuncs.opMode;
69 0 : queNum_ = tempVTopo_[0].size();
70 0 : CHK_PRT_RET(
71 : queNum_ != tempInsQues.size(),
72 : HCCL_ERROR("[CollAlgFactory] [InsTempAllReduceMesh1DOneShot] Rank [%d], requiredQue Error.", myRank_),
73 : HcclResult::HCCL_E_INTERNAL);
74 :
75 0 : RankSliceInfo sliceInfoVec;
76 0 : CHK_RET(CalcSlice(tempAlgParams.sliceSize, sliceInfoVec));
77 :
78 0 : HCCL_INFO("[InsTempAllReduceMesh1DOneShot][PreCopy] write userIn data directly to the ScratchBuffer, skip precopy");
79 0 : CHK_RET(RunAllReduce(tempAlgParams, sliceInfoVec, tempLinks, tempInsQues));
80 0 : HCCL_INFO("[InsTempAllReduceMesh1DOneShot][PostCopy] data is already in the userOut, skip postcopy");
81 :
82 0 : HCCL_INFO("[InsTempAllReduceMesh1DOneShot][Run] AllReduceMesh1DOneShot finished: rank[%d] end", myRank_);
83 0 : return HcclResult::HCCL_SUCCESS;
84 0 : }
85 :
86 0 : HcclResult InsTempAllReduceMesh1DOneShot::RunAllReduce(
87 : const TemplateDataParams& tempAlgParams, const RankSliceInfo& sliceInfoVec, const ResLinks& tempLinks,
88 : std::vector<InsQuePtr>& tempInsQues)
89 : {
90 0 : HCCL_INFO("[InsTempAllReduceMesh1DOneShot][RunAllReduce] send/recv: rank[%d]", myRank_);
91 :
92 : // semaphore sync
93 0 : if (tempVTopo_[0].size() > 1) {
94 0 : CHK_RET(PreSyncInterQueues(tempInsQues));
95 : }
96 :
97 0 : DataSlice usrInSlices = DataSlice(BufferType::INPUT, tempAlgParams.buffInfo.inBuffBaseOff, tempAlgParams.sliceSize);
98 : DataSlice usrOutSlices
99 0 : = DataSlice(BufferType::OUTPUT, tempAlgParams.buffInfo.outBuffBaseOff, tempAlgParams.sliceSize);
100 :
101 : // 主流动作
102 0 : CHK_RET(LocalCopy(tempInsQues[0], usrInSlices, usrOutSlices));
103 :
104 : // 从流动作
105 0 : for (u32 queIdx = 1; queIdx < queNum_; queIdx++) {
106 0 : u32 nextRank = (myRank_ + queIdx) % tempRankSize_; // 让rank和que对应上
107 0 : RankId fromRank = nextRank;
108 0 : RankId toRank = nextRank;
109 :
110 0 : const std::vector<LinkData>& linkRecv = tempLinks.at(fromRank);
111 0 : const std::vector<LinkData>& linkSend = tempLinks.at(toRank);
112 :
113 0 : std::vector<DataSlice> txSrcSlices;
114 0 : std::vector<DataSlice> txDstSlices;
115 :
116 0 : u64 txDstOffset = sliceInfoVec[myRank_][0].offset + tempAlgParams.buffInfo.scratchBuffBaseOff;
117 0 : u64 txDstSize = sliceInfoVec[myRank_][0].size;
118 0 : DataSlice txSrcSlice = usrInSlices;
119 0 : DataSlice txDstSlice = DataSlice(BufferType::SCRATCH, txDstOffset, txDstSize);
120 0 : txSrcSlices.push_back(txSrcSlice);
121 0 : txDstSlices.push_back(txDstSlice);
122 :
123 0 : std::vector<DataSlice> rxSrcSlices;
124 0 : std::vector<DataSlice> rxDstSlices;
125 0 : u64 rxDstOffset = sliceInfoVec[fromRank][0].offset + tempAlgParams.buffInfo.scratchBuffBaseOff;
126 0 : u64 rxDstSize = sliceInfoVec[fromRank][0].size;
127 0 : DataSlice rxSrcSlice = usrInSlices;
128 0 : DataSlice rxDstSlice = DataSlice(BufferType::SCRATCH, rxDstOffset, rxDstSize);
129 0 : rxSrcSlices.push_back(rxSrcSlice);
130 0 : rxDstSlices.push_back(rxDstSlice);
131 :
132 0 : TxRxLinks sendRecvLinks(linkSend[0], linkRecv[0]);
133 0 : TxRxSlicesList sendRecvSlicesList({txSrcSlices, txDstSlices}, {rxSrcSlices, rxDstSlices});
134 :
135 0 : SendRecvInfo sendRecvInfo(sendRecvLinks, sendRecvSlicesList);
136 0 : CHK_PRT_RET(
137 : SendRecv(sendRecvInfo, tempInsQues[queIdx], 0, true, DmaMode::PUT),
138 : HCCL_ERROR("[InsTempAllReduceMesh1DOneShot] RunAllReduce SendRecv failed"), HcclResult::HCCL_E_INTERNAL);
139 0 : }
140 :
141 : // semaphore sync
142 0 : if (tempVTopo_[0].size() > 1) {
143 0 : CHK_RET(PostSyncInterQueues(tempInsQues));
144 : }
145 :
146 0 : HCCL_INFO("[InsTempAllReduceMesh1DOneShot][RunAllReduce] reduce: rank[%d]", myRank_);
147 :
148 0 : for (u32 rankIdx = 0; rankIdx < tempVTopo_[0].size(); rankIdx++) {
149 0 : RankId myRank = myRank_;
150 0 : RankId curRank = rankIdx;
151 0 : if (curRank == myRank) {
152 0 : continue;
153 : }
154 :
155 0 : u64 curSrcOffset = sliceInfoVec[curRank][0].offset + tempAlgParams.buffInfo.scratchBuffBaseOff;
156 0 : u64 curSrcSize = sliceInfoVec[curRank][0].size;
157 0 : DataSlice curSrcSlice = DataSlice(BufferType::SCRATCH, curSrcOffset, curSrcSize);
158 0 : DataSlice curDstSlice = usrOutSlices;
159 :
160 0 : CHK_RET(LocalReduce(tempInsQues[0], curSrcSlice, curDstSlice, dataType_, redOp_));
161 : }
162 :
163 0 : return HcclResult::HCCL_SUCCESS;
164 : }
165 :
166 : } // namespace Hccl
|