LCOV - code coverage report
Current view: top level - legacy/ascend910/algorithm/impl/coll_executor/coll_reduce_scatter - coll_reduce_scatter_ring_zerocopy_exchange_executor.cc (source / functions) Coverage Total Hit
Test: coverage.info Lines: 0.0 % 105 0
Test Date: 2026-08-18 17:47:01 Functions: 0.0 % 6 0

            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 "coll_reduce_scatter_ring_zerocopy_exchange_executor.h"
      12              : 
      13              : namespace hccl {
      14              : 
      15            0 : CollReduceScatterRingZerocopyExchangeExecutor::CollReduceScatterRingZerocopyExchangeExecutor(
      16            0 :     const HcclDispatcher dispatcher, std::unique_ptr<TopoMatcher>& topoMatcher)
      17            0 :     : CollReduceScatterRingZerocopyExecutor(dispatcher, topoMatcher)
      18            0 : {}
      19              : 
      20            0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::CalcCommInfo(std::vector<LevelNSubCommTransport>& opTransport)
      21              : {
      22              :     // 调用父类编排函数建链关系计算函数
      23            0 :     CHK_RET(CollReduceScatterRingZerocopyExecutor::CalcCommInfo(opTransport));
      24              :     // 额外增加数据交换的建链
      25            0 :     CHK_RET(CalcExchangeCommInfo(opTransport));
      26            0 :     return HCCL_SUCCESS;
      27              : }
      28              : 
      29              : HcclResult
      30            0 : CollReduceScatterRingZerocopyExchangeExecutor::CalcExchangeCommInfo(std::vector<LevelNSubCommTransport>& opTransport)
      31              : {
      32            0 :     std::set<u32> commTargetUserRankSet;
      33            0 :     u32 remoteRankSend = 0;
      34            0 :     u32 remoteRankRecv = 0;
      35              : 
      36            0 :     CHK_RET(CalExchangeRemoteRankForReduceScatter(remoteRankSend, remoteRankRecv));
      37            0 :     commTargetUserRankSet.insert(remoteRankSend);
      38            0 :     commTargetUserRankSet.insert(remoteRankRecv);
      39              :     CommParaInfo commParaInfo(
      40              :         COMM_COMBINE_ORDER, CommType::COMM_TAG_PARTIAL_MESH_COMBINED, INVALID_VALUE_RANKID, INVALID_VALUE_RANKID, false,
      41            0 :         false, commTargetUserRankSet);
      42              : 
      43            0 :     TransportMemType inputType = TransportMemType::CCL_INPUT;
      44            0 :     TransportMemType outputType = TransportMemType::CCL_OUTPUT;
      45              : 
      46            0 :     CHK_RET(CalcCommPlaneInfo(tag_, commParaInfo, opTransport[COMM_COMBINE_ORDER], inputType, outputType));
      47            0 :     LevelNSubCommTransport& commTransport = opTransport[COMM_COMBINE_ORDER];
      48            0 :     for (u32 subCommIndex = 0; subCommIndex < commTransport.size(); subCommIndex++) {
      49            0 :         for (auto& transportRequest : commTransport[subCommIndex].transportRequests) {
      50              :             transportRequest.isUsedRdma
      51            0 :                 = (topoAttr_.superPodNum > 1
      52            0 :                    || (static_cast<bool>(topoMatcher_->GetExternalInputInterHccsDisable()) && topoAttr_.serverNum > 1));
      53              :         }
      54              :     }
      55            0 :     return HCCL_SUCCESS;
      56            0 : }
      57              : 
      58            0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::KernelRunInterServerPostProcess(
      59              :     const OpParam& param, const ExecMem& execMem)
      60              : {
      61              :     // 计算需要交换数据的通信对端
      62            0 :     u32 remoteRankSend = 0;
      63            0 :     u32 remoteRankRecv = 0;
      64            0 :     CHK_RET(CalExchangeRemoteRankForReduceScatter(remoteRankSend, remoteRankRecv));
      65              : 
      66            0 :     Stream stream = param.stream;
      67            0 :     u64 outputMemSize = execMem.outputMem.size();
      68            0 :     if (remoteRankSend != topoAttr_.userRank && remoteRankRecv != topoAttr_.userRank) { // 需要交换数据
      69              :         // 获取通信对端的link
      70            0 :         LINK sendLink;
      71            0 :         LINK recvLink;
      72            0 :         CHK_RET(GetTransportForExchange(remoteRankSend, sendLink));
      73            0 :         CHK_RET(GetTransportForExchange(remoteRankRecv, recvLink));
      74            0 :         CHK_PTR_NULL(sendLink);
      75            0 :         CHK_PTR_NULL(recvLink);
      76              :         // 当通信对端恰好是同server的邻居时,复用Level0的建链,其注册的内存是UserMem,需要特殊处理
      77              :         // 否则,在CommCombineOrder上建链,其注册内存是CCL Buffer
      78            0 :         if (IsLevel0Neighbor(remoteRankSend, level0RankSize_)) {
      79              :             // ccl in -> user in
      80            0 :             u64 memOffset = (level1Rank_ * level2RankSize_ + level2Rank_) * outputMemSize;
      81            0 :             DeviceMem dstMem = DeviceMem::create(static_cast<u8*>(param.inputPtr), execMem.outputMem.size());
      82            0 :             DeviceMem srcMem = execMem.inputMem.range(memOffset, outputMemSize);
      83            0 :             CHK_SMART_PTR_NULL(dstMem);
      84            0 :             HcclResult ret = HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream);
      85            0 :             CHK_PRT_RET(
      86              :                 ret != HCCL_SUCCESS,
      87              :                 HCCL_ERROR(
      88              :                     "[CollReduceScatterRingZerocopyExchangeExecutor][ExchangeData]ReduceScatter double "
      89              :                     "ring memcpy Failed, Offset[%llu], Size[%llu]",
      90              :                     memOffset, outputMemSize),
      91              :                 ret);
      92              :             // user in send to remote user out
      93            0 :             recvLink->TxAck(stream);
      94            0 :             sendLink->RxAck(stream);
      95            0 :             sendLink->TxAsync(UserMemType::OUTPUT_MEM, 0, param.inputPtr, outputMemSize, stream);
      96            0 :             if (IsLevel0Neighbor(remoteRankRecv, level0RankSize_)) {
      97            0 :                 recvLink->RxAsync(UserMemType::INPUT_MEM, 0, execMem.outputPtr, outputMemSize, stream);
      98              :             } else {
      99            0 :                 u32 remoteLevel1Index = remoteRankRecv % (level0RankSize_ * level1RankSize_) / level0RankSize_;
     100            0 :                 u32 remoteLevel2Index = remoteRankRecv / level0RankSize_ / level1RankSize_;
     101            0 :                 u64 rxSrcOffset = (remoteLevel1Index * level2RankSize_ + remoteLevel2Index) * outputMemSize;
     102            0 :                 recvLink->RxAsync(UserMemType::INPUT_MEM, rxSrcOffset, execMem.outputMem.ptr(), outputMemSize, stream);
     103              :             }
     104            0 :         } else {
     105            0 :             recvLink->TxAck(stream);
     106            0 :             sendLink->RxAck(stream);
     107            0 :             u64 txDstOffset = (level1Rank_ * level2RankSize_ + level2Rank_) * outputMemSize;
     108            0 :             sendLink->TxAsync(
     109            0 :                 UserMemType::OUTPUT_MEM, 0, static_cast<u8*>(execMem.inputMem.ptr()) + txDstOffset, outputMemSize,
     110              :                 stream);
     111            0 :             if (IsLevel0Neighbor(remoteRankRecv, level0RankSize_)) {
     112            0 :                 recvLink->RxAsync(UserMemType::INPUT_MEM, 0, execMem.outputPtr, outputMemSize, stream);
     113              :             } else {
     114            0 :                 u32 remoteLevel1Index = remoteRankRecv % (level0RankSize_ * level1RankSize_) / level0RankSize_;
     115            0 :                 u32 remoteLevel2Index = remoteRankRecv / level0RankSize_ / level1RankSize_;
     116            0 :                 u64 rxSrcOffset = (remoteLevel1Index * level2RankSize_ + remoteLevel2Index) * outputMemSize;
     117            0 :                 recvLink->RxAsync(UserMemType::INPUT_MEM, rxSrcOffset, execMem.outputMem.ptr(), outputMemSize, stream);
     118              :             }
     119              :         }
     120              :         // 交换数据的两端之间Barrier,确认收发完成
     121            0 :         CHK_RET(recvLink->TxAck(stream));
     122            0 :         CHK_RET(sendLink->RxAck(stream));
     123            0 :         CHK_RET(sendLink->TxDataSignal(stream));
     124            0 :         CHK_RET(recvLink->RxDataSignal(stream));
     125            0 :     } else { // 不需要交换数据,将数据从ccl in拷到ccl out
     126            0 :         u64 memOffset = (level1Rank_ * level2RankSize_ + level2Rank_) * outputMemSize;
     127            0 :         DeviceMem dstMem = execMem.outputMem;
     128            0 :         DeviceMem srcMem = execMem.inputMem.range(memOffset, outputMemSize);
     129            0 :         CHK_RET(HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream));
     130            0 :     }
     131              : 
     132              :     // 如果不需要交换数据,或者收端不是邻居,那么结果数据还在CCL Out上,需要搬到User Out上去
     133            0 :     if ((remoteRankRecv == topoAttr_.userRank) || !IsLevel0Neighbor(remoteRankRecv, level0RankSize_)) {
     134            0 :         DeviceMem srcMem = execMem.outputMem;
     135            0 :         DeviceMem dstMem = DeviceMem::create(static_cast<u8*>(execMem.outputPtr), execMem.outputMem.size());
     136            0 :         CHK_RET(HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream));
     137            0 :     }
     138              : 
     139            0 :     return HCCL_SUCCESS;
     140            0 : }
     141              : 
     142            0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::CalcLevel0DataSlices(
     143              :     const OpParam& param, const ExecMem& execMem, std::vector<Slice>& dataSegsSlice)
     144              : {
     145            0 :     return CalcIntraServerDataSlicesContinuous(
     146            0 :         param, execMem, level0RankSize_, level1RankSize_, level2RankSize_, dataSegsSlice);
     147              : }
     148              : 
     149            0 : HcclResult CollReduceScatterRingZerocopyExchangeExecutor::KernelRunInterServerPreProcess(
     150              :     const OpParam& param, const ExecMem& execMem)
     151              : {
     152            0 :     HCCL_CONFIG_INFO(
     153              :         HCCL_ALG, "[CollReduceScatterRingZerocopyExchangeExecutor][KernelRun] userRank[%u] starts.",
     154              :         topoAttr_.userRank);
     155            0 :     u32 unitSize = 0;
     156            0 :     CHK_RET(SalGetDataTypeSize(param.DataDes.dataType, unitSize));
     157              : 
     158            0 :     DeviceMem dstMem;
     159            0 :     DeviceMem srcMem;
     160            0 :     u64 curSize = execMem.outputMem.size();
     161            0 :     Stream stream = param.stream;
     162            0 :     for (u32 i = 0; i < level1RankSize_ * level2RankSize_; i++) {
     163              :         // 拷贝input上每个slice的数据到中转内存,源端每个slice的size固定为output的size
     164            0 :         dstMem = execMem.inputMem.range(i * curSize, curSize);
     165            0 :         srcMem = DeviceMem::create(
     166            0 :             static_cast<u8*>(execMem.inputPtr)
     167            0 :                 + param.DataDes.count * unitSize * level1RankSize_ * level2RankSize_ * level0Rank_
     168            0 :                 + param.DataDes.count * unitSize * i,
     169            0 :             curSize);
     170            0 :         CHK_RET(HcclD2DMemcpyAsync(dispatcher_, dstMem, srcMem, stream));
     171              :     }
     172            0 :     return HCCL_SUCCESS;
     173            0 : }
     174              : 
     175              : REGISTER_EXEC(
     176              :     "ReduceScatterRingZerocopyExchangeExecutor", ReduceScatterRingZerocopy,
     177              :     CollReduceScatterRingZerocopyExchangeExecutor);
     178              : } // namespace hccl
        

Generated by: LCOV version 2.0-1