LCOV - code coverage report
Current view: top level - acl/acl_tdt_queue - queue_process.cpp (source / functions) Coverage Total Hit
Test: coverage.info Lines: 93.7 % 459 430
Test Date: 2026-07-28 10:53:01 Functions: 100.0 % 34 34

            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              : #include "queue_process.h"
      11              : #include "log_inner.h"
      12              : #include "runtime/rt_mem_queue.h"
      13              : #include "runtime/dev.h"
      14              : #include "queue_schedule/qs_client.h"
      15              : 
      16              : namespace acl {
      17              : constexpr uint32_t RT_MQ_DEPTH_DEFAULT = 8U;
      18              : constexpr uint16_t MBUF_ENHANCED_QS = 2U;
      19              : constexpr uint16_t MBUF_ENHANCED_ACL = 1U;
      20              : bool QueueProcessor::isInitQs_ = false;
      21              : bool QueueProcessor::isMbufInit_ = false;
      22              : 
      23            4 : aclError QueueProcessor::acltdtEnqueue(const uint32_t qid, const acltdtBuf buf, const int32_t timeout)
      24              : {
      25            4 :     ACL_REQUIRES_NOT_NULL_WITH_INPUT_REPORT(buf);
      26            3 :     const QueueDataMutexPtr muPtr = GetMutexForData(qid);
      27            3 :     ACL_CHECK_MALLOC_RESULT(muPtr);
      28            3 :     const uint64_t startTime = GetTimestamp();
      29            3 :     uint64_t endTime = 0U;
      30            3 :     bool continueFlag = false;
      31              :     do {
      32          541 :         const std::lock_guard<std::mutex> lk(muPtr->muForDequeue);
      33          541 :         constexpr int32_t deviceId = 0;
      34          541 :         const rtError_t rtRet = rtMemQueueEnQueue(deviceId, qid, buf);
      35          541 :         if (rtRet == RT_ERROR_NONE) {
      36            1 :             return ACL_SUCCESS;
      37              :         }
      38          540 :         if (rtRet != ACL_ERROR_RT_QUEUE_FULL) {
      39            1 :             return rtRet;
      40              :         }
      41          539 :         (void)mmSleep(1U); // sleep 1ms
      42          539 :         endTime = GetTimestamp();
      43          539 :         continueFlag =
      44          539 :             ((endTime - startTime) <= (static_cast<uint64_t>(timeout) * static_cast<uint64_t>(MSEC_TO_USEC)));
      45         1080 :     } while (continueFlag || (timeout < 0));
      46            1 :     return ACL_ERROR_FAILURE;
      47            3 : }
      48              : 
      49            3 : aclError QueueProcessor::acltdtDequeue(const uint32_t qid, acltdtBuf* const buf, const int32_t timeout)
      50              : {
      51            3 :     ACL_REQUIRES_NOT_NULL_WITH_INPUT_REPORT(buf);
      52            3 :     const QueueDataMutexPtr muPtr = GetMutexForData(qid);
      53            3 :     ACL_CHECK_MALLOC_RESULT(muPtr);
      54            3 :     const uint64_t startTime = GetTimestamp();
      55            3 :     uint64_t endTime = 0U;
      56            3 :     bool continueFlag = false;
      57              :     do {
      58          542 :         const std::lock_guard<std::mutex> lk(muPtr->muForEnqueue);
      59          542 :         constexpr int32_t deviceId = 0;
      60          542 :         const rtError_t rtRet = rtMemQueueDeQueue(deviceId, qid, buf);
      61          542 :         if (rtRet == RT_ERROR_NONE) {
      62            1 :             return ACL_SUCCESS;
      63              :         }
      64          541 :         if (rtRet != ACL_ERROR_RT_QUEUE_EMPTY) {
      65            1 :             return rtRet;
      66              :         }
      67          540 :         (void)mmSleep(1U); // sleep 1ms
      68          540 :         endTime = GetTimestamp();
      69          540 :         continueFlag =
      70          540 :             ((endTime - startTime) <= (static_cast<uint64_t>(timeout) * static_cast<uint64_t>(MSEC_TO_USEC)));
      71         1082 :     } while (continueFlag || (timeout < 0));
      72            1 :     return ACL_ERROR_FAILURE;
      73            3 : }
      74              : 
      75            1 : aclError QueueProcessor::acltdtGrantQueue(
      76              :     const uint32_t qid, const int32_t pid, const uint32_t permission, const int32_t timeout)
      77              : {
      78              :     (void)(qid);
      79              :     (void)(pid);
      80              :     (void)(permission);
      81              :     (void)(timeout);
      82            1 :     ACL_LOG_ERROR("[Unsupport][Feature]acltdtGrantQueue is not supported in this version. Please check.");
      83            1 :     const char_t* argList[] = {"func"};
      84            1 :     const char_t* argVal[] = {__func__};
      85            1 :     acl::AclErrorLogManager::ReportInputErrorWithChar(acl::UNSUPPORTED_SYSTEM_MSG, argList, argVal, 1UL);
      86            1 :     return ACL_ERROR_FEATURE_UNSUPPORTED;
      87              : }
      88              : 
      89            1 : aclError QueueProcessor::acltdtAttachQueue(const uint32_t qid, const int32_t timeout, uint32_t* const permission)
      90              : {
      91              :     (void)(qid);
      92              :     (void)(permission);
      93              :     (void)(timeout);
      94            1 :     ACL_LOG_ERROR("[Unsupport][Feature]acltdtAttachQueue is not supported in this version. Please check.");
      95            1 :     const char_t* argList[] = {"func"};
      96            1 :     const char_t* argVal[] = {__func__};
      97            1 :     acl::AclErrorLogManager::ReportInputErrorWithChar(acl::UNSUPPORTED_SYSTEM_MSG, argList, argVal, 1UL);
      98            1 :     return ACL_ERROR_FEATURE_UNSUPPORTED;
      99              : }
     100              : 
     101            2 : aclError QueueProcessor::acltdtDestroyQueueOndevice(const uint32_t qid, const bool isThreadMode)
     102              : {
     103            2 :     ACL_LOG_INFO("Start to destroy queue %u", qid);
     104            2 :     constexpr int32_t deviceId = 0;
     105              :     // get qs id
     106            2 :     int32_t dstPid = 0;
     107            2 :     size_t routeNum = 0UL;
     108            2 :     const std::lock_guard<std::recursive_mutex> lk(muForQueueCtrl_);
     109            2 :     if (GetDstInfo(deviceId, QS_PID, dstPid, isThreadMode) == ACL_SUCCESS) {
     110            2 :         ACL_LOG_INFO("find qs pid %d", dstPid);
     111            2 :         rtEschedEventSummary_t eventSum = {0, 0U, 0, 0U, 0U, nullptr, 0U, 0};
     112            2 :         rtEschedEventReply_t ack = {nullptr, 0U, 0U};
     113            2 :         bqs::QsProcMsgRsp qsRsp = {0UL, 0, 0U, 0U, 0U, {0}};
     114            2 :         eventSum.pid = dstPid;
     115            2 :         eventSum.grpId = bqs::BIND_QUEUE_GROUP_ID;
     116            2 :         eventSum.eventId = RT_MQ_SCHED_EVENT_QS_MSG; // qs EVENT_ID
     117            2 :         eventSum.dstEngine = static_cast<uint32_t>(RT_MQ_DST_ENGINE_CCPU_DEVICE);
     118            2 :         ack.buf = reinterpret_cast<char_t*>(&qsRsp);
     119            2 :         ack.bufLen = sizeof(qsRsp);
     120            2 :         const acltdtQueueRouteQueryInfo queryInfo = {bqs::BQS_QUERY_TYPE_SRC_OR_DST, qid, qid, true, true, true};
     121            2 :         ACL_REQUIRES_OK(GetQueueRouteNum(&queryInfo, deviceId, eventSum, ack, routeNum));
     122              :     }
     123            2 :     if (routeNum > 0U) {
     124            0 :         ACL_LOG_ERROR("qid [%u] can not be destroyed, it need to be unbinded first.", qid);
     125            0 :         return ACL_ERROR_FAILURE;
     126              :     }
     127            2 :     ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(rtMemQueueDestroy(deviceId, qid), rtMemQueueDestroy);
     128            2 :     DeleteMutexForData(qid);
     129            2 :     ACL_LOG_INFO("successfully to execute destroy queue %u", qid);
     130            2 :     return ACL_SUCCESS;
     131            2 : }
     132              : 
     133            2 : aclError QueueProcessor::GetQueuePermission(
     134              :     const int32_t deviceId, uint32_t qid, rtMemQueueShareAttr_t& permission) const
     135              : {
     136            2 :     uint32_t outLen = sizeof(permission);
     137            2 :     if (rtMemQueueQuery(deviceId, RT_MQ_QUERY_QUE_ATTR_OF_CUR_PROC, &qid, sizeof(qid), &permission, &outLen) !=
     138              :         RT_ERROR_NONE) {
     139            0 :         return ACL_ERROR_FAILURE;
     140              :     }
     141            2 :     return ACL_SUCCESS;
     142              : }
     143              : 
     144            5 : aclError QueueProcessor::InitQueueSchedule(const int32_t devId) const
     145              : {
     146            5 :     if (!isInitQs_) {
     147            1 :         ACL_LOG_INFO("need to init queue schedule");
     148            1 :         ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(rtMemQueueInitQS(devId, nullptr), rtMemQueueInitQS);
     149            1 :         isInitQs_ = true;
     150              :     }
     151            5 :     return ACL_SUCCESS;
     152              : }
     153              : 
     154           20 : aclError QueueProcessor::GetDstInfo(
     155              :     const int32_t deviceId, const PID_QUERY_TYPE type, int32_t& dstPid, const bool isThreadMode) const
     156              : {
     157           20 :     if (isThreadMode && isInitQs_) {
     158            1 :         dstPid = mmGetPid();
     159            1 :         return ACL_SUCCESS;
     160              :     }
     161           19 :     rtBindHostpidInfo_t info = {0, 0U, 0U, 0};
     162           19 :     info.hostPid = mmGetPid();
     163           19 :     if (type == CP_PID) {
     164           13 :         info.cpType = RT_DEV_PROCESS_CP1;
     165              :     } else {
     166            6 :         info.cpType = RT_DEV_PROCESS_QS;
     167              :     }
     168           19 :     info.chipId = static_cast<uint32_t>(deviceId);
     169           19 :     ACL_LOG_INFO("start to get dst pid, deviceId is %d, type is %d", deviceId, type);
     170           19 :     const auto ret = rtQueryDevPid(&info, &dstPid);
     171           19 :     if (ret != ACL_RT_SUCCESS) {
     172            1 :         ACL_LOG_INFO("can not query device pid");
     173            1 :         return ret;
     174              :     }
     175           18 :     ACL_LOG_INFO("get dst pid %d success, type is %d", dstPid, type);
     176           18 :     return ACL_SUCCESS;
     177              : }
     178              : 
     179           10 : static aclError AllocMBufOnDevice(void** const devPtr, void** const mBuf, const size_t size)
     180              : {
     181           10 :     ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(rtMbufAlloc(mBuf, size), rtMbufAlloc);
     182           10 :     ACL_CHECK_MALLOC_RESULT(*mBuf);
     183           10 :     if (rtMbufGetBuffAddr(*mBuf, devPtr) != RT_ERROR_NONE) {
     184            0 :         (void)rtMbufFree(*mBuf);
     185            0 :         return ACL_ERROR_BAD_ALLOC;
     186              :     }
     187           10 :     if (*devPtr == nullptr) {
     188            0 :         (void)rtMbufFree(*mBuf);
     189            0 :         ACL_LOG_INNER_ERROR("[Get][mbuf]get dataPtr failed.");
     190            0 :         return ACL_ERROR_BAD_ALLOC;
     191              :     }
     192           10 :     (void)memset_s(*devPtr, size, 0, size);
     193           10 :     return ACL_SUCCESS;
     194              : }
     195              : 
     196            5 : aclError QueueProcessor::SendConnectQsMsg(
     197              :     const int32_t deviceId, rtEschedEventSummary_t& eventSum, rtEschedEventReply_t& ack)
     198              : {
     199              :     // send contact msg
     200            5 :     ACL_LOG_INFO("start to send contact msg");
     201            5 :     bqs::QsBindInit qsInitMsg = {0U, 0, 0U, MBUF_ENHANCED_ACL, {0}};
     202            5 :     qsInitMsg.pid = mmGetPid();
     203            5 :     qsInitMsg.grpId = 0U;
     204            5 :     eventSum.subeventId = bqs::ACL_BIND_QUEUE_INIT;
     205            5 :     eventSum.msgLen = sizeof(qsInitMsg);
     206            5 :     eventSum.msg = reinterpret_cast<char_t*>(&qsInitMsg);
     207            5 :     const rtError_t ret = rtEschedSubmitEventSync(deviceId, &eventSum, &ack);
     208            5 :     eventSum.msgLen = 0U;
     209            5 :     eventSum.msg = nullptr;
     210            5 :     ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(ret, rtEschedSubmitEventSync);
     211            5 :     bqs::QsProcMsgRsp* const rsp = reinterpret_cast<bqs::QsProcMsgRsp*>(ack.buf);
     212            5 :     if (rsp->retCode != 0) {
     213            1 :         ACL_LOG_INNER_ERROR("send connect qs failed, ret code is %d", rsp->retCode);
     214            1 :         return ACL_ERROR_FAILURE;
     215              :     }
     216            4 :     qsContactId_ = rsp->retValue;
     217            4 :     if (rsp->majorVersion >= MBUF_ENHANCED_QS) {
     218            1 :         isMbufEnhanced_ = true;
     219              :     }
     220            4 :     ACL_LOG_INFO("successfully execute to SendConnectQsMsg");
     221            4 :     return ACL_SUCCESS;
     222              : }
     223              : 
     224            6 : aclError QueueProcessor::SendBindUnbindMsgOnDevice(
     225              :     acltdtQueueRouteList* const qRouteList, const bool isBind, rtEschedEventSummary_t& eventSum,
     226              :     rtEschedEventReply_t& ack) const
     227              : {
     228            6 :     ACL_LOG_INFO("start to send bind or unbind msg");
     229              :     // send bind or unbind msg
     230            6 :     const size_t routeSize = sizeof(bqs::QsRouteHead) + (qRouteList->routeList.size() * sizeof(bqs::QueueRoute));
     231            6 :     ACL_LOG_INFO("route size is %zu, queue route num is %zu", routeSize, qRouteList->routeList.size());
     232            6 :     void* devPtr = nullptr;
     233            6 :     void* mBuf = nullptr;
     234            6 :     ACL_REQUIRES_OK(AllocMBufOnDevice(&devPtr, &mBuf, routeSize));
     235            6 :     bqs::QsRouteHead* const head = reinterpret_cast<bqs::QsRouteHead*>(devPtr);
     236            6 :     head->length = routeSize;
     237            6 :     head->routeNum = qRouteList->routeList.size();
     238            6 :     head->subEventId =
     239            6 :         isBind ? static_cast<uint32_t>(bqs::ACL_BIND_QUEUE) : static_cast<uint32_t>(bqs::ACL_UNBIND_QUEUE);
     240            6 :     size_t offset = sizeof(bqs::QsRouteHead);
     241            9 :     for (size_t i = 0UL; i < qRouteList->routeList.size(); ++i) {
     242            3 :         bqs::QueueRoute* const tmp = reinterpret_cast<bqs::QueueRoute*>(static_cast<uint8_t*>(devPtr) + offset);
     243            3 :         tmp->srcId = qRouteList->routeList[i].srcId;
     244            3 :         tmp->dstId = qRouteList->routeList[i].dstId;
     245            3 :         offset += sizeof(bqs::QueueRoute);
     246              :     }
     247              :     // device need to use mbuff
     248            6 :     auto ret = rtMemQueueEnQueue(0, qsContactId_, mBuf);
     249            6 :     if (ret != RT_ERROR_NONE) {
     250            1 :         (void)rtMbufFree(mBuf);
     251            1 :         mBuf = nullptr;
     252            1 :         devPtr = nullptr;
     253            1 :         return ret;
     254              :     }
     255              : 
     256            5 :     bqs::QueueRouteList bqsBindUnbindMsg = {0U, 0U, {0}};
     257            5 :     eventSum.subeventId =
     258            5 :         isBind ? static_cast<uint32_t>(bqs::ACL_BIND_QUEUE) : static_cast<uint32_t>(bqs::ACL_UNBIND_QUEUE);
     259            5 :     eventSum.msgLen = sizeof(bqsBindUnbindMsg);
     260            5 :     eventSum.msg = reinterpret_cast<char_t*>(&bqsBindUnbindMsg);
     261            5 :     ret = rtEschedSubmitEventSync(0, &eventSum, &ack);
     262            5 :     eventSum.msgLen = 0U;
     263            5 :     eventSum.msg = nullptr;
     264            5 :     if (ret != RT_ERROR_NONE) {
     265            1 :         if (!isMbufEnhanced_) {
     266            1 :             (void)rtMbufFree(mBuf);
     267            1 :             mBuf = nullptr;
     268            1 :             devPtr = nullptr;
     269              :         }
     270            1 :         return ret;
     271              :     }
     272            4 :     if (isMbufEnhanced_) {
     273              :         // after event sync mbuf need to be dequeue to be used as mbuf can not be free by enqueue side
     274            0 :         ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(rtMemQueueDeQueue(0, qsContactId_, &mBuf), rtMemQueueDeQueue);
     275            0 :         (void)rtMbufGetBuffAddr(mBuf, &devPtr);
     276              :     }
     277            4 :     bqs::QsProcMsgRsp* const rsp = reinterpret_cast<bqs::QsProcMsgRsp*>(ack.buf);
     278            4 :     if (rsp->retCode != 0) {
     279            1 :         ACL_LOG_INNER_ERROR("send connect qs failed, ret code is %d", rsp->retCode);
     280            1 :         (void)rtMbufFree(mBuf);
     281            1 :         mBuf = nullptr;
     282            1 :         devPtr = nullptr;
     283            1 :         return ACL_ERROR_FAILURE;
     284              :     }
     285            3 :     offset = sizeof(bqs::QsRouteHead);
     286            4 :     for (size_t i = 0UL; i < qRouteList->routeList.size(); ++i) {
     287            1 :         bqs::QueueRoute* const tmp = reinterpret_cast<bqs::QueueRoute*>(static_cast<uint8_t*>(devPtr) + offset);
     288            1 :         qRouteList->routeList[i].status = tmp->status;
     289            1 :         ACL_LOG_INFO(
     290              :             "route %zu, srcqid is %u, dst pid is %u, status is %d", i, qRouteList->routeList[i].srcId,
     291              :             qRouteList->routeList[i].dstId, qRouteList->routeList[i].status);
     292            1 :         offset += sizeof(bqs::QueueRoute);
     293              :     }
     294            3 :     (void)rtMbufFree(mBuf);
     295            3 :     devPtr = nullptr;
     296            3 :     mBuf = nullptr;
     297            3 :     return ret;
     298              : }
     299              : 
     300            7 : aclError QueueProcessor::GetQueueRouteNum(
     301              :     const acltdtQueueRouteQueryInfo* const queryInfo, const int32_t deviceId, rtEschedEventSummary_t& eventSum,
     302              :     rtEschedEventReply_t& ack, size_t& routeNum) const
     303              : {
     304            7 :     ACL_LOG_INFO("start to get queue route num");
     305            7 :     bqs::QueueRouteQuery routeQuery = {0UL, 0U, 0U, 0U, 0, 0, 0UL, {0}};
     306            7 :     routeQuery.queryType = static_cast<uint32_t>(queryInfo->mode);
     307            7 :     routeQuery.srcId = queryInfo->srcId;
     308            7 :     routeQuery.dstId = queryInfo->dstId;
     309              : 
     310            7 :     eventSum.subeventId = bqs::ACL_QUERY_QUEUE_NUM;
     311            7 :     eventSum.msgLen = sizeof(routeQuery);
     312            7 :     eventSum.msg = reinterpret_cast<char_t*>(&routeQuery);
     313            7 :     const rtError_t ret = rtEschedSubmitEventSync(deviceId, &eventSum, &ack);
     314            7 :     eventSum.msgLen = 0U;
     315            7 :     eventSum.msg = nullptr;
     316            7 :     ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(ret, rtEschedSubmitEventSync);
     317            7 :     bqs::QsProcMsgRsp* const rsp = reinterpret_cast<bqs::QsProcMsgRsp*>(ack.buf);
     318            7 :     if (rsp->retCode != 0) {
     319            1 :         ACL_LOG_INNER_ERROR("get queue route num failed, ret code is %d", rsp->retCode);
     320            1 :         return ACL_ERROR_FAILURE;
     321              :     }
     322            6 :     routeNum = rsp->retValue;
     323            6 :     ACL_LOG_INFO("successfully to get queue route num %zu.", routeNum);
     324            6 :     return ACL_SUCCESS;
     325              : }
     326              : 
     327            6 : aclError QueueProcessor::QueryQueueRoutesOnDevice(
     328              :     const acltdtQueueRouteQueryInfo* const queryInfo, const size_t routeNum, rtEschedEventSummary_t& eventSum,
     329              :     rtEschedEventReply_t& ack, acltdtQueueRouteList* const qRouteList) const
     330              : {
     331            6 :     ACL_LOG_INFO("start to query queue route %zu", routeNum);
     332            6 :     if (routeNum == 0U) {
     333            2 :         return ACL_SUCCESS;
     334              :     }
     335            4 :     const size_t routeSize =
     336            4 :         sizeof(bqs::QsRouteHead) + sizeof(bqs::QueueRouteQuery) + (routeNum * sizeof(bqs::QueueRoute));
     337            4 :     ACL_LOG_INFO("route size is %zu, queue route num is %zu", routeSize, qRouteList->routeList.size());
     338            4 :     void* devPtr = nullptr;
     339            4 :     void* mBuf = nullptr;
     340            4 :     ACL_REQUIRES_OK(AllocMBufOnDevice(&devPtr, &mBuf, routeSize));
     341            4 :     bqs::QsRouteHead* const head = reinterpret_cast<bqs::QsRouteHead*>(devPtr);
     342            4 :     head->length = routeSize;
     343            4 :     head->routeNum = routeNum;
     344            4 :     head->subEventId = bqs::ACL_QUERY_QUEUE;
     345            4 :     bqs::QueueRouteQuery* const routeQuery =
     346            4 :         reinterpret_cast<bqs::QueueRouteQuery*>(static_cast<uint8_t*>(devPtr) + sizeof(bqs::QsRouteHead));
     347            4 :     routeQuery->queryType = static_cast<uint32_t>(queryInfo->mode);
     348            4 :     routeQuery->srcId = queryInfo->srcId;
     349            4 :     routeQuery->dstId = queryInfo->dstId;
     350              :     // device need to use mbuff
     351            4 :     auto ret = rtMemQueueEnQueue(0, qsContactId_, mBuf);
     352            4 :     if (ret != RT_ERROR_NONE) {
     353            1 :         (void)rtMbufFree(mBuf);
     354            1 :         devPtr = nullptr;
     355            1 :         mBuf = nullptr;
     356            1 :         return ret;
     357              :     }
     358            3 :     bqs::QueueRouteList qsCommonMsg = {0U, 0U, {0}};
     359            3 :     eventSum.subeventId = bqs::ACL_QUERY_QUEUE;
     360            3 :     eventSum.msgLen = sizeof(qsCommonMsg);
     361            3 :     eventSum.msg = reinterpret_cast<char_t*>(&qsCommonMsg);
     362              : 
     363            3 :     ret = rtEschedSubmitEventSync(0, &eventSum, &ack);
     364            3 :     eventSum.msgLen = 0U;
     365            3 :     eventSum.msg = nullptr;
     366            3 :     if (ret != RT_ERROR_NONE) {
     367            1 :         if (!isMbufEnhanced_) {
     368            1 :             (void)rtMbufFree(mBuf);
     369            1 :             mBuf = nullptr;
     370            1 :             devPtr = nullptr;
     371              :         }
     372            1 :         return ret;
     373              :     }
     374            2 :     if (isMbufEnhanced_) {
     375              :         // after event sync mbuf need to be dequeue to be used as mbuf can not be free by enqueue side
     376            0 :         ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(rtMemQueueDeQueue(0, qsContactId_, &mBuf), rtMemQueueDeQueue);
     377            0 :         (void)rtMbufGetBuffAddr(mBuf, &devPtr);
     378              :     }
     379              : 
     380            2 :     bqs::QsProcMsgRsp* const rsp = reinterpret_cast<bqs::QsProcMsgRsp*>(ack.buf);
     381            2 :     if (rsp->retCode != 0) {
     382            1 :         ACL_LOG_INNER_ERROR("Failed to query queue routes, ret code is %d.", rsp->retCode);
     383            1 :         (void)rtMbufFree(mBuf);
     384            1 :         devPtr = nullptr;
     385            1 :         mBuf = nullptr;
     386            1 :         return ACL_ERROR_FAILURE;
     387              :     }
     388            1 :     size_t offset = sizeof(bqs::QsRouteHead) + sizeof(bqs::QueueRouteQuery);
     389            2 :     for (size_t i = 0UL; i < routeNum; ++i) {
     390            1 :         bqs::QueueRoute* const tmp = reinterpret_cast<bqs::QueueRoute*>(static_cast<uint8_t*>(devPtr) + offset);
     391            1 :         const acltdtQueueRoute tmpQueueRoute = {tmp->srcId, tmp->dstId, tmp->status};
     392            1 :         qRouteList->routeList.push_back(tmpQueueRoute);
     393            1 :         ACL_LOG_INFO(
     394              :             "route %zu, srcqid is %u, dst pid is %u, status is %d", i, qRouteList->routeList[i].srcId,
     395              :             qRouteList->routeList[i].dstId, qRouteList->routeList[i].status);
     396            1 :         offset += sizeof(bqs::QueueRoute);
     397              :     }
     398            1 :     (void)rtMbufFree(mBuf);
     399            1 :     devPtr = nullptr;
     400            1 :     mBuf = nullptr;
     401            1 :     ACL_LOG_INFO("Successfully to execute acltdtQueryQueueRoutes, queue route is %zu", qRouteList->routeList.size());
     402            1 :     return ACL_SUCCESS;
     403              : }
     404              : 
     405            5 : aclError QueueProcessor::QueryAllocGroup()
     406              : {
     407            5 :     ACL_LOG_INFO("Start to QueryAllocGroup.");
     408              :     static bool isGroupQuery = false;
     409            5 :     if (isGroupQuery) {
     410            4 :         return ACL_SUCCESS;
     411              :     }
     412            1 :     const std::lock_guard<std::mutex> lk(muForQueryGroup_);
     413            1 :     if (isGroupQuery) {
     414            0 :         return ACL_SUCCESS;
     415              :     }
     416              :     size_t grpNum;
     417              :     uint32_t alloc;
     418            1 :     std::string grpName;
     419            1 :     const auto pid = mmGetPid();
     420            1 :     rtMemGrpQueryInput_t input = {};
     421            1 :     input.cmd = RT_MEM_GRP_QUERY_GROUPS_OF_PROCESS;
     422            1 :     input.grpQueryByProc.pid = pid;
     423            1 :     rtMemGrpQueryOutput_t output = {};
     424            1 :     rtMemGrpOfProc_t outputInfo[QUERY_BUFF_GRP_MAX_NUM] = {{}};
     425            1 :     output.groupsOfProc = outputInfo;
     426            1 :     output.maxNum = QUERY_BUFF_GRP_MAX_NUM;
     427              : 
     428            1 :     ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(rtMemGrpQuery(&input, &output), rtMemGrpQuery);
     429            1 :     grpNum = output.resultNum;
     430            1 :     if ((grpNum == 0U) || (output.groupsOfProc == nullptr)) {
     431            0 :         ACL_LOG_ERROR("[Check] grpNum is zero or groupsOfProc is nullptr, grpNum is %zu", grpNum);
     432            0 :         return ACL_ERROR_FAILURE;
     433              :     }
     434            1 :     for (size_t num = 0U; num < grpNum; num++) {
     435            1 :         alloc = output.groupsOfProc[num].attr.alloc;
     436            1 :         grpName = std::string(output.groupsOfProc[num].groupName);
     437            1 :         ACL_LOG_INFO("This proc [%d] has [%zu] group, alloc is %u, name is %s", pid, grpNum, alloc, grpName.c_str());
     438            1 :         if (alloc == 1U) {
     439            1 :             ACL_REQUIRES_OK(QueryGroupId(grpName));
     440            1 :             isGroupQuery = true;
     441            1 :             return ACL_SUCCESS;
     442              :         }
     443              :     }
     444            0 :     ACL_LOG_ERROR("[Check] has no alloc");
     445            0 :     return ACL_ERROR_FAILURE;
     446            1 : }
     447              : 
     448            2 : aclError QueueProcessor::QueryGroupId(const std::string& grpName)
     449              : {
     450            2 :     ACL_LOG_INFO("Query groupId from name = %s", grpName.c_str());
     451            2 :     rtMemGrpQueryInput_t input = {};
     452            2 :     input.cmd = RT_MEM_GRP_QUERY_GROUP_ID;
     453              :     const auto strcpyRet =
     454            2 :         strcpy_s(input.grpQueryGroupId.grpName, sizeof(input.grpQueryGroupId.grpName), grpName.c_str());
     455            2 :     if (strcpyRet != EOK) {
     456            0 :         ACL_LOG_INNER_ERROR("[strcpy]copy group name to input failed, result = %d.", strcpyRet);
     457            0 :         return ACL_ERROR_FAILURE;
     458              :     }
     459            2 :     rtMemGrpQueryOutput_t output = {};
     460            2 :     rtMemGrpQueryGroupIdInfo_t outputInfo = {};
     461            2 :     output.groupIdInfo = &outputInfo;
     462              : 
     463            2 :     ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(rtMemGrpQuery(&input, &output), rtMemGrpQuery);
     464            2 :     qsGroupId_ = output.groupIdInfo->groupId;
     465            2 :     ACL_LOG_INFO("This groupId is %d, name is %s", qsGroupId_, grpName.c_str());
     466            2 :     return ACL_SUCCESS;
     467              : }
     468              : 
     469            6 : aclError QueueProcessor::acltdtAllocBufData(const size_t size, const uint32_t type, acltdtBuf* const buf)
     470              : {
     471              :     rtError_t ret;
     472            6 :     if (!isMbufInit_) {
     473            1 :         rtMemBuffCfg_t cfg = {{}};
     474            1 :         ret = rtMbufInit(&cfg);
     475            1 :         if ((ret != ACL_RT_SUCCESS) && (ret != ACL_ERROR_RT_REPEATED_INIT)) {
     476            0 :             return ret;
     477              :         }
     478            1 :         isMbufInit_ = true;
     479              :     }
     480            6 :     ACL_REQUIRES_RTS_OK(rtMbufAllocEx(buf, size, type, qsGroupId_));
     481            6 :     return ACL_SUCCESS;
     482              : }
     483              : 
     484            6 : aclError QueueProcessor::acltdtFreeBuf(acltdtBuf buf) { return aclrtFreeBuf(buf); }
     485              : 
     486            1 : aclError QueueProcessor::acltdtSetBufDataLen(const acltdtBuf buf, const size_t len)
     487              : {
     488            1 :     return aclrtSetBufDataLen(buf, len);
     489              : }
     490              : 
     491            1 : aclError QueueProcessor::acltdtGetBufDataLen(const acltdtBuf buf, size_t* const len)
     492              : {
     493            1 :     return aclrtGetBufDataLen(buf, len);
     494              : }
     495              : 
     496            4 : aclError QueueProcessor::acltdtGetBufData(const acltdtBuf buf, void** const dataPtr, size_t* const size)
     497              : {
     498            4 :     return aclrtGetBufData(buf, dataPtr, size);
     499              : }
     500              : 
     501            1 : aclError QueueProcessor::acltdtGetBufUserData(
     502              :     const acltdtBuf buf, void* dataPtr, const size_t size, const size_t offset)
     503              : {
     504            1 :     return aclrtGetBufUserData(buf, dataPtr, size, offset);
     505              : }
     506              : 
     507            1 : aclError QueueProcessor::acltdtSetBufUserData(
     508              :     acltdtBuf buf, const void* dataPtr, const size_t size, const size_t offset)
     509              : {
     510            1 :     return aclrtSetBufUserData(buf, dataPtr, size, offset);
     511              : }
     512              : 
     513            1 : aclError QueueProcessor::acltdtCopyBufRef(const acltdtBuf buf, acltdtBuf* const newBuf)
     514              : {
     515            1 :     return aclrtCopyBufRef(buf, newBuf);
     516              : }
     517              : 
     518            1 : aclError QueueProcessor::acltdtAppendBufChain(const acltdtBuf headBuf, const acltdtBuf buf)
     519              : {
     520            1 :     return aclrtAppendBufChain(headBuf, buf);
     521              : }
     522              : 
     523            1 : aclError QueueProcessor::acltdtGetBufChainNum(const acltdtBuf headBuf, uint32_t* const num)
     524              : {
     525            1 :     return aclrtGetBufChainNum(headBuf, num);
     526              : }
     527              : 
     528            1 : aclError QueueProcessor::acltdtGetBufFromChain(const acltdtBuf headBuf, const uint32_t index, acltdtBuf* const buf)
     529              : {
     530            1 :     return aclrtGetBufFromChain(headBuf, index, buf);
     531              : }
     532              : 
     533           19 : QueueDataMutexPtr QueueProcessor::GetMutexForData(const uint32_t qid)
     534              : {
     535           19 :     const std::lock_guard<std::mutex> lk(muForQueueMap_);
     536           19 :     const auto it = muForQueue_.find(qid);
     537           19 :     if (it != muForQueue_.end()) {
     538           13 :         return it->second;
     539              :     } else {
     540            6 :         const QueueDataMutexPtr queueDataMutex = std::make_shared<QueueDataMutex>();
     541            6 :         muForQueue_[qid] = queueDataMutex;
     542            6 :         return queueDataMutex;
     543            6 :     }
     544           19 : }
     545              : 
     546            3 : void QueueProcessor::DeleteMutexForData(const uint32_t qid)
     547              : {
     548            3 :     const std::lock_guard<std::mutex> lk(muForQueueMap_);
     549            3 :     const auto it = muForQueue_.find(qid);
     550            3 :     if (it != muForQueue_.end()) {
     551            0 :         (void)muForQueue_.erase(it);
     552              :     }
     553            6 :     return;
     554            3 : }
     555              : 
     556         1085 : uint64_t QueueProcessor::GetTimestamp() const
     557              : {
     558         1085 :     mmTimeval tv{};
     559         1085 :     const auto ret = mmGetTimeOfDay(&tv, nullptr);
     560         1085 :     if (ret != EN_OK) {
     561            0 :         ACL_LOG_WARN("Func mmGetTimeOfDay did not return success, errorCode = %d", ret);
     562              :     }
     563              :     // 1000000: seconds to microseconds
     564         1085 :     const uint64_t totalUseTime = static_cast<size_t>(tv.tv_usec) + (static_cast<uint64_t>(tv.tv_sec) * 1000000UL);
     565         1085 :     return totalUseTime;
     566              : }
     567              : 
     568           13 : aclError QueueProcessor::GetDeviceId(int32_t& deviceId) const
     569              : {
     570           13 :     deviceId = 0;
     571           13 :     aclrtRunMode aclRunMode = ACL_HOST;
     572           13 :     const aclError getRunModeRet = aclrtGetRunMode(&aclRunMode);
     573           13 :     if (getRunModeRet != ACL_SUCCESS) {
     574            0 :         ACL_LOG_CALL_ERROR("[Get][RunMode]get run mode failed, errorCode = %d.", getRunModeRet);
     575            0 :         return getRunModeRet;
     576              :     }
     577              : 
     578           13 :     if (aclRunMode == ACL_HOST) {
     579            0 :         ACL_REQUIRES_RTS_OK(rtGetDevice(&deviceId));
     580              :     }
     581           13 :     return ACL_SUCCESS;
     582              : }
     583              : 
     584           10 : aclError QueueProcessor::acltdtEnqueueData(
     585              :     const uint32_t qid, const void* const data, const size_t dataSize, const void* const userData,
     586              :     const size_t userDataSize, const int32_t timeout, const uint32_t rsv)
     587              : {
     588           10 :     ACL_LOG_INFO(
     589              :         "Start to enqueue data qid is %u, dataSize is %zu, userDataSize is %zu, "
     590              :         "timeout is %d, rsv is %u",
     591              :         qid, dataSize, userDataSize, timeout, rsv);
     592           10 :     ACL_REQUIRES_NOT_NULL_WITH_INPUT_REPORT(data);
     593           10 :     ACL_REQUIRES_POSITIVE_REPORT(dataSize);
     594              : 
     595              :     int32_t deviceId;
     596           10 :     ACL_REQUIRES_OK(GetDeviceId(deviceId));
     597           10 :     const QueueDataMutexPtr muPtr = GetMutexForData(qid);
     598           10 :     ACL_CHECK_MALLOC_RESULT(muPtr);
     599           10 :     const std::lock_guard<std::mutex> lk(muPtr->muForDequeue);
     600              : 
     601           10 :     rtMemQueueBuff_t queueBuf = {nullptr, 0U, nullptr, 0U};
     602           10 :     rtMemQueueBuffInfo queueBufInfo = {const_cast<void*>(data), dataSize};
     603           10 :     queueBuf.buffCount = 1U;
     604           10 :     queueBuf.buffInfo = &queueBufInfo;
     605           10 :     queueBuf.contextAddr = const_cast<void*>(userData);
     606           10 :     queueBuf.contextLen = userDataSize;
     607              : 
     608           10 :     const rtError_t ret = rtMemQueueEnQueueBuff(deviceId, qid, &queueBuf, timeout);
     609           10 :     if (ret == ACL_ERROR_RT_QUEUE_FULL) {
     610            3 :         ACL_LOG_INFO("queue is full, device is %d, qid is %u", deviceId, qid);
     611            3 :         return ret;
     612              :     }
     613            7 :     if (ret != RT_ERROR_NONE) {
     614            3 :         return ret;
     615              :     }
     616              : 
     617            4 :     ACL_LOG_INFO("success to execute acltdtEnqueueData, device is %d, qid is %u", deviceId, qid);
     618            4 :     return ACL_SUCCESS;
     619           10 : }
     620              : 
     621            3 : aclError QueueProcessor::acltdtDequeueData(
     622              :     const uint32_t qid, void* const data, const size_t dataSize, size_t* const retDataSize, void* const userData,
     623              :     const size_t userDataSize, const int32_t timeout)
     624              : {
     625            3 :     ACL_LOG_INFO(
     626              :         "Start to dequeue data qid is %u, dataSize is %zu, userDataSize is %zu, "
     627              :         "timeout is %d,",
     628              :         qid, dataSize, userDataSize, timeout);
     629            3 :     ACL_REQUIRES_NOT_NULL_WITH_INPUT_REPORT(data);
     630            3 :     ACL_REQUIRES_NOT_NULL_WITH_INPUT_REPORT(retDataSize);
     631            3 :     ACL_REQUIRES_POSITIVE_REPORT(dataSize);
     632              : 
     633              :     int32_t deviceId;
     634            3 :     ACL_REQUIRES_OK(GetDeviceId(deviceId));
     635              : 
     636            3 :     const QueueDataMutexPtr muPtr = GetMutexForData(qid);
     637            3 :     ACL_CHECK_MALLOC_RESULT(muPtr);
     638            3 :     const std::lock_guard<std::mutex> lk(muPtr->muForEnqueue);
     639              : 
     640            3 :     rtError_t ret = rtMemQueuePeek(deviceId, qid, retDataSize, timeout);
     641            3 :     if (ret == ACL_ERROR_RT_QUEUE_EMPTY) {
     642            1 :         ACL_LOG_INFO("queue is empty, device is %d, qid is %u", deviceId, qid);
     643            1 :         return ret;
     644              :     }
     645            2 :     if (ret != RT_ERROR_NONE) {
     646            1 :         ACL_LOG_ERROR("peek queue [%u] failed, device is %d", qid, deviceId);
     647            1 :         return ret;
     648              :     }
     649              : 
     650            1 :     rtMemQueueBuff_t queueBuf = {nullptr, 0U, nullptr, 0U};
     651            1 :     rtMemQueueBuffInfo queueBufInfo = {data, dataSize};
     652            1 :     queueBuf.buffCount = 1U;
     653            1 :     queueBuf.buffInfo = &queueBufInfo;
     654            1 :     queueBuf.contextAddr = userData;
     655            1 :     queueBuf.contextLen = userDataSize;
     656              : 
     657            1 :     ret = rtMemQueueDeQueueBuff(deviceId, qid, &queueBuf, timeout);
     658            1 :     if (ret == ACL_ERROR_RT_QUEUE_EMPTY) {
     659            0 :         ACL_LOG_INFO("queue is empty, device is %d, qid is %u", deviceId, qid);
     660            0 :         return ret;
     661              :     }
     662              : 
     663            1 :     if (ret != RT_ERROR_NONE) {
     664            0 :         ACL_LOG_ERROR("Failed to rtMemQueueDeQueueBuf, device is %d, qid is %u.", deviceId, qid);
     665            0 :         return ret;
     666              :     }
     667              : 
     668            1 :     ACL_LOG_INFO(
     669              :         "success to execute acltdtDequeueData, device is %d, qid is %u, retDataSize is %zu", deviceId, qid,
     670              :         *retDataSize);
     671            1 :     return ACL_SUCCESS;
     672            3 : }
     673              : 
     674            3 : void QueueProcessor::acltdtSetDefaultQueueAttr(acltdtQueueAttr& attr)
     675              : {
     676            3 :     (void)memset_s(attr.name, static_cast<size_t>(RT_MQ_MAX_NAME_LEN), 0, sizeof(attr.name));
     677            3 :     attr.depth = RT_MQ_DEPTH_DEFAULT;
     678            3 :     attr.workMode = static_cast<uint32_t>(RT_MQ_MODE_DEFAULT);
     679            3 :     attr.flowCtrlFlag = false;
     680            3 :     attr.flowCtrlDropTime = 0U;
     681            3 :     attr.overWriteFlag = false;
     682            3 :     return;
     683              : }
     684              : 
     685            4 : aclError QueueProcessor::acltdtCreateQueueWithAttr(
     686              :     const int32_t deviceId, const acltdtQueueAttr* const attr, uint32_t* const qid) const
     687              : {
     688            4 :     if (attr == nullptr) {
     689            2 :         acltdtQueueAttr tmpAttr{};
     690            2 :         acltdtSetDefaultQueueAttr(tmpAttr);
     691            2 :         ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(rtMemQueueCreate(deviceId, &tmpAttr, qid), rtMemQueueCreate);
     692              :     } else {
     693            2 :         ACL_REQUIRES_RTS_OK_WARN_NOT_SUPPORT(rtMemQueueCreate(deviceId, attr, qid), rtMemQueueCreate);
     694              :     }
     695            4 :     return ACL_SUCCESS;
     696              : }
     697              : } // namespace acl
        

Generated by: LCOV version 2.0-1