LCOV - code coverage report
Current view: top level - legacy/ascend910/framework/device/aicpu_kfc/algorithm - aicpu_allreduce.cc (source / functions) Coverage Total Hit
Test: coverage.info Lines: 58.6 % 444 260
Test Date: 2026-08-18 17:47:01 Functions: 73.3 % 30 22

            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 "aicpu_allreduce.h"
      12              : #include <cmath>
      13              : #include <algorithm>
      14              : #include "common/aicpu_hccl_common.h"
      15              : 
      16           41 : HcclResult AicpuAllreduce::RunAlgorithm(
      17              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, u64 dataCount, HcclDataType dataType,
      18              :     [[maybe_unused]] u64 strideLen, AivAicpuOpParam* /* nextTask */)
      19              : {
      20           41 :     CHK_PTR_NULL(ctx_);
      21           41 :     if (CC_EXE_ONE_SHOT_8_STREAM == ctx_->commOpType) {
      22           34 :         return RunAllReduceReduceBcast(opType, sendBuffer, recvBuffer, dataCount * ctx_->unitSize, dataType);
      23            7 :     } else if (ctx_->commOpType == CC_EXE_ONE_SHOT_1_STREAM) {
      24            1 :         return RunAllReduceOneShot1Stream(opType, sendBuffer, recvBuffer, dataCount * ctx_->unitSize, dataType);
      25            6 :     } else if (ctx_->commOpType == CC_EXE_TWO_SHOT_1_STREAM) {
      26            1 :         return RunAllReduceTwoShot1Stream(opType, sendBuffer, recvBuffer, dataCount, dataType);
      27            5 :     } else if (ctx_->commOpType == CC_EXE_ONE_SHOT_HD) {
      28            1 :         return RunAllReduceOneshotHD(opType, sendBuffer, recvBuffer, dataCount * ctx_->unitSize, dataType);
      29            4 :     } else if (ctx_->commOpType == CC_EXE_ONE_SHOT_SINGLE_RING) {
      30            2 :         return RunAllReduceRing(opType, sendBuffer, recvBuffer, dataCount, dataType);
      31              :     }
      32              : 
      33            2 :     if (ctx_->useBufferType == MC2_BUFFER_TYPE_WINDOW_IN) {
      34            1 :         if (HCCL_SUCCESS == RunAllReduceAlignWin2Win(opType, recvBuffer, dataCount, dataType)) {
      35            0 :             return HCCL_SUCCESS;
      36              :         }
      37              :     } else {
      38            1 :         if (HCCL_SUCCESS == RunAllReduceAlign(opType, sendBuffer, recvBuffer, dataCount, dataType)) {
      39            0 :             return HCCL_SUCCESS;
      40              :         }
      41              :     }
      42              : 
      43            2 :     return RunAllReduce(opType, sendBuffer, recvBuffer, dataCount, dataType);
      44              : }
      45              : 
      46            1 : int64_t AicpuAllreduce::RoundUpWithDivisor(u64 value, u64 divisor) const
      47              : {
      48            1 :     if ((value == 0) || (divisor == 0)) {
      49            0 :         return divisor;
      50              :     }
      51              :     // divisor必须大于等于1, 返回value向上取divisor的整数倍的值
      52            1 :     return ((value + (divisor - 1)) / divisor) * divisor;
      53              : }
      54              : 
      55              : HcclResult
      56            1 : AicpuAllreduce::PrepareSlice(u64 dataCount, HcclDataType dataType, u32 sliceNum, std::vector<Slice>& dataSlice) const
      57              : {
      58            1 :     Slice temp;
      59            1 :     u32 unitSize = DataUnitSize(dataType);
      60            1 :     u64 totalSize = dataCount * unitSize;
      61            1 :     dataSlice.clear();
      62            1 :     dataSlice.reserve(sliceNum);
      63            1 :     if (sliceNum == 0) {
      64            0 :         HCCL_ERROR("[Prepare][SliceData]data slice prepare, sliceNum is 0.");
      65            0 :         return HCCL_E_PARA;
      66              :     }
      67            1 :     u64 sizePerSlice = (totalSize + sliceNum - 1) / sliceNum; /* 1是为了向上取整 */
      68            1 :     sizePerSlice = RoundUpWithDivisor(sizePerSlice, HCCL_MIN_SLICE_ALIGN);
      69            1 :     u64 residueSize = totalSize;
      70            1 :     u32 i = 0;
      71            3 :     while (residueSize > 0) {
      72            2 :         u64 sliceSize = sizePerSlice < residueSize ? sizePerSlice : residueSize;
      73            2 :         temp.size = sliceSize;
      74            2 :         temp.offset = totalSize - residueSize;
      75            2 :         i++;
      76            2 :         if (sliceSize <= 0) {
      77            0 :             HCCL_ERROR("[Prepare][SliceData]data_slice_prepare sliceSize[%llu]", sliceSize);
      78            0 :             return HCCL_E_PARA;
      79              :         }
      80            2 :         residueSize -= sliceSize;
      81            2 :         dataSlice.push_back(temp);
      82              :     }
      83            1 :     while (i < sliceNum) {
      84            0 :         temp.size = 0;
      85            0 :         temp.offset = totalSize;
      86            0 :         i++;
      87            0 :         dataSlice.push_back(temp);
      88              :     }
      89            1 :     return HCCL_SUCCESS;
      90              : }
      91              : 
      92            3 : void AicpuAllreduce::GetDataSizes16K(std::vector<u64>& dataSizes, u64 allDataSize) const
      93              : {
      94            3 :     u64 num16k = allDataSize / HCCL_COPY_ALIGN;
      95            3 :     u64 tailSize = allDataSize % HCCL_COPY_ALIGN;
      96              : 
      97            3 :     u64 baseNum = num16k / ctx_->rankNum;
      98            3 :     u64 tailNum = num16k % ctx_->rankNum;
      99              : 
     100            7 :     for (u32 i = 0; i < ctx_->rankNum; i++) {
     101            4 :         dataSizes[i] = baseNum * HCCL_COPY_ALIGN;
     102              :     }
     103            3 :     for (u32 i = 0; i < tailNum; i++) {
     104            0 :         dataSizes[i] += HCCL_COPY_ALIGN;
     105              :     }
     106            3 :     dataSizes[ctx_->rankNum - 1] += tailSize;
     107            3 : }
     108              : 
     109            1 : HcclResult AicpuAllreduce::RunAllReduceAlignWin2Win(
     110              :     HcclReduceOp opType, void* recvBuffer, u64 dataCount, HcclDataType dataType) const
     111              : {
     112            1 :     u32 unitSize = ctx_->unitSize;
     113            1 :     u64 allDataSize = dataCount * unitSize;
     114            1 :     u8* curOutputPtr = static_cast<u8*>(recvBuffer);
     115              : 
     116            1 :     if (ctx_->rankNum == 0) {
     117            0 :         return HCCL_E_PARA;
     118              :     }
     119            1 :     std::vector<u64> dataSizes(ctx_->rankNum, 0);
     120            1 :     GetDataSizes16K(dataSizes, allDataSize);
     121              : 
     122            2 :     if (dataSizes[0] > ctx_->windowSize || dataSizes[ctx_->rankNum - 1] > ctx_->windowSize
     123            2 :         || (ctx_->rankNum - 1) * HCCL_COPY_ALIGN >= allDataSize) {
     124            1 :         return HCCL_E_PARA;
     125              :     }
     126              : 
     127            0 :     std::vector<u64> dataOffsets(ctx_->rankNum, 0);
     128            0 :     for (u32 i = 1; i < ctx_->rankNum; i++) {
     129            0 :         dataOffsets[i] = dataOffsets[i - 1] + dataSizes[i - 1];
     130              :     }
     131              : 
     132            0 :     u64 winOffset = ctx_->winOffset;
     133              : 
     134              :     // 1. 前同步
     135            0 :     TaskOrchestrator::DoPreSync();
     136              : 
     137              :     // 2. 跨片SDMA,分批拷贝 + 分批结束同步
     138            0 :     TaskOrchestrator::IpcCpyWin2Win(dataSizes, winOffset, dataOffsets, opType, dataType);
     139              : 
     140              :     // 3. 后同步
     141            0 :     TaskOrchestrator::DoPostSync();
     142              : 
     143              :     // 4. 前同步
     144            0 :     TaskOrchestrator::DoPreSync();
     145              : 
     146              :     // 5. 片内数据 Win拷贝到Rcv
     147            0 :     TaskOrchestrator::SelfCpyWin2RcvEx1(
     148            0 :         curOutputPtr, dataSizes[ctx_->rankId], dataOffsets[ctx_->rankId], winOffset, HCCL_REDUCE_RESERVED, dataType);
     149              : 
     150              :     // 6. 跨片SDMA,分批拷贝 + 分批结束同步
     151            0 :     TaskOrchestrator::IpcCpyWin2RcvEx(curOutputPtr, dataSizes, dataOffsets, winOffset, HCCL_REDUCE_RESERVED, dataType);
     152              : 
     153              :     // 7. 后同步
     154            0 :     TaskOrchestrator::DoPostSync();
     155              : 
     156            0 :     TaskOrchestrator::LaunchTasks();
     157              : 
     158            0 :     return HCCL_SUCCESS;
     159            1 : }
     160              : 
     161            1 : HcclResult AicpuAllreduce::RunAllReduceAlign(
     162              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, u64 dataCount, HcclDataType dataType) const
     163              : {
     164            1 :     u32 unitSize = ctx_->unitSize;
     165            1 :     u64 allDataSize = dataCount * unitSize;
     166              : 
     167            1 :     u8* curInputPtr = static_cast<u8*>(sendBuffer);
     168            1 :     u8* curOutputPtr = static_cast<u8*>(recvBuffer);
     169              : 
     170            1 :     if (ctx_->rankNum == 0) {
     171            0 :         return HCCL_E_PARA;
     172              :     }
     173            1 :     std::vector<u64> dataSizes(ctx_->rankNum, 0);
     174            1 :     GetDataSizes16K(dataSizes, allDataSize);
     175              : 
     176            2 :     if (dataSizes[0] > ctx_->windowSize || dataSizes[ctx_->rankNum - 1] > ctx_->windowSize
     177            2 :         || (ctx_->rankNum - 1) * HCCL_COPY_ALIGN >= allDataSize) {
     178            1 :         return HCCL_E_PARA;
     179              :     }
     180              : 
     181            0 :     std::vector<u64> dataOffsets(ctx_->rankNum, 0);
     182            0 :     for (u32 i = 1; i < ctx_->rankNum; i++) {
     183            0 :         dataOffsets[i] = dataOffsets[i - 1] + dataSizes[i - 1];
     184              :     }
     185              : 
     186              :     // 1. 片内数据 Snd拷贝到Window
     187            0 :     TaskOrchestrator::SelfCpySnd2Win(
     188            0 :         curInputPtr, dataSizes[ctx_->rankId], dataOffsets[ctx_->rankId], 0, HCCL_REDUCE_RESERVED, dataType);
     189              : 
     190              :     // 2. 前同步
     191            0 :     TaskOrchestrator::DoPreSync();
     192              : 
     193              :     // 3. 跨片SDMA,分批拷贝 + 分批结束同步
     194            0 :     TaskOrchestrator::IpcCpySnd2Win(curInputPtr, dataSizes, dataOffsets, nullptr, opType, dataType);
     195              : 
     196              :     // 4. 后同步
     197            0 :     TaskOrchestrator::DoPostSync();
     198              : 
     199              :     // 5. 前同步
     200            0 :     TaskOrchestrator::DoPreSync();
     201              : 
     202              :     // 6. 片内数据 Win拷贝到Rcv
     203            0 :     TaskOrchestrator::SelfCpyWin2Rcv(
     204            0 :         curOutputPtr, dataSizes[ctx_->rankId], 0, dataOffsets[ctx_->rankId], HCCL_REDUCE_RESERVED, dataType);
     205              : 
     206              :     // 7. 跨片SDMA,分批拷贝 + 分批结束同步
     207            0 :     TaskOrchestrator::IpcCpyWin2Rcv(curOutputPtr, dataSizes, nullptr, dataOffsets, HCCL_REDUCE_RESERVED, dataType);
     208              : 
     209              :     // 8. 后同步
     210            0 :     TaskOrchestrator::DoPostSync();
     211              : 
     212            0 :     TaskOrchestrator::LaunchTasks();
     213              : 
     214            0 :     return HCCL_SUCCESS;
     215            1 : }
     216              : 
     217            3 : HcclResult AicpuAllreduce::RunAllReduce(
     218              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, u64 dataCount, HcclDataType dataType) const
     219              : {
     220            3 :     u64 windowSize = ctx_->windowSize; // window size default is 200M, maybe need read from cfg/env.
     221            3 :     u32 unitSize = ctx_->unitSize;
     222            3 :     u64 maxCountPerLoop = (windowSize / unitSize) * ctx_->rankNum; // 中转内存单次最多能够接受的output count
     223              : 
     224            3 :     u8* curInputPtr = static_cast<u8*>(sendBuffer);
     225            3 :     u8* curOutputPtr = static_cast<u8*>(recvBuffer);
     226            3 :     u64 inputOffset = 0;
     227            3 :     u64 outputOffset = 0;
     228            3 :     u64 countLeft = dataCount;
     229              : 
     230            3 :     u64 dataSlice[AC_MAX_RANK_NUM] = {0};
     231            3 :     u64 sliceSize[AC_MAX_RANK_NUM] = {0};
     232            3 :     if (ctx_->rankNum <= 0) {
     233            1 :         return HCCL_E_UNAVAIL;
     234              :     }
     235              : 
     236            2 :     while (countLeft > 0) {
     237            0 :         curInputPtr += inputOffset;
     238            0 :         curOutputPtr += outputOffset;
     239            0 :         u64 curCount = (countLeft > maxCountPerLoop) ? maxCountPerLoop : countLeft;
     240            0 :         u64 curSize = curCount * unitSize; // 单位 byte
     241            0 :         u64 curRankCnt = curCount / ctx_->rankNum;
     242            0 :         for (u32 i = 0; i < ctx_->rankNum; i++) {
     243            0 :             dataSlice[i] = i * curRankCnt * unitSize;
     244            0 :             sliceSize[i] = curRankCnt * unitSize;
     245              :         }
     246            0 :         sliceSize[ctx_->rankNum - 1] += (curCount - curRankCnt * ctx_->rankNum) * unitSize;
     247              : 
     248            0 :         HCCL_DEBUG(
     249              :             "RunAllReducev:curInputPtr[%p], curOutputPtr[%p], curCount[%llu], curSize[%llu]", curInputPtr, curOutputPtr,
     250              :             curCount, curSize);
     251              : 
     252            0 :         if (ctx_->useBufferType != MC2_BUFFER_TYPE_WINDOW_IN) {
     253            0 :             RunAllReduceSlice(curOutputPtr, curInputPtr, sliceSize, dataSlice, opType, dataType);
     254              :         } else {
     255            0 :             RunAllReduceSliceWin2Win(curOutputPtr, sliceSize, dataSlice, opType, dataType);
     256              :         }
     257              : 
     258            0 :         countLeft -= curCount;
     259            0 :         inputOffset = curSize;
     260            0 :         outputOffset = curSize;
     261              :     }
     262              : 
     263            2 :     return HCCL_SUCCESS;
     264              : }
     265              : 
     266            1 : void AicpuAllreduce::RunAllReduceSliceWin2Win(
     267              :     u8* curOutputPtr, u64* sliceSize, u64* dataSlice, HcclReduceOp opType, HcclDataType dataType) const
     268              : {
     269            1 :     u64 winOffset = ctx_->winOffset;
     270              :     // 1. 前同步
     271            1 :     TaskOrchestrator::DoPreSync();
     272              : 
     273              :     // 2. 跨片SDMA,分批拷贝 + 分批结束同步
     274            1 :     TaskOrchestrator::IpcCpyWin2Win(sliceSize, dataSlice, opType, winOffset, dataType);
     275              : 
     276              :     // 3. 后同步
     277            1 :     TaskOrchestrator::DoPostSync();
     278              : 
     279              :     // 4. 前同步
     280            1 :     TaskOrchestrator::DoPreSync();
     281              : 
     282              :     // 5. 片内数据 Win拷贝到Rcv
     283            1 :     TaskOrchestrator::SelfCpyWin2RcvEx1(
     284            1 :         curOutputPtr, sliceSize[ctx_->rankId], dataSlice[ctx_->rankId], winOffset, HCCL_REDUCE_RESERVED, dataType);
     285              : 
     286              :     // 6. 跨片SDMA,分批拷贝 + 分批结束同步
     287            1 :     TaskOrchestrator::IpcCpyWin2RcvEx(curOutputPtr, sliceSize, dataSlice, winOffset, HCCL_REDUCE_RESERVED, dataType);
     288              : 
     289              :     // 7. 后同步
     290            1 :     TaskOrchestrator::DoPostSync();
     291              : 
     292            1 :     TaskOrchestrator::LaunchTasks();
     293            1 : }
     294              : 
     295            1 : void AicpuAllreduce::RunAllReduceSlice(
     296              :     u8* curOutputPtr, u8* curInputPtr, u64* sliceSize, u64* dataSlice, HcclReduceOp opType, HcclDataType dataType) const
     297              : {
     298              :     // 1. 片内数据 Snd拷贝到Window
     299            1 :     TaskOrchestrator::SelfCpySnd2Win(
     300            1 :         curInputPtr, sliceSize[ctx_->rankId], dataSlice[ctx_->rankId], 0, HCCL_REDUCE_RESERVED, dataType);
     301              : 
     302              :     // 2. 前同步
     303            1 :     TaskOrchestrator::DoPreSync();
     304              : 
     305              :     // 3. 跨片SDMA,分批拷贝 + 分批结束同步
     306            1 :     TaskOrchestrator::IpcCpySnd2Win(curInputPtr, sliceSize, dataSlice, nullptr, opType, dataType);
     307              : 
     308              :     // 4. 后同步
     309            1 :     TaskOrchestrator::DoPostSync();
     310              : 
     311              :     // 5. 前同步
     312            1 :     TaskOrchestrator::DoPreSync();
     313              : 
     314              :     // 6. 片内数据 Win拷贝到Rcv
     315            1 :     TaskOrchestrator::SelfCpyWin2Rcv(
     316            1 :         curOutputPtr, sliceSize[ctx_->rankId], 0, dataSlice[ctx_->rankId], HCCL_REDUCE_RESERVED, dataType);
     317              : 
     318              :     // 7. 跨片SDMA,分批拷贝 + 分批结束同步
     319            1 :     TaskOrchestrator::IpcCpyWin2Rcv(curOutputPtr, sliceSize, nullptr, dataSlice, HCCL_REDUCE_RESERVED, dataType);
     320              : 
     321              :     // 8. 后同步
     322            1 :     TaskOrchestrator::DoPostSync();
     323              : 
     324            1 :     TaskOrchestrator::LaunchTasks();
     325            1 : }
     326              : 
     327            1 : HcclResult AicpuAllreduce::RunAllReduceOneShot4Stream(
     328              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, u64 dataSize, HcclDataType dataType) const
     329              : {
     330              :     // 第一轮第一组
     331            1 :     u32 mainRankId = ctx_->rankId;
     332            1 :     u32 maxStreamNum = ctx_->rankNum / 2;
     333            1 :     u32 startRank = 0;
     334            1 :     u32 endRank = maxStreamNum - 1;
     335              : 
     336            1 :     if (mainRankId >= maxStreamNum) {
     337            1 :         startRank = maxStreamNum;
     338            1 :         endRank = ctx_->rankNum - 1;
     339              :     }
     340              : 
     341              :     // 第1轮
     342              :     // 1. 片内数据 拷贝到Window
     343            1 :     TaskOrchestrator::SelfCpySnd2WinEx(
     344              :         mainRankId, sendBuffer, dataSize, 0, 0, HCCL_REDUCE_RESERVED, dataType, maxStreamNum);
     345            1 :     TaskOrchestrator::MainSubPreSync(mainRankId, startRank, endRank, maxStreamNum);
     346              : 
     347            1 :     TaskOrchestrator::IpcPreSyncEx(startRank, endRank, maxStreamNum, false);
     348              :     // 2. 跨片SDMA 片内Send拷贝到对端Window
     349            1 :     TaskOrchestrator::IpcCpySnd2WinEx(
     350              :         sendBuffer, dataSize, nullptr, nullptr, opType, dataType, startRank, endRank, maxStreamNum, false);
     351            1 :     TaskOrchestrator::IpcPostSyncEx(startRank, endRank, maxStreamNum, false);
     352              : 
     353            1 :     TaskOrchestrator::MainSubPostSync(mainRankId, startRank, endRank, maxStreamNum);
     354              : 
     355              :     // 第2轮
     356            1 :     u32 remoteRank = (ctx_->rankNum - 1) - mainRankId; // 0-7; 1-6; 2-5; 3-4
     357              :     // 3. 片内数据 拷贝到recv
     358            1 :     TaskOrchestrator::SelfCpyWin2RcvEx(
     359              :         mainRankId, recvBuffer, dataSize, 0, 0, HCCL_REDUCE_RESERVED, dataType, maxStreamNum);
     360            1 :     TaskOrchestrator::IpcPreSyncEx(remoteRank, remoteRank, maxStreamNum, true);
     361              :     // 4. 跨片SDMA Window拷贝到对端Recv
     362            1 :     TaskOrchestrator::IpcCpyWin2RcvEx(
     363              :         recvBuffer, dataSize, nullptr, nullptr, opType, dataType, remoteRank, remoteRank, maxStreamNum, true);
     364            1 :     TaskOrchestrator::IpcPostSyncEx(remoteRank, remoteRank, maxStreamNum, true);
     365              : 
     366              :     // 5. 下发sqe
     367            1 :     TaskOrchestrator::LaunchTasksEx(0, maxStreamNum - 1, maxStreamNum);
     368              : 
     369            1 :     return HCCL_SUCCESS;
     370              : }
     371              : 
     372           34 : HcclResult AicpuAllreduce::RunReduceBcastOnMainSq(
     373              :     u32 mainRankId, u32 maxStreamNum, u32 /* startRank */, u32 /* endRank */, void* sendBuffer, void* recvBuffer,
     374              :     u64 dataSize, HcclDataType dataType) const
     375              : {
     376           34 :     HCCL_DEBUG("run RunReduceBcastOnMainSq start");
     377           34 :     if (ctx_->useBufferType != MC2_BUFFER_TYPE_WINDOW_IN) {
     378              :         // 1. reduce
     379           34 :         TaskOrchestrator::SelfCpySnd2WinEx(
     380              :             mainRankId, sendBuffer, dataSize, 0, 0, HCCL_REDUCE_RESERVED, dataType, maxStreamNum);
     381              :     }
     382              : 
     383              :     // 2. 前同步
     384           34 :     TaskOrchestrator::MainSubPreSync();
     385           34 :     TaskOrchestrator::IpcPreRecordEx(0, maxStreamNum - 1, maxStreamNum, false);
     386              :     // 3. 后同步
     387           34 :     TaskOrchestrator::IpcPostWaitEx(0, maxStreamNum - 1, maxStreamNum, true);
     388              : 
     389              :     // 4. 前同步
     390           34 :     TaskOrchestrator::MainSubPreSync();
     391           34 :     TaskOrchestrator::IpcPreRecordEx(0, maxStreamNum - 1, maxStreamNum, false);
     392              :     // 5. bcast
     393           34 :     TaskOrchestrator::SelfCpyWin2RcvEx(
     394           34 :         mainRankId, recvBuffer, dataSize, ctx_->useBufferType != MC2_BUFFER_TYPE_WINDOW_IN ? 0 : ctx_->winOffset, 0,
     395              :         HCCL_REDUCE_RESERVED, dataType, maxStreamNum);
     396              :     // 6. 后同步
     397           34 :     TaskOrchestrator::IpcPostWaitEx(0, maxStreamNum - 1, maxStreamNum, true);
     398              : 
     399           34 :     HCCL_DEBUG("run RunReduceBcastOnMainSq end");
     400           34 :     return HCCL_SUCCESS;
     401              : }
     402              : 
     403            0 : HcclResult AicpuAllreduce::RunReduceBcastOnOtherSq(
     404              :     HcclReduceOp opType, u32 mainRankId, u32 maxStreamNum, void* sendBuffer, void* recvBuffer, u64 dataSize,
     405              :     HcclDataType dataType) const
     406              : {
     407            0 :     HCCL_DEBUG("run RunReduceBcastOnOtherSq start");
     408              : 
     409              :     // reduce
     410            0 :     TaskOrchestrator::IpcPreWaitEx(mainRankId, mainRankId, maxStreamNum, true);
     411            0 :     if (ctx_->useBufferType != MC2_BUFFER_TYPE_WINDOW_IN) {
     412            0 :         TaskOrchestrator::SelfCpySnd2WinEx(mainRankId, sendBuffer, dataSize, 0, 0, opType, dataType, maxStreamNum);
     413              :     } else {
     414            0 :         TaskOrchestrator::IpcCpyWin2WinEx(mainRankId, dataSize, ctx_->winOffset, opType, dataType, maxStreamNum);
     415              :     }
     416            0 :     TaskOrchestrator::IpcPostRecordEx(mainRankId, mainRankId, maxStreamNum, true);
     417              : 
     418              :     // bcast
     419            0 :     TaskOrchestrator::IpcPreWaitEx(mainRankId, mainRankId, maxStreamNum, true);
     420            0 :     TaskOrchestrator::SelfCpyWin2RcvEx(
     421            0 :         mainRankId, recvBuffer, dataSize, ctx_->useBufferType != MC2_BUFFER_TYPE_WINDOW_IN ? 0 : ctx_->winOffset, 0,
     422              :         HCCL_REDUCE_RESERVED, dataType, maxStreamNum);
     423            0 :     TaskOrchestrator::IpcPostRecordEx(mainRankId, mainRankId, maxStreamNum, true);
     424            0 :     HCCL_DEBUG("run RunReduceBcastOnOtherSq end");
     425            0 :     return HCCL_SUCCESS;
     426              : }
     427              : 
     428           34 : HcclResult AicpuAllreduce::RunAllReduceReduceBcast(
     429              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, u64 dataSize, HcclDataType dataType) const
     430              : {
     431           34 :     HCCL_DEBUG("run RunAllReduceReduceBcast start");
     432              : 
     433           34 :     u32 mainRankId = 0;
     434           34 :     u32 maxStreamNum = ctx_->rankNum;
     435           34 :     u32 startRank = 0;
     436           34 :     u32 endRank = ctx_->rankNum - 1;
     437              : 
     438           34 :     if (ctx_->rankId == mainRankId) {
     439           34 :         RunReduceBcastOnMainSq(
     440              :             mainRankId, maxStreamNum, startRank, endRank, sendBuffer, recvBuffer, dataSize, dataType);
     441              :     } else {
     442            0 :         RunReduceBcastOnOtherSq(opType, mainRankId, maxStreamNum, sendBuffer, recvBuffer, dataSize, dataType);
     443              :     }
     444              : 
     445              :     // 下发sqe
     446           34 :     TaskOrchestrator::LaunchTasksEx(0, maxStreamNum - 1, maxStreamNum);
     447              : 
     448           34 :     HCCL_DEBUG("run RunAllReduceReduceBcast end");
     449           34 :     return HCCL_SUCCESS;
     450              : }
     451              : 
     452            1 : HcclResult AicpuAllreduce::RunAllReduceOneShot1Stream(
     453              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, u64 dataSize, HcclDataType dataType) const
     454              : {
     455            1 :     HCCL_INFO("run RunAllReduceOneShot1Stream start");
     456            1 :     u32 maxStreamNum = ctx_->rankNum;
     457            1 :     u8* curOutputPtr = static_cast<u8*>(recvBuffer);
     458            1 :     u32 startRank = 0;
     459            1 :     u32 endRank = maxStreamNum - 1;
     460              : 
     461            1 :     TaskOrchestrator::SelfCpySnd2WinEx1(sendBuffer, dataSize, 0, 0, HCCL_REDUCE_RESERVED, dataType, maxStreamNum);
     462              : 
     463            1 :     TaskOrchestrator::SelfCpySnd2RcvEx(sendBuffer, recvBuffer, 0, 0, dataSize, HCCL_REDUCE_RESERVED, dataType);
     464              : 
     465            1 :     TaskOrchestrator::IpcPreSyncEx(startRank, endRank, maxStreamNum, true);
     466              : 
     467            1 :     TaskOrchestrator::IpcCpyWin2RcvEx(
     468              :         curOutputPtr, dataSize, nullptr, nullptr, opType, dataType, startRank, endRank, maxStreamNum, true);
     469              : 
     470            1 :     TaskOrchestrator::IpcPostSyncEx(startRank, endRank, maxStreamNum, true);
     471              : 
     472              :     // 下发sqe
     473            1 :     TaskOrchestrator::LaunchTasksEx(0, maxStreamNum - 1, maxStreamNum);
     474              : 
     475            1 :     HCCL_INFO("run RunAllReduceOneShot1Stream end");
     476            1 :     return HCCL_SUCCESS;
     477              : }
     478              : 
     479            2 : HcclResult AicpuAllreduce::RunAllReduceTwoShot1Stream(
     480              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, u64 dataCount, HcclDataType dataType) const
     481              : {
     482            2 :     HCCL_INFO("run RunAllReduceTwoShot1Stream start");
     483            2 :     u32 maxStreamNum = ctx_->rankNum;
     484            2 :     u32 startRank = 0;
     485            2 :     u32 endRank = maxStreamNum - 1;
     486            2 :     u8* curInputPtr = static_cast<u8*>(sendBuffer);
     487            2 :     u8* curOutputPtr = static_cast<u8*>(recvBuffer);
     488            2 :     u32 unitSize = ctx_->unitSize;
     489            2 :     u64 inputOffset = 0;
     490            2 :     u64 outputOffset = 0;
     491            2 :     u64 windowSize = ctx_->windowSize; // window size default is 200M, maybe need read from cfg/env.
     492            2 :     u64 maxCountPerLoop = windowSize / unitSize * ctx_->rankNum; // 中转内存单次最多能够接受的output count
     493            2 :     u64 countLeft = dataCount;
     494              : 
     495            2 :     while (countLeft > 0) {
     496            0 :         curInputPtr += inputOffset;
     497            0 :         curOutputPtr += outputOffset;
     498            0 :         u64 curCount = (countLeft > maxCountPerLoop) ? maxCountPerLoop : countLeft;
     499            0 :         u64 curSize = curCount * unitSize; // 单位 byte
     500            0 :         std::vector<Slice> dataSlice;
     501            0 :         PrepareSlice(curCount, dataType, maxStreamNum, dataSlice);
     502            0 :         u64 sliceSize = dataSlice[ctx_->rankId].size;
     503              : 
     504            0 :         HCCL_INFO(
     505              :             "RunAllReducev:curInputPtr[%p], curOutputPtr[%p], curCount[%llu], curSize[%llu]", curInputPtr, curOutputPtr,
     506              :             curCount, curSize);
     507              : 
     508            0 :         TaskOrchestrator::SelfCpySnd2WinEx1(
     509            0 :             curInputPtr, sliceSize, dataSlice[ctx_->rankId].offset, 0, HCCL_REDUCE_RESERVED, dataType, maxStreamNum);
     510              : 
     511            0 :         TaskOrchestrator::IpcPreSyncEx(startRank, endRank, maxStreamNum, true);
     512              : 
     513            0 :         TaskOrchestrator::IpcCpySnd2WinSliceEx(
     514              :             curInputPtr, dataSlice, nullptr, opType, dataType, startRank, endRank, maxStreamNum, true);
     515              : 
     516            0 :         TaskOrchestrator::IpcCpyWin2RcvSliceEx(
     517              :             curOutputPtr, dataSlice, nullptr, HCCL_REDUCE_RESERVED, dataType, startRank, endRank, maxStreamNum, true);
     518              : 
     519            0 :         TaskOrchestrator::IpcPostSyncEx(startRank, endRank, maxStreamNum, true);
     520              : 
     521            0 :         TaskOrchestrator::SelfCpyWin2Rcv(
     522            0 :             curOutputPtr, sliceSize, 0, dataSlice[ctx_->rankId].offset, HCCL_REDUCE_RESERVED, dataType);
     523              : 
     524              :         // 下发sqe
     525            0 :         TaskOrchestrator::LaunchTasksEx(0, maxStreamNum - 1, maxStreamNum);
     526              : 
     527            0 :         countLeft -= curCount;
     528            0 :         inputOffset = curSize;
     529            0 :         outputOffset = curSize;
     530            0 :     }
     531              : 
     532            2 :     HCCL_INFO("run RunAllReduceTwoShot1Stream end");
     533            2 :     return HCCL_SUCCESS;
     534              : }
     535              : 
     536              : // 计算HD算法给定轮中给定rank的对端的rank号
     537            1 : u32 AicpuAllreduce::GetHdPeer(const u32 hdRound, const u32 curRank) const
     538              : {
     539              :     // 将所有的设备分成若干组,相邻的组之间的对位节点相互通信
     540              :     // 每个组中的设备数量为 2^当前轮次,即:1,2,4,8 ....
     541            1 :     u32 groupSize = std::pow(2, hdRound);
     542              :     // 获取当前节点在组内的位置 (同时也是对端节点在组内的位置)
     543            1 :     u32 rankOffset = curRank % groupSize;
     544              :     // 获取当前节点在第几组
     545            1 :     u32 curGroupIdx = curRank / groupSize;
     546              :     // 偶数号组内的节点与下一组的节点通信,对应的,奇数号组内的节点与上一组的节点通信
     547            1 :     u32 peerGroupIdx = (curGroupIdx % 2 == 0) ? curGroupIdx + 1 : curGroupIdx - 1;
     548              : 
     549            1 :     u32 peerRank = peerGroupIdx * groupSize + rankOffset;
     550            1 :     return peerRank;
     551              : }
     552              : 
     553              : // OneshotHD 算法, 使用了DMA消减
     554              : // 只支持卡数大于 2 且为 2 的幂数的场景
     555              : // datasize 需小于 ccl buffer 大小
     556            1 : HcclResult AicpuAllreduce::RunAllReduceOneshotHD(
     557              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, u64 dataSize, HcclDataType dataType) const
     558              : {
     559              :     /* 分3个阶段:
     560              :      * 第一阶段:将 input 的数据拷贝到 window 上
     561              :      * 第二阶段:将 input 数据发送到对端 window 进行 reduce 并将 reduce 完的
     562              :      * window 拷贝到 output
     563              :      * 第三阶段:不断将对端 window 数据读取到本端 output,
     564              :      * 如果不是最后一轮,则将 reduce 完的 output 拷贝到 window
     565              :      */
     566            1 :     HCCL_INFO("run RunAllReduceOneshotHD start");
     567              : 
     568            1 :     u8* curOutputPtr = static_cast<u8*>(recvBuffer);
     569            1 :     const u32 curRank = ctx_->rankId;
     570              :     // 第一阶段:
     571              :     // 将输入数据拷贝到 Window
     572            1 :     CHK_RET(TaskOrchestrator::SelfCpySnd2Win(sendBuffer, dataSize, 0, 0, HCCL_REDUCE_RESERVED, dataType));
     573              :     // 第二阶段:
     574              :     // 片内拷贝 片内Send拷贝到对端Window - 卡内双 die 间 allreduce
     575            1 :     u32 hdRound = 0;
     576            1 :     u32 peerRank = GetHdPeer(hdRound, curRank);
     577            1 :     CHK_RET(TaskOrchestrator::IpcPreSyncEx(peerRank, peerRank, ctx_->rankNum, true));
     578              : 
     579              :     // 将输入数据写到对端
     580            1 :     CHK_RET(TaskOrchestrator::IpcCpySnd2WinP2P(sendBuffer, peerRank, dataSize, 0, 0, opType, dataType));
     581            1 :     TaskOrchestrator::IpcPostSyncEx(peerRank, peerRank, ctx_->rankNum, true);
     582              : 
     583              :     // 片内拷贝 Win拷贝到Rcv
     584            1 :     CHK_RET(TaskOrchestrator::SelfCpyWin2Rcv(curOutputPtr, dataSize, 0, 0, HCCL_REDUCE_RESERVED, dataType));
     585              :     // 第三阶段:
     586              :     // 循环 log2(rankNum) - 1 次
     587            1 :     u32 remainingHdRounds = ctx_->rankNum >> 1;
     588            1 :     while (remainingHdRounds >>= 1) { // 使用位移代替 log2
     589            0 :         hdRound++;
     590            0 :         peerRank = GetHdPeer(hdRound, curRank);
     591              :         // 跨片拷贝 对端 window 拷贝到 rcv
     592            0 :         CHK_RET(TaskOrchestrator::IpcPreSyncEx(peerRank, peerRank, ctx_->rankNum, true));
     593            0 :         CHK_RET(TaskOrchestrator::IpcCpyWin2RcvP2PMainStream(curOutputPtr, peerRank, dataSize, 0, 0, opType, dataType));
     594            0 :         CHK_RET(TaskOrchestrator::IpcPostSyncEx(peerRank, peerRank, ctx_->rankNum, true));
     595              :         // 如果这不是最后一轮,则需要将rcv里的数据同步到 win 里
     596            0 :         if (remainingHdRounds > 1) {
     597            0 :             CHK_RET(TaskOrchestrator::SelfCpyRcv2Win(curOutputPtr, dataSize, 0, 0, HCCL_REDUCE_RESERVED, dataType));
     598              :         }
     599              :     }
     600              :     // 下发sqe
     601            1 :     TaskOrchestrator::LaunchTasksEx(0, ctx_->rankNum - 1, ctx_->rankNum);
     602            1 :     HCCL_INFO("run RunAllReduceOneshotHD end");
     603            1 :     return HCCL_SUCCESS;
     604              : }
     605              : 
     606              : // 将 oriValue 对调整为 alignValue 的倍数
     607            2 : u64 AicpuAllreduce::AlignWith(u64 oriValue, u64 alignValue) const
     608              : {
     609            2 :     if (oriValue <= alignValue || alignValue == 0) {
     610            2 :         return oriValue;
     611              :     }
     612            0 :     u64 remain = oriValue % alignValue;
     613            0 :     return oriValue - remain;
     614              : }
     615              : 
     616              : // 按照 cclBuffer 大小将数据切分
     617            2 : HcclResult AicpuAllreduce::GetBurstDataCounts(u64 windowSize, u64 dataCount, std::vector<u64>& burstDataCounts) const
     618              : {
     619            2 :     u64 alignedWindowSize = AlignWith(windowSize, HCCL_COPY_ALIGN);
     620            2 :     u32 unitSize = ctx_->unitSize;
     621            2 :     CHK_PRT_RET(unitSize == 0, HCCL_ERROR("UnitSize is 0"), HCCL_E_UNAVAIL);
     622            2 :     u32 maxDataPerBurst = alignedWindowSize / unitSize;
     623            2 :     CHK_PRT_RET(maxDataPerBurst == 0, HCCL_ERROR("maxDataPerBurst is 0"), HCCL_E_UNAVAIL);
     624            0 :     burstDataCounts.insert(burstDataCounts.end(), dataCount / maxDataPerBurst, maxDataPerBurst);
     625            0 :     u64 tailSize = dataCount % maxDataPerBurst;
     626            0 :     if (tailSize != 0) {
     627            0 :         burstDataCounts.push_back(tailSize);
     628              :     }
     629            0 :     return HCCL_SUCCESS;
     630              : }
     631              : 
     632            1 : std::vector<std::vector<u32>> AicpuAllreduce::GetRingOrders() const
     633              : {
     634              :     // simple ring
     635            1 :     std::vector<std::vector<u32>> ringOrders;
     636            1 :     std::vector<u32> ringOrder(ctx_->rankNum);
     637            3 :     for (u32 i = 0; i < ctx_->rankNum; i++) {
     638            2 :         ringOrder[i] = i;
     639              :     }
     640            1 :     ringOrders.push_back(ringOrder);
     641            1 :     return ringOrders;
     642            1 : }
     643              : 
     644              : // 当前只支持从0开始的rankID
     645            1 : HcclResult AicpuAllreduce::reorderRingSlice(
     646              :     const std::vector<u32>& ringOrder, const std::vector<Slice>& ringSlices,
     647              :     std::vector<Slice>& orderedRingSlices) const
     648              : {
     649            1 :     size_t ringSize = ringOrder.size();
     650            1 :     CHK_PRT_RET(ringSize == 0, HCCL_ERROR("ringSize is 0"), HCCL_E_UNAVAIL);
     651            1 :     orderedRingSlices.resize(ringSize);
     652            3 :     for (size_t rankIdx = 0; rankIdx < ringSize; rankIdx++) {
     653            2 :         u32 currentRank = ringOrder[rankIdx];
     654            2 :         u32 previousRankIdx = (rankIdx + ringSize - 1) % ringSize;
     655            2 :         u32 previousRank = ringOrder[previousRankIdx];
     656            2 :         Slice currentSlice = ringSlices[previousRank];
     657            2 :         orderedRingSlices[currentRank] = currentSlice;
     658              :     }
     659            1 :     return HCCL_SUCCESS;
     660              : }
     661              : 
     662            1 : HcclResult AicpuAllreduce::PrepareRingSlice(
     663              :     const std::vector<std::vector<u32>>& ringOrders, u64 dataCount, HcclDataType dataType,
     664              :     std::vector<std::vector<Slice>>& orderedAllRingSlice) const
     665              : {
     666            1 :     u32 ringNum = ringOrders.size();
     667            1 :     u32 sliceNum = ctx_->rankNum * ringNum;
     668            1 :     std::vector<Slice> dataSlices;
     669            1 :     PrepareSlice(dataCount, dataType, sliceNum, dataSlices);
     670            1 :     std::vector<std::vector<Slice>> allRingSlices(ringNum);
     671            1 :     orderedAllRingSlice.resize(ringNum);
     672              : 
     673            3 :     for (size_t i = 0; i < dataSlices.size(); i++) {
     674            2 :         allRingSlices[i % ringNum].push_back(dataSlices[i]);
     675              :     }
     676            2 :     for (size_t i = 0; i < ringNum; i++) {
     677            1 :         reorderRingSlice(ringOrders[i], allRingSlices[i], orderedAllRingSlice[i]);
     678              :     }
     679            1 :     return HCCL_SUCCESS;
     680            1 : }
     681              : 
     682            0 : HcclResult AicpuAllreduce::RingIPCPreSync(const u32 stream, const u32 prevRank, const u32 nextRank) const
     683              : {
     684              :     // 通知下游
     685            0 :     CHK_RET(AicpuDispatcher::SignalRecord(stream, nextRank, AicpuDispatcher::IPC, AicpuDispatcher::PRE_SYNC));
     686              :     // 等待上游通知
     687            0 :     CHK_RET(AicpuDispatcher::SignalWait(stream, prevRank, AicpuDispatcher::IPC, AicpuDispatcher::PRE_SYNC));
     688            0 :     return HCCL_SUCCESS;
     689              : }
     690              : 
     691            0 : HcclResult AicpuAllreduce::RingIPCPostSync(const u32 stream, const u32 prevRank, const u32 nextRank) const
     692              : {
     693              :     // 回复上游通知
     694            0 :     CHK_RET(AicpuDispatcher::SignalRecord(stream, prevRank, AicpuDispatcher::IPC, AicpuDispatcher::POST_SYNC));
     695              :     // 等待下游回复,回收 notify
     696            0 :     CHK_RET(AicpuDispatcher::SignalWait(stream, nextRank, AicpuDispatcher::IPC, AicpuDispatcher::POST_SYNC));
     697            0 :     return HCCL_SUCCESS;
     698              : }
     699              : 
     700            0 : HcclResult AicpuAllreduce::GetPrevRankList(const std::vector<u32>& ringOrder, std::vector<u32>& previousRankList) const
     701              : {
     702            0 :     size_t ringSize = ringOrder.size();
     703            0 :     for (size_t rankIdx = 0; rankIdx < ringSize; rankIdx++) {
     704            0 :         u32 currRank = ringOrder[rankIdx];
     705            0 :         u32 previousRankIdx = (rankIdx + ringSize - 1) % ringSize;
     706            0 :         previousRankList[currRank] = ringOrder[previousRankIdx];
     707              :     }
     708            0 :     return HCCL_SUCCESS;
     709              : }
     710              : 
     711            0 : size_t AicpuAllreduce::FindNextRank(const std::vector<u32>& previousRankList, const u32 localRank) const
     712              : {
     713            0 :     return std::find(previousRankList.begin(), previousRankList.end(), localRank) - previousRankList.begin();
     714              : }
     715              : 
     716            0 : Slice* AicpuAllreduce::GetNextRingSlice(
     717              :     const std::vector<u32>& previousRankList, std::vector<Slice>& orderedRingSlices, u32& curSliceIdx) const
     718              : {
     719            0 :     curSliceIdx = previousRankList[curSliceIdx];
     720            0 :     Slice* nextSlice = &orderedRingSlices[curSliceIdx];
     721            0 :     return nextSlice;
     722              : }
     723              : 
     724            0 : HcclResult AicpuAllreduce::RunAllReduceRingAlg(
     725              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, std::vector<Slice>& orderedRingSlices,
     726              :     std::vector<u32>& ringOrder, HcclDataType dataType) const
     727              : {
     728            0 :     size_t ringSize = ringOrder.size();
     729            0 :     std::vector<u32> previousRankList(ringSize);
     730              :     // 计算每一个rank在ring环上的前一个rank
     731            0 :     CHK_RET(GetPrevRankList(ringOrder, previousRankList));
     732            0 :     const u32 prevRank = previousRankList[ctx_->rankId];
     733            0 :     const u32 nextRank = FindNextRank(previousRankList, ctx_->rankId);
     734            0 :     const u32 subStream = prevRank;
     735            0 :     u32 curSliceIdx = ctx_->rankId;
     736            0 :     Slice* localSlice = &orderedRingSlices[curSliceIdx];
     737            0 :     Slice* remoteSlice = nullptr;
     738              :     // 第一轮:主流准备前2片数据
     739            0 :     CHK_RET(TaskOrchestrator::SelfCpySnd2Win(
     740              :         sendBuffer, localSlice->size, localSlice->offset, localSlice->offset, HCCL_REDUCE_RESERVED, dataType));
     741            0 :     localSlice = GetNextRingSlice(previousRankList, orderedRingSlices, curSliceIdx);
     742            0 :     CHK_RET(TaskOrchestrator::SelfCpySnd2Win(
     743              :         sendBuffer, localSlice->size, localSlice->offset, localSlice->offset, HCCL_REDUCE_RESERVED, dataType));
     744            0 :     for (size_t rankOffset = 0; rankOffset < ringSize - 1; rankOffset++) {
     745            0 :         remoteSlice = localSlice;
     746              :         // 从流开始 reduce 操作
     747            0 :         CHK_RET(TaskOrchestrator::MainSubPreSync(subStream)); // 主流启动从流
     748              :         // 从流跨片 reduce
     749            0 :         CHK_RET(RingIPCPreSync(subStream, prevRank, nextRank));
     750            0 :         CHK_RET(TaskOrchestrator::IpcCpyWin2WinP2P(
     751              :             prevRank, remoteSlice->size, remoteSlice->offset, remoteSlice->offset, opType, dataType));
     752            0 :         CHK_RET(RingIPCPostSync(subStream, prevRank, nextRank));
     753            0 :         if (rankOffset < ringSize - 2) { // ringSize - 2: 最后一轮不进行本地搬运
     754              :             // 主流继续准备数据
     755            0 :             localSlice = GetNextRingSlice(previousRankList, orderedRingSlices, curSliceIdx);
     756            0 :             CHK_RET(TaskOrchestrator::SelfCpySnd2Win(
     757              :                 sendBuffer, localSlice->size, localSlice->offset, localSlice->offset, HCCL_REDUCE_RESERVED, dataType));
     758              :         }
     759            0 :         CHK_RET(TaskOrchestrator::MainSubPostSync(subStream)); // 从流通知主流,回收 notify,主流继续执行
     760            0 :         TaskOrchestrator::LaunchTasksEx(0, ctx_->rankNum - 1, ctx_->rankNum); // 下发sqe
     761              :     }
     762            0 :     for (size_t rankOffset = 0; rankOffset < ringSize - 1; rankOffset++) {
     763            0 :         localSlice = remoteSlice; // 主流搬运上一轮从流准备的数据
     764            0 :         remoteSlice = GetNextRingSlice(
     765              :             previousRankList, orderedRingSlices, curSliceIdx); // 从流继续往下循环,搬运reduce好的数据
     766            0 :         CHK_RET(TaskOrchestrator::MainSubPreSync(subStream));  // 主流通知从流回收notify资源
     767            0 :         TaskOrchestrator::SelfCpyWin2Rcv(
     768              :             recvBuffer, localSlice->size, localSlice->offset, localSlice->offset, HCCL_REDUCE_RESERVED, dataType);
     769            0 :         CHK_RET(RingIPCPreSync(subStream, prevRank, nextRank));
     770            0 :         if (rankOffset < ringSize - 2) { // < RingSize - 2: 不是最后一轮,搬到 window 上,让下游读
     771            0 :             CHK_RET(TaskOrchestrator::IpcCpyWin2WinP2P(
     772              :                 prevRank, remoteSlice->size, remoteSlice->offset, remoteSlice->offset, HCCL_REDUCE_RESERVED, dataType));
     773              :         } else { // 最后一轮,dma 消减,直接从对端读入 rcv
     774            0 :             CHK_RET(TaskOrchestrator::IpcCpyWin2RcvP2P(
     775              :                 recvBuffer, prevRank, remoteSlice->size, remoteSlice->offset, remoteSlice->offset, HCCL_REDUCE_RESERVED,
     776              :                 dataType));
     777              :         }
     778            0 :         CHK_RET(RingIPCPostSync(subStream, prevRank, nextRank));
     779            0 :         CHK_RET(TaskOrchestrator::MainSubPostSync(subStream));                // 通知主流继续
     780            0 :         TaskOrchestrator::LaunchTasksEx(0, ctx_->rankNum - 1, ctx_->rankNum); // 下发sqe
     781              :     }
     782            0 :     return HCCL_SUCCESS;
     783            0 : }
     784              : 
     785            0 : HcclResult AicpuAllreduce::RunAllReduceRingSingleBurst(
     786              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, u64 dataCount, HcclDataType dataType,
     787              :     std::vector<std::vector<u32>>& ringOrders) const
     788              : {
     789            0 :     std::vector<std::vector<Slice>> orderedRingSlices;
     790            0 :     PrepareRingSlice(ringOrders, dataCount, dataType, orderedRingSlices);
     791            0 :     for (size_t i = 0; i < ringOrders.size(); i++) {
     792            0 :         CHK_RET(RunAllReduceRingAlg(opType, sendBuffer, recvBuffer, orderedRingSlices[i], ringOrders[i], dataType));
     793              :     }
     794            0 :     return HCCL_SUCCESS;
     795            0 : }
     796              : 
     797              : // 单 ring 算法, 使用了DMA消减
     798            2 : HcclResult AicpuAllreduce::RunAllReduceRing(
     799              :     HcclReduceOp opType, void* sendBuffer, void* recvBuffer, u64 dataCount, HcclDataType dataType) const
     800              : {
     801            2 :     HCCL_INFO("run RunAllReduceRing start");
     802              :     // 数据准备阶段
     803            2 :     std::vector<u64> burstDataCounts;
     804            2 :     u64 windowSize = ctx_->windowSize;
     805            2 :     u32 unitSize = ctx_->unitSize;
     806            2 :     CHK_RET(GetBurstDataCounts(windowSize, dataCount, burstDataCounts));
     807            0 :     u8* currSendBuffer = static_cast<u8*>(sendBuffer);
     808            0 :     u8* currRecvBuffer = static_cast<u8*>(recvBuffer);
     809              :     u64 burstSize;
     810            0 :     std::vector<std::vector<u32>> ringOrders = GetRingOrders();
     811            0 :     for (u64 burstDataCount : burstDataCounts) {
     812            0 :         burstSize = burstDataCount * unitSize;
     813            0 :         CHK_RET(
     814              :             RunAllReduceRingSingleBurst(opType, currSendBuffer, currRecvBuffer, burstDataCount, dataType, ringOrders));
     815            0 :         currSendBuffer += burstSize;
     816            0 :         currRecvBuffer += burstSize;
     817              :     }
     818            0 :     HCCL_INFO("run RunAllReduceRing end");
     819            0 :     return HCCL_SUCCESS;
     820            2 : }
        

Generated by: LCOV version 2.0-1