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

Generated by: LCOV version 2.0-1