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

Generated by: LCOV version 2.0-1