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

Generated by: LCOV version 1.14