LCOV - code coverage report
Current view: top level - legacy/ascend910/algorithm/impl/coll_executor/coll_send_receive - coll_batch_send_recv_retry_executor.cc (source / functions) Coverage Total Hit
Test: coverage.info Lines: 0.0 % 162 0
Test Date: 2026-08-04 10:52:23 Functions: 0.0 % 11 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_batch_send_recv_retry_executor.h"
      12              : namespace hccl {
      13              : constexpr u32 PAIRSIZE_TWO = 2;
      14              : 
      15            0 : CollBatchSendRecvRetryExecutor::CollBatchSendRecvRetryExecutor(const HcclDispatcher dispatcher,
      16            0 :     std::unique_ptr<TopoMatcher> &topoMatcher)
      17            0 :     : CollBatchSendRecvExecutor(dispatcher, topoMatcher)
      18              : {
      19            0 : }
      20              : 
      21            0 : HcclResult CollBatchSendRecvRetryExecutor::CreatePairWiseList(HcclSendRecvItem *sendRecvInfo, u32 itemNum)
      22              : {
      23            0 :     HCCL_INFO("[CollBatchSendRecvRetryExecutor][GetPairWiseList] Start sort the batchSendRecv tasklist.");
      24            0 :     CHK_PTR_NULL(sendRecvInfo);
      25              : 
      26            0 :     for (u32 i = 0; i < itemNum; i++) {
      27            0 :         HCCL_INFO("[CollBatchSendRecvRetryExecutor][GetPairWiseList] index is %u, itemNum is %u, localRankID is %u, "\
      28              :             "remoteRank is %u, sendRecvType is %u, rankSize is %u.", i, itemNum, topoAttr_.userRank,
      29              :             sendRecvInfo->remoteRank, static_cast<u32>(sendRecvInfo->sendRecvType), topoAttr_.userRankSize);
      30            0 :         CHK_PTR_NULL(sendRecvInfo->buf);
      31              : 
      32            0 :         if (sendRecvInfo->sendRecvType == HcclSendRecvType::HCCL_SEND) {
      33            0 :             sendDeque_.push_back(sendRecvInfo);
      34            0 :         } else if (sendRecvInfo->sendRecvType == HcclSendRecvType::HCCL_RECV) {
      35            0 :             recvDeque_.push_back(sendRecvInfo);
      36              :         } else {
      37            0 :             HCCL_ERROR("[CollBatchSendRecvRetryExecutor][GetPairWiseList] sendRecvType wrong sendrecvType is %d, "\
      38              :                 "rankID is %u, remoteRank is %u.", sendRecvInfo->sendRecvType, topoAttr_.userRank,
      39              :                 sendRecvInfo->remoteRank);
      40            0 :             return HCCL_E_PARA;
      41              :         }
      42            0 :         sendRecvInfo++;
      43              :     }
      44              :     /* 此处的排序逻辑(pair-wise算法):
      45              :         1.sendDeque元素顺序是:先放remoteRank号小于等于root rank的第一个任务,依次减小(循环索引)直至放完
      46              :         2.recvDeque元素顺序是:先放remoteRank号大于等于root rank的第一个任务,依次增大(循环索引)直至放完
      47              :     */
      48            0 :     auto sendCompare = [this](HcclSendRecvItem* a, HcclSendRecvItem* b) {
      49            0 :         u32 aFlag = (a->remoteRank <= topoAttr_.userRank) ? (a->remoteRank + topoAttr_.userRankSize) : a->remoteRank;
      50            0 :         u32 bFlag = (b->remoteRank <= topoAttr_.userRank) ? (b->remoteRank + topoAttr_.userRankSize) : b->remoteRank;
      51            0 :         return aFlag > bFlag;
      52            0 :     };
      53              : 
      54            0 :     auto recvCompare = [this](HcclSendRecvItem* a, HcclSendRecvItem* b) {
      55            0 :         u32 aFlag = (a->remoteRank < topoAttr_.userRank) ? (a->remoteRank + topoAttr_.userRankSize) : a->remoteRank;
      56            0 :         u32 bFlag = (b->remoteRank < topoAttr_.userRank) ? (b->remoteRank + topoAttr_.userRankSize) : b->remoteRank;
      57            0 :         return aFlag < bFlag;
      58            0 :     };
      59              : 
      60            0 :     std::sort(sendDeque_.begin(), sendDeque_.end(), sendCompare);
      61            0 :     std::sort(recvDeque_.begin(), recvDeque_.end(), recvCompare);
      62              : 
      63              :     // 生成SendRecvPair
      64            0 :     u32 pairNum = std::max(sendDeque_.size(), recvDeque_.size());
      65            0 :     for (u32 pairIndex = 0; pairIndex < pairNum; pairIndex++) {
      66            0 :         std::vector<HcclSendRecvItem*> sendRecvPair;
      67            0 :         if (sendDeque_.size() > pairIndex) {
      68            0 :            sendRecvPair.push_back(sendDeque_[pairIndex]);
      69              :         }
      70            0 :         if (recvDeque_.size() > pairIndex) {
      71            0 :             sendRecvPair.push_back(recvDeque_[pairIndex]);
      72              :         }
      73            0 :         sendRecvPairList_.push_back(sendRecvPair);
      74            0 :     }
      75            0 :     HCCL_INFO("[CollBatchSendRecvRetryExecutor][GetPairWiseList] End sort the batchSendRecv tasklist.");
      76            0 :     return HCCL_SUCCESS;
      77              : }
      78              : 
      79            0 : HcclResult CollBatchSendRecvRetryExecutor::GetPairWiseList(std::vector<std::vector<HcclSendRecvItem*>> &sendRecvPairList)
      80              : {
      81            0 :     sendRecvPairList = sendRecvPairList_;
      82            0 :     return HCCL_SUCCESS;
      83              : }
      84              : 
      85            0 : HcclResult CollBatchSendRecvRetryExecutor::CheckSendRecvPair(const std::vector<HcclSendRecvItem*> &sendRecvPair)
      86              : {
      87            0 :     if (sendRecvPair.empty()) {
      88            0 :         HCCL_ERROR("[CollBatchSendRecvRetryExecutor] please check the pair list.");
      89            0 :         return HCCL_E_PARA;
      90              :     }
      91            0 :     if (sendRecvPair.size() == 1 && sendRecvPair[0]->remoteRank == topoAttr_.userRank) {
      92            0 :         HCCL_ERROR("[CollBatchSendRecvRetryExecutor] SendTask and Recv Task to rank itself do not match,"\
      93              :             "please check the task list.");
      94            0 :         return HCCL_E_PARA;
      95              :     }
      96            0 :     return HCCL_SUCCESS;
      97              : }
      98              : 
      99            0 : HcclResult CollBatchSendRecvRetryExecutor::Orchestrate(OpParam& param, AlgResourceResponse& algResource)
     100              : {
     101            0 :     HcclUs startut = TIME_NOW();
     102            0 :     HCCL_CONFIG_INFO(HCCL_ALG, "[CollBatchSendRecvRetryExecutor] batchsendrecv retry starts.");
     103            0 :     algResResp_ = &algResource;
     104            0 :     CHK_RET(CheckCommSize(COMM_COMBINE_ORDER, COMM_SIZE_TWO));
     105              : 
     106              :     // 校验当前sendRecvPair
     107            0 :     std::vector<HcclSendRecvItem*> sendRecvPair;
     108            0 :     if (param.BatchSendRecvDataDes.curIterNum < sendRecvPairList_.size()) {
     109            0 :         sendRecvPair = sendRecvPairList_[param.BatchSendRecvDataDes.curIterNum];
     110              :     } else {
     111            0 :         HCCL_ERROR("[CollBatchSendRecvRetryExecutor] the curIterNum[%u] is out of range[0, %zu].",
     112              :             param.BatchSendRecvDataDes.curIterNum, sendRecvPairList_.size());
     113            0 :         return HCCL_E_PARA;
     114              :     }
     115            0 :     CHK_RET(CheckSendRecvPair(sendRecvPair));
     116              : 
     117              :     // 自发自收场景
     118            0 :     if (sendRecvPair.size() == PAIRSIZE_TWO && sendRecvPair[0]->remoteRank == topoAttr_.userRank &&
     119            0 :         sendRecvPair[1]->remoteRank == topoAttr_.userRank) {
     120            0 :         if (sendRecvPair[0]->count == sendRecvPair[1]->count && sendRecvPair[0]->dataType == sendRecvPair[1]->dataType) {
     121            0 :             u64 dataSize = sendRecvPair[0]->count * SIZE_TABLE[sendRecvPair[0]->dataType];
     122            0 :             DeviceMem inUserMem = DeviceMem::create(static_cast<u8*>(sendRecvPair[0]->buf), dataSize);
     123            0 :             DeviceMem outUserMem = DeviceMem::create(static_cast<u8*>(sendRecvPair[1]->buf), dataSize);
     124            0 :             CHK_RET(HcclD2DMemcpyAsync(dispatcher_, outUserMem, inUserMem, param.stream));
     125            0 :             return HCCL_SUCCESS;
     126            0 :         } else {
     127            0 :              HCCL_ERROR("[HcclBatchSendRecvRetry] Send task and recv task to self : data size do not equal, please"\
     128              :                 "check the task list.");
     129            0 :             return HCCL_E_PARA;
     130              :         }
     131              :     }
     132              : 
     133              :     // 重执行正常执行场景,前后需和控制流做同步
     134            0 :     if (param.BatchSendRecvDataDes.curMode == BatchSendRecvCurMode::SEND_RECV) {
     135            0 :         HCCL_INFO("[BatchSendRecv] Stream sync: main stream record, subStream wait.");
     136            0 :         CHK_RET(LocalNotify::Post(param.stream, dispatcher_, algResResp_->notifiesAux[STREAM_INDEX_0], PROF_STAGE_0));
     137            0 :         CHK_RET(LocalNotify::Wait(algResResp_->slaveStreams[STREAM_INDEX_0], dispatcher_,
     138              :             algResResp_->notifiesAux[STREAM_INDEX_0], PROF_STAGE_0));
     139            0 :         CHK_RET(LocalNotify::Post(param.stream, dispatcher_, algResResp_->notifiesAux[STREAM_INDEX_1], PROF_STAGE_1));
     140            0 :         CHK_RET(LocalNotify::Wait(algResResp_->slaveStreams[STREAM_INDEX_1], dispatcher_,
     141              :             algResResp_->notifiesAux[STREAM_INDEX_1], PROF_STAGE_1));
     142              :     }
     143              :     // run sendrecv
     144            0 :     CHK_RET(RunLoop(param, algResource, sendRecvPair));
     145              : 
     146            0 :     if (param.BatchSendRecvDataDes.curMode == BatchSendRecvCurMode::SEND_RECV) {
     147            0 :         HCCL_INFO("[BatchSendRecv] Stream sync: subStream record, main stream wait.");
     148            0 :         CHK_RET(LocalNotify::Post(algResResp_->slaveStreams[STREAM_INDEX_0], dispatcher_,
     149              :             algResResp_->notifiesMain[STREAM_INDEX_0], PROF_STAGE_0));
     150            0 :         CHK_RET(LocalNotify::Wait(param.stream, dispatcher_, algResResp_->notifiesMain[STREAM_INDEX_0],
     151              :             PROF_STAGE_0));
     152            0 :         CHK_RET(LocalNotify::Post(algResResp_->slaveStreams[STREAM_INDEX_1], dispatcher_,
     153              :             algResResp_->notifiesMain[STREAM_INDEX_1], PROF_STAGE_1));
     154            0 :         CHK_RET(LocalNotify::Wait(param.stream, dispatcher_, algResResp_->notifiesMain[STREAM_INDEX_1],
     155              :             PROF_STAGE_1));
     156            0 :     } else if (param.BatchSendRecvDataDes.curMode == BatchSendRecvCurMode::SEND) {
     157            0 :         CHK_RET(LocalNotify::Post(algResResp_->slaveStreams[STREAM_INDEX_0], dispatcher_,
     158              :             algResResp_->notifiesMain[STREAM_INDEX_0], PROF_STAGE_0));
     159            0 :     } else if (param.BatchSendRecvDataDes.curMode == BatchSendRecvCurMode::RECV) {
     160            0 :         CHK_RET(LocalNotify::Post(algResResp_->slaveStreams[STREAM_INDEX_1], dispatcher_,
     161              :             algResResp_->notifiesMain[STREAM_INDEX_1], PROF_STAGE_1));
     162              :     }
     163              : 
     164            0 :     CHK_RET(LaunchTaskExtend(dispatcher_, param.stream, algResResp_->slaveStreams));
     165            0 :     HCCL_INFO("[info][print] LaunchTaskExtend success.");
     166            0 :     HCCL_INFO("tag[%s] BatchSendRecv Executor orchestrate success, take time [%lld]us.",
     167              :         param.tag.c_str(), DURATION_US(TIME_NOW() - startut));
     168            0 :     return HCCL_SUCCESS;
     169            0 : }
     170              : 
     171            0 : HcclResult CollBatchSendRecvRetryExecutor::RunLoop(OpParam &param, AlgResourceResponse &algRes,
     172              :     const std::vector<HcclSendRecvItem*> &sendRecvPair)
     173              : {
     174              :     // 判断当前需执行的算子
     175            0 :     std::vector<HcclSendRecvItem*> curSendRecvPair;
     176            0 :     if (param.BatchSendRecvDataDes.curMode == BatchSendRecvCurMode::SEND) {
     177            0 :         curSendRecvPair.push_back(sendRecvPair[0]);
     178            0 :     } else if (param.BatchSendRecvDataDes.curMode == BatchSendRecvCurMode::RECV) {
     179            0 :         curSendRecvPair.push_back(sendRecvPair[sendRecvPair.size() - 1]);
     180              :     } else {
     181            0 :         curSendRecvPair = sendRecvPair;
     182              :     }
     183              : 
     184              :     // 执行当前需执行的算子
     185            0 :     for (const auto& itemPtr : curSendRecvPair) {
     186            0 :         if (static_cast<bool>(param.BatchSendRecvDataDes.isDirectRemoteRank[itemPtr->remoteRank])) {
     187              :             // device direct链路的任务会在host侧下发,此处需要跳过
     188            0 :             continue;
     189              :         }
     190            0 :         HCCL_INFO("[CollBatchSendRecvRetryExecutor][RunLoop] remoteRank %u", itemPtr->remoteRank);
     191            0 :         if (itemPtr->sendRecvType == HcclSendRecvType::HCCL_SEND) {
     192            0 :             CHK_RET(CalcSendSlices(algRes, itemPtr));
     193            0 :         } else if (itemPtr->sendRecvType == HcclSendRecvType::HCCL_RECV) {
     194            0 :             CHK_RET(CalcRecvSlices(algRes, itemPtr));
     195              :         } else {
     196            0 :             HCCL_ERROR("[CollBatchSendRecvRetryExecutor][RunLoop] sendRecvType is Wrong.");
     197            0 :             return HCCL_E_PARA;
     198              :         }
     199              :     }
     200              : 
     201            0 :     u32 loopInOnceLaunch = 0;
     202              :     // 每隔200个loop launch一次
     203            0 :     while (!sendDataSilces_.empty() || !recvDataSilces_.empty()) {
     204            0 :         if(!sendDataSilces_.empty()) {
     205            0 :             CHK_RET(ProcessSendDataSlice(algResResp_->slaveStreams[STREAM_INDEX_0], false, true)); 
     206            0 :             sendDataSilces_.pop_front();
     207              :         }
     208            0 :         if(!recvDataSilces_.empty()) {
     209            0 :             CHK_RET(ProcessRecvDataSlice(algResResp_->slaveStreams[STREAM_INDEX_1], true));
     210            0 :             recvDataSilces_.pop_front();
     211              :         }
     212            0 :         loopInOnceLaunch++;
     213            0 :         if (loopInOnceLaunch == MAX_LOOP_IN_ONCE_LAUNCH || (sendDataSilces_.empty() && recvDataSilces_.empty())) {
     214            0 :             CHK_RET(LaunchTaskExtend(dispatcher_, param.stream, algResResp_->slaveStreams));
     215            0 :             HCCL_INFO("[BatchSendRecv] LaunchTaskExtend, unprocessed send slices[%u], recv slices[%u].",
     216              :                 sendDataSilces_.size(), recvDataSilces_.size());
     217            0 :             loopInOnceLaunch = 0;
     218              :         }
     219              :     }
     220            0 :     return HCCL_SUCCESS;
     221            0 : }
     222              : 
     223            0 : HcclResult CollBatchSendRecvRetryExecutor::CalcStreamNum(u32& streamNum)
     224              : {
     225            0 :     streamNum = LEVEL0_PLANE_NUM_IN_NPRING_DOUBLE;
     226            0 :     HCCL_INFO("[CollBatchSendRecvRetryExecutor][CalcScratchMemSize] tag_[%s], streamNum[%u].", tag_.c_str(), streamNum);
     227            0 :     return HCCL_SUCCESS;
     228              : }
     229              : 
     230            0 : HcclResult CollBatchSendRecvRetryExecutor::CalcSendSlices(AlgResourceResponse& algRes, HcclSendRecvItem* sendRecvItem)
     231              : {
     232            0 :     HCCL_INFO("[CollBatchSendRecvExecutor][CalcSendSlices] tag[%s], remoteRank[%u], buf[%p], count[%llu],"\
     233              :         "dataType[%s], sendRecvType[%d].", tag_.c_str(), sendRecvItem->remoteRank, sendRecvItem->buf,
     234              :         sendRecvItem->count, GetDataTypeEnumStr(sendRecvItem->dataType).c_str(), sendRecvItem->sendRecvType);
     235            0 :     u8 *curInputPtr = static_cast<u8 *>(sendRecvItem->buf);
     236            0 :     CHK_PTR_NULL(curInputPtr);
     237            0 :     u32 unitSize = SIZE_TABLE[sendRecvItem->dataType];
     238            0 :     u64 maxCountPerLoop = CalcSendLoopMaxCount(const_cast<DeviceMem&>(algRes.cclInputMem), unitSize);
     239              : 
     240            0 :     for (u64 countLeft = sendRecvItem->count, curCount = 0, curOffset = 0; countLeft > 0;
     241            0 :         countLeft -= curCount) {
     242            0 :         curInputPtr += curOffset;
     243            0 :         curCount = (countLeft > maxCountPerLoop) ? maxCountPerLoop : countLeft;
     244            0 :         u64 curSize = curCount * unitSize; // 单位:字节
     245            0 :         sendDataSilces_.emplace_back(curInputPtr, curSize, sendRecvItem->remoteRank);
     246            0 :         HCCL_DEBUG("[CollBatchSendRecvExecutor][CalcSendSlices] tag[%s], slice userAddr[%p], slice size[%llu].",
     247              :             tag_.c_str(), curInputPtr, curSize);
     248            0 :         curOffset = curSize;
     249              :     }
     250            0 :     return HCCL_SUCCESS;
     251              : }
     252              : 
     253            0 : HcclResult CollBatchSendRecvRetryExecutor::CalcRecvSlices(AlgResourceResponse& algRes, HcclSendRecvItem* sendRecvItem)
     254              : {
     255            0 :     HCCL_INFO("[CollBatchSendRecvRetryExecutor][CalcSendSlices] tag[%s], remoteRank[%u], buf[%p], count[%llu],"\
     256              :         "dataType[%s], sendRecvType[%d].", tag_.c_str(), sendRecvItem ->remoteRank, sendRecvItem ->buf, sendRecvItem->count,
     257              :         GetDataTypeEnumStr(sendRecvItem->dataType).c_str(), sendRecvItem->sendRecvType);
     258            0 :     u8 *curOutputPtr = static_cast<u8*>(sendRecvItem->buf);
     259            0 :     CHK_PTR_NULL(curOutputPtr);
     260            0 :     u32 unitSize = SIZE_TABLE[sendRecvItem->dataType];
     261            0 :     u64 maxCountPerLoop = CalcRecvLoopMaxCount(const_cast<DeviceMem&>(algRes.cclOutputMem), unitSize);
     262            0 :     HCCL_DEBUG("[CollBatchSendRecvRetryExecutor][CalcSendSlices]maxCountPerLoop is %llu", maxCountPerLoop);
     263              : 
     264            0 :     for (u64 countLeft = sendRecvItem->count, curCount = 0, curOffset = 0; countLeft > 0;
     265            0 :         countLeft -= curCount) {
     266            0 :         curOutputPtr += curOffset;
     267            0 :         curCount = (countLeft > maxCountPerLoop) ? maxCountPerLoop : countLeft;
     268            0 :         u64 curSize = curCount * unitSize; // 单位:字节
     269            0 :         recvDataSilces_.emplace_back(curOutputPtr, curSize, sendRecvItem->remoteRank);
     270            0 :         HCCL_DEBUG("[CollBatchSendRecvRetryExecutor][CalcRecvSlices] tag[%s], slice userAddr[%p], slice size[%llu].",
     271              :             tag_.c_str(), curOutputPtr, curSize);
     272            0 :         curOffset = curSize;
     273              :     }
     274            0 :     return HCCL_SUCCESS;
     275              : }
     276              : 
     277              : REGISTER_EXEC("BatchSendRecvRetry", BatchSendRecvRetryExecutor, CollBatchSendRecvRetryExecutor);
     278              : } // namespace hccl
        

Generated by: LCOV version 2.0-1