LCOV - code coverage report
Current view: top level - server - router_server.cpp (source / functions) Coverage Total Hit
Test: coverage.info Lines: 87.6 % 800 701
Test Date: 2026-08-12 11:05:07 Functions: 97.6 % 42 41

            Line data    Source code
       1              : /**
       2              :  * Copyright (c) 2025 Huawei Technologies Co., Ltd.
       3              :  * This program is free software, you can redistribute it and/or modify it under the terms and conditions of
       4              :  * CANN Open Software License Agreement Version 2.0 (the "License").
       5              :  * Please refer to the License for details. You may not use this file except in compliance with the License.
       6              :  * THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
       7              :  * INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE.
       8              :  * See LICENSE in the root of the software repository for the full text of the License.
       9              :  */
      10              : 
      11              : #include "router_server.h"
      12              : 
      13              : #include <map>
      14              : #include <csignal>
      15              : #include <sstream>
      16              : #include <securec.h>
      17              : #include <sys/types.h>
      18              : #include <unistd.h>
      19              : #include "statistic_manager.h"
      20              : #include "bqs_log.h"
      21              : #include "common/bqs_util.h"
      22              : #include "hccl_process.h"
      23              : #include "hccl/hccl_so_manager.h"
      24              : #include "common/type_def.h"
      25              : #include "queue_manager.h"
      26              : #include "queue_schedule_hal_interface_ref.h"
      27              : #include "qs_interface_process.h"
      28              : #include "queue_schedule_sub_module_interface.h"
      29              : #include "entity_manager.h"
      30              : 
      31              : namespace bqs {
      32              : namespace {
      33              : const std::string PIPELINE_QUEUE_NAME = "QsPipeQueue";
      34              : constexpr const uint16_t MAJOR_VERSION = 3U;
      35              : constexpr const uint16_t MINOR_VERSION = 0U;
      36              : constexpr const QueueShareAttr ADMIN_QUEUE_ATTR = {1U, 1U, 1U, 0U};
      37              : constexpr const char_t* ROUTER_SERVER_THREAD_NAME_PREFIX = "router_server";
      38              : 
      39              : // mapping of request subeventId and response subeventId
      40              : const std::map<int32_t, int32_t> g_reqRspMapping = {
      41              :     {AICPU_BIND_QUEUE, AICPU_BIND_QUEUE_RES},
      42              :     {AICPU_BIND_QUEUE_INIT, AICPU_BIND_QUEUE_INIT_RES},
      43              :     {AICPU_UNBIND_QUEUE, AICPU_UNBIND_QUEUE_RES},
      44              :     {AICPU_QUERY_QUEUE, AICPU_QUERY_QUEUE_RES},
      45              :     {AICPU_QUERY_QUEUE_NUM, AICPU_QUERY_QUEUE_NUM_RES}};
      46              : // mapping of request subeventId and qs operate type
      47              : const std::map<uint32_t, QsOperType> g_qsOperation = {
      48              :     {static_cast<uint32_t>(AICPU_BIND_QUEUE_INIT), QsOperType::BIND_INIT},
      49              :     {static_cast<uint32_t>(AICPU_BIND_QUEUE_INIT), QsOperType::BIND_INIT},
      50              :     {static_cast<uint32_t>(AICPU_QUERY_QUEUE_NUM), QsOperType::QUERY_NUM},
      51              :     {static_cast<uint32_t>(AICPU_QUEUE_RELATION_PROCESS), QsOperType::RELATION_PROCESS},
      52              :     {static_cast<uint32_t>(ACL_BIND_QUEUE_INIT), QsOperType::BIND_INIT},
      53              :     {static_cast<uint32_t>(ACL_BIND_QUEUE), QsOperType::RELATION_PROCESS},
      54              :     {static_cast<uint32_t>(ACL_UNBIND_QUEUE), QsOperType::RELATION_PROCESS},
      55              :     {static_cast<uint32_t>(ACL_QUERY_QUEUE_NUM), QsOperType::QUERY_NUM},
      56              :     {static_cast<uint32_t>(ACL_QUERY_QUEUE), QsOperType::RELATION_PROCESS},
      57              :     {static_cast<uint32_t>(DGW_CREATE_HCOM_HANDLE), QsOperType::CREATE_HCOM_HANDLE},
      58              :     {static_cast<uint32_t>(DGW_DESTORY_HCOM_HANDLE), QsOperType::DESTROY_HCOM_HANDLE},
      59              :     {static_cast<uint32_t>(BIND_HOSTPID), QsOperType::BIND_HOST_PID},
      60              :     // new operation type for flow gateway
      61              :     {static_cast<uint32_t>(UPDATE_CONFIG), QsOperType::UPDATE_CONFIG},
      62              :     {static_cast<uint32_t>(QUERY_CONFIG_NUM), QsOperType::QUERY_CONFIG_NUM},
      63              :     {static_cast<uint32_t>(QUERY_CONFIG), QsOperType::QUERY_CONFIG},
      64              :     {static_cast<uint32_t>(QUERY_LINKSTATUS), QsOperType::QUERY_LINKSTATUS},
      65              :     {static_cast<uint32_t>(QUERY_LINKSTATUS_V2), QsOperType::QUERY_LINKSTATUS_V2},
      66              : };
      67              : 
      68              : // 老版驱动结构体——解析数据时需要使用与驱动相同的结构体
      69              : struct old_event_sync_msg {
      70              :     int pid;                     /* local pid */
      71              :     unsigned int dst_engine : 4; /* local engine */
      72              :     unsigned int gid : 6;
      73              :     unsigned int event_id : 6;
      74              :     unsigned int subevent_id : 16; /* Range: 0 ~ 4095 */
      75              :     char msg[];
      76              : };
      77              : 
      78              : } // namespace
      79              : 
      80            2 : RouterServer::RouterServer()
      81            2 :     : processing_(false),
      82            2 :       done_(false),
      83            2 :       processingExtra_(false),
      84            2 :       doneExtra_(true),
      85            2 :       bindQueueGroupId_(static_cast<uint32_t>(BIND_QUEUE_GROUP_ID)),
      86            2 :       running_(false),
      87            2 :       deviceId_(0U),
      88            2 :       srcPid_(-1),
      89            2 :       srcVersion_(0U),
      90            2 :       srcGroupId_(-1),
      91            2 :       pipelineQueueId_(MAX_QUEUE_ID_NUM),
      92            2 :       subEventId_(0U),
      93            2 :       deployMode_(QueueSchedulerRunMode::MULTI_PROCESS),
      94            2 :       retCode_(static_cast<int32_t>(BQS_STATUS_OK)),
      95            2 :       attachedFlag_(false),
      96            2 :       isAicpuEvent_(false),
      97            2 :       qsRouteListPtr_(nullptr),
      98            2 :       qsRouterHeadPtr_(nullptr),
      99            2 :       qsRouterQueryPtr_(nullptr),
     100            2 :       drvSyncMsg_(nullptr),
     101            2 :       aicpuRspHead_(0UL),
     102            2 :       f2nfGroupId_(0U),
     103            2 :       schedPolicy_(0UL),
     104            2 :       cfgInfoOperator_(nullptr),
     105            2 :       callHcclFlag_(false),
     106            2 :       numaFlag_(false),
     107            2 :       readyToHandleMsg_(false),
     108            2 :       manageThreadStatus_(ThreadStatus::NOT_INIT),
     109            2 :       needAttachGroup_(false),
     110            6 :       compatMsg_(false)
     111            2 : {}
     112              : 
     113          176 : void RouterServer::Destroy()
     114              : {
     115          176 :     if (!running_) {
     116          170 :         return;
     117              :     }
     118            6 :     BQS_LOG_INFO("[RouterServer]QS Server destroy.");
     119            6 :     running_ = false;
     120            6 :     queueRouteQueryList_.clear();
     121            6 :     cv_.notify_all();
     122            6 :     if (monitorQsEvent_.joinable()) {
     123            5 :         monitorQsEvent_.join();
     124              :     }
     125            6 :     manageThreadStatus_ = ThreadStatus::NOT_INIT;
     126            6 :     if (pipelineQueueId_ < MAX_QUEUE_ID_NUM) {
     127            3 :         const auto ret = halQueueDestroy(deviceId_, pipelineQueueId_);
     128            3 :         BQS_LOG_ERROR_WHEN(
     129              :             ret != DRV_ERROR_NONE, "[RouterServer]Destroy relation event buff queue error, queue id[%u], ret[%d]",
     130              :             pipelineQueueId_.load(), static_cast<int32_t>(ret));
     131              :     }
     132            6 :     cfgInfoOperator_ = nullptr;
     133            6 :     BQS_LOG_RUN_INFO("[RouterServer]QS Server finish destroy.");
     134              : }
     135              : 
     136            2 : RouterServer::~RouterServer() { Destroy(); }
     137              : 
     138            0 : bool RouterServer::GetCallHcclFlag() const { return callHcclFlag_; }
     139              : 
     140            3 : uint32_t RouterServer::GetPipelineQueueId() const { return pipelineQueueId_; }
     141              : 
     142          534 : RouterServer& RouterServer::GetInstance()
     143              : {
     144          534 :     static RouterServer instance;
     145          534 :     return instance;
     146              : }
     147              : 
     148          106 : void RouterServer::HandleBqsMsg(event_info& info)
     149              : {
     150          106 :     if (info.comm.event_id != EVENT_QS_MSG) {
     151            1 :         BQS_LOG_ERROR(
     152              :             "[RouterServer]Queue schedule does not support [%d] event", static_cast<int32_t>(info.comm.event_id));
     153            1 :         return;
     154              :     }
     155          105 :     subEventId_ = info.comm.subevent_id;
     156              :     // check aicpu event
     157          105 :     isAicpuEvent_ = (subEventId_ <= static_cast<uint32_t>(AICPU_RELATED_MESSAGE_SPLIT)) ? true : false;
     158          105 :     if (!isAicpuEvent_) {
     159              :         // msg is char array, not need check nullptr
     160           87 :         drvSyncMsg_ = PtrToPtr<char_t, event_sync_msg>(info.priv.msg);
     161              :     }
     162              : 
     163          105 :     if (!readyToHandleMsg_) {
     164            1 :         BQS_LOG_WARN("[RouterServer] is not ready to HandleBqsMsg");
     165            1 :         SendRspEvent(static_cast<int32_t>(BQS_STATUS_NOT_INIT));
     166            1 :         return;
     167              :     }
     168              : 
     169              :     // check subEventId, SendRspEvent need drvSyncMsg if event from acl
     170          104 :     const auto iter = g_qsOperation.find(subEventId_);
     171          104 :     if (iter == g_qsOperation.end()) {
     172            2 :         BQS_LOG_RUN_INFO("[RouterServer]SubEventId is invalid[%d]", subEventId_);
     173            2 :         SendRspEvent(static_cast<int32_t>(BQS_STATUS_PARAM_INVALID));
     174            2 :         return;
     175              :     }
     176              : 
     177              :     // aicpu message is not allowed in thread mode.
     178          102 :     if (isAicpuEvent_ && (deployMode_ == QueueSchedulerRunMode::MULTI_THREAD)) {
     179            1 :         BQS_LOG_ERROR(
     180              :             "[RouterServer]Thread mode[%u] does not sopport event[%u] from aicpu.", static_cast<int32_t>(deployMode_),
     181              :             subEventId_);
     182            1 :         return;
     183              :     }
     184          101 :     PreProcessEvent(info);
     185          101 :     BQS_LOG_INFO("[RouterServer]HandleBqsMsg end, isAicpuEvent[%d].", static_cast<int32_t>(isAicpuEvent_));
     186          101 :     return;
     187              : }
     188              : 
     189          106 : void RouterServer::PreProcessEvent(const event_info& info)
     190              : {
     191              :     // get event head for responseing by sync interface
     192          106 :     BQS_LOG_INFO("[RouterServer]PreProcess start operation[%u]", subEventId_);
     193          106 :     const auto iter = g_qsOperation.find(subEventId_);
     194          106 :     if (iter == g_qsOperation.end()) {
     195            1 :         SendRspEvent(static_cast<int32_t>(BQS_STATUS_PARAM_INVALID));
     196            1 :         BQS_LOG_ERROR("[RouterServer] RouterServer receive unsupported msg type:%d", subEventId_);
     197            1 :         return;
     198              :     }
     199              : 
     200          105 :     const QsOperType operType = iter->second;
     201          105 :     switch (operType) {
     202            4 :         case QsOperType::BIND_INIT:
     203            4 :             SendRspEvent(static_cast<int32_t>(ProcessBindInit(info)));
     204            4 :             break;
     205            4 :         case QsOperType::QUERY_NUM:
     206            4 :             StatisticManager::GetInstance().GetBindStat();
     207            4 :             ParseGetBindNumMsg(info);
     208            4 :             break;
     209           10 :         case QsOperType::RELATION_PROCESS: {
     210              :             // get message from mbuf, eventid also in mbuf
     211           10 :             Mbuf* mBuf = nullptr;
     212           10 :             const auto resultCode = ParseRelationInfo(&mBuf);
     213           10 :             if (resultCode != BQS_STATUS_OK) {
     214            1 :                 BQS_LOG_ERROR(
     215              :                     "[RouterServer]Get detail message from mbuf failed ret[%d].", static_cast<int32_t>(resultCode));
     216            1 :                 SendRspEvent(static_cast<int32_t>(resultCode));
     217            1 :                 if ((mBuf != nullptr) && (srcVersion_ != 0U)) {
     218            0 :                     const auto freeRet = halMbufFree(mBuf);
     219            0 :                     BQS_LOG_ERROR_WHEN(
     220              :                         freeRet != static_cast<int32_t>(DRV_ERROR_NONE), "Free mbuf failed, ret is %d", freeRet);
     221              :                 }
     222            1 :                 return;
     223              :             }
     224              : 
     225            9 :             BQS_LOG_INFO("[RouterServer]Start to process relation event[%u].", subEventId_);
     226            9 :             ProcessQueueRelationEvent(mBuf);
     227            9 :             break;
     228              :         }
     229           85 :         case QsOperType::UPDATE_CONFIG:
     230              :         case QsOperType::CREATE_HCOM_HANDLE:
     231              :         case QsOperType::DESTROY_HCOM_HANDLE:
     232              :         case QsOperType::QUERY_CONFIG:
     233              :         case QsOperType::QUERY_CONFIG_NUM: {
     234           85 :             ProcessConfigEvent(operType);
     235           85 :             break;
     236              :         }
     237            0 :         case QsOperType::QUERY_LINKSTATUS: {
     238            0 :             ProcessQueryLinkStatusEvent();
     239            0 :             break;
     240              :         }
     241            1 :         case QsOperType::QUERY_LINKSTATUS_V2: {
     242            1 :             ProcessQueryLinkStatusEvent();
     243            1 :             break;
     244              :         }
     245            1 :         default: {
     246            1 :             SendRspEvent(static_cast<int32_t>(BQS_STATUS_PARAM_INVALID));
     247            1 :             BQS_LOG_RUN_INFO("[RouterServer]RouterServer receive unsupported msg type:%u", subEventId_);
     248            1 :             break;
     249              :         }
     250              :     }
     251          104 :     return;
     252              : }
     253              : 
     254           87 : void RouterServer::ProcessConfigEvent(const QsOperType operType)
     255              : {
     256           87 :     void* mbuf = nullptr;
     257           87 :     auto drvRet = halQueueDeQueue(deviceId_, pipelineQueueId_, &mbuf);
     258           87 :     if ((drvRet != DRV_ERROR_NONE) || (mbuf == nullptr)) {
     259            2 :         BQS_LOG_ERROR(
     260              :             "halQueueDeQueue from queue[%u] in device[%u] failed, error[%d]", pipelineQueueId_.load(), deviceId_,
     261              :             static_cast<int32_t>(drvRet));
     262            1 :         SendRspEvent(static_cast<int32_t>(BQS_STATUS_DRIVER_ERROR));
     263            2 :         return;
     264              :     }
     265           86 :     auto resultCode = (cfgInfoOperator_ == nullptr) ?
     266              :                           BQS_STATUS_INNER_ERROR :
     267           85 :                           cfgInfoOperator_->ParseConfigEvent(subEventId_, pipelineQueueId_, mbuf, srcVersion_);
     268           86 :     if (resultCode == BQS_STATUS_WAIT) {
     269           37 :         resultCode = WaitSyncMsgProc();
     270              :     }
     271           86 :     if (operType == QsOperType::CREATE_HCOM_HANDLE) {
     272            8 :         callHcclFlag_ = true;
     273              :     }
     274              : 
     275           86 :     if (srcVersion_ != 0U) {
     276           86 :         BQS_LOG_INFO("Enque mbuf back");
     277           86 :         drvRet = halQueueEnQueue(deviceId_, pipelineQueueId_, mbuf);
     278           86 :         if (drvRet != DRV_ERROR_NONE) {
     279            2 :             BQS_LOG_ERROR(
     280              :                 "halQueueEnQueue into queue[%u] in device[%u] failed, error[%d]", pipelineQueueId_.load(), deviceId_,
     281              :                 static_cast<int32_t>(drvRet));
     282            1 :             SendRspEvent(static_cast<int32_t>(BQS_STATUS_DRIVER_ERROR));
     283            1 :             const auto freeRet = halMbufFree(PtrToPtr<void, Mbuf>(mbuf));
     284            1 :             BQS_LOG_ERROR_WHEN(freeRet != static_cast<int32_t>(DRV_ERROR_NONE), "Free mbuf failed, ret is %d", freeRet);
     285            1 :             return;
     286              :         }
     287              :     }
     288           85 :     BQS_LOG_INFO(
     289              :         "config for operate[%d] resultCode is %d.", static_cast<int32_t>(operType), static_cast<int32_t>(resultCode));
     290           85 :     SendRspEvent(static_cast<int32_t>(resultCode));
     291              : }
     292              : 
     293            2 : void RouterServer::ProcessQueryLinkStatusEvent()
     294              : {
     295            2 :     int32_t ret = static_cast<int32_t>(dgw::EntityManager::Instance(0U).CheckLinkStatus());
     296            2 :     if ((ret == 0) && numaFlag_) {
     297            0 :         ret = static_cast<int32_t>(dgw::EntityManager::Instance(1U).CheckLinkStatus());
     298              :     }
     299            2 :     SendRspEvent(ret);
     300            2 : }
     301              : 
     302            9 : void RouterServer::ProcessQueueRelationEvent(Mbuf* mbuf)
     303              : {
     304            9 :     BqsStatus ret = BQS_STATUS_INNER_ERROR;
     305            9 :     switch (subEventId_) {
     306            0 :         case AICPU_BIND_QUEUE: {
     307            0 :             StatisticManager::GetInstance().BindStat();
     308            0 :             ret = WaitBindMsgProc();
     309            0 :             break;
     310              :         }
     311            2 :         case ACL_BIND_QUEUE: {
     312            2 :             StatisticManager::GetInstance().BindStat();
     313            2 :             ret = WaitBindMsgProc();
     314            2 :             break;
     315              :         }
     316            1 :         case AICPU_UNBIND_QUEUE: {
     317            1 :             StatisticManager::GetInstance().UnbindStat();
     318            1 :             ret = WaitBindMsgProc();
     319            1 :             break;
     320              :         }
     321            0 :         case ACL_UNBIND_QUEUE: {
     322            0 :             StatisticManager::GetInstance().UnbindStat();
     323            0 :             ret = WaitBindMsgProc();
     324            0 :             break;
     325              :         }
     326            4 :         case AICPU_QUERY_QUEUE: {
     327            4 :             StatisticManager::GetInstance().GetBindStat();
     328            4 :             ret = ParseGetBindDetailMsg();
     329            4 :             break;
     330              :         }
     331            0 :         case ACL_QUERY_QUEUE: {
     332            0 :             StatisticManager::GetInstance().GetBindStat();
     333            0 :             ret = ParseGetBindDetailMsg();
     334            0 :             break;
     335              :         }
     336            2 :         default:
     337            2 :             BQS_LOG_ERROR("[RouterServer]unsupport subEventId[%u] in bind relation procedure", subEventId_);
     338            2 :             break;
     339              :     }
     340              : 
     341            9 :     if (srcVersion_ != 0U) {
     342            7 :         BQS_LOG_INFO("Enque mbuf back.");
     343            7 :         const auto drvRet = halQueueEnQueue(deviceId_, pipelineQueueId_, mbuf);
     344            7 :         if (drvRet != DRV_ERROR_NONE) {
     345            2 :             BQS_LOG_ERROR(
     346              :                 "halQueueEnQueue into queue[%u] in device[%u] failed, error[%d].", pipelineQueueId_.load(), deviceId_,
     347              :                 static_cast<int32_t>(drvRet));
     348            1 :             SendRspEvent(static_cast<int32_t>(BQS_STATUS_DRIVER_ERROR));
     349            1 :             const auto freeRet = halMbufFree(mbuf);
     350            1 :             BQS_LOG_ERROR_WHEN(freeRet != static_cast<int32_t>(DRV_ERROR_NONE), "Free mbuf failed, ret is %d", freeRet);
     351            1 :             return;
     352              :         }
     353              :     }
     354            8 :     SendRspEvent(static_cast<int32_t>(ret));
     355            8 :     return;
     356              : }
     357              : 
     358            3 : BqsStatus RouterServer::WaitBindMsgProc()
     359              : {
     360            3 :     BQS_LOG_INFO("[RouterServer]Bind relation [add/del], stage [wait]");
     361            3 :     auto ret = ParseBindUnbindMsg();
     362            3 :     if (ret == BQS_STATUS_OK) {
     363            0 :         ret = WaitSyncMsgProc();
     364              :     }
     365            3 :     BQS_LOG_INFO("[RouterServer]RouterServer WaitBindMsgProc end");
     366            3 :     return ret;
     367              : }
     368              : 
     369            2 : BqsStatus RouterServer::AttachAndInitGroup()
     370              : {
     371            2 :     BQS_LOG_INFO("[RouterServer]Attach and init group begin.");
     372            2 :     int32_t drvRet = 0;
     373              :     // 针对aicpusd 与qs合设的情况,如果qs以模块方式启动,则不需要重复加组,因为aicpusd已经加组
     374            2 :     if (qsInitGroupName_.empty() && (!SubModuleInterface::GetInstance().GetStartFlag())) {
     375            1 :         const std::unique_ptr<GroupQueryOutput> groupInfoPtr(new (std::nothrow) GroupQueryOutput());
     376            1 :         if (groupInfoPtr == nullptr) {
     377            0 :             BQS_LOG_ERROR("[RouterServer] Fail to allocate GroupQueryOutput");
     378            0 :             return BQS_STATUS_INNER_ERROR;
     379              :         }
     380            1 :         GroupQueryOutput& groupInfo = *(groupInfoPtr.get());
     381            1 :         uint32_t groupInfoLen = 0U;
     382            1 :         pid_t curPid = drvDeviceGetBareTgid();
     383              :         // query group info for current qs process
     384            1 :         drvRet = halGrpQuery(
     385              :             GRP_QUERY_GROUPS_OF_PROCESS, &curPid, static_cast<uint32_t>(sizeof(curPid)),
     386              :             reinterpret_cast<void*>(&groupInfo), &groupInfoLen);
     387            1 :         if (drvRet != static_cast<int32_t>(DRV_ERROR_NONE)) {
     388            0 :             BQS_LOG_ERROR("[RouterServer]halGrpQuery of qs[%d] failed before attached,ret[%d]", curPid, drvRet);
     389            0 :             return BQS_STATUS_DRIVER_ERROR;
     390              :         }
     391              :         // not in any group, cannot do attach process
     392            1 :         if (groupInfoLen == 0U) {
     393            0 :             BQS_LOG_ERROR("[RouterServer]QS should be add sharepool group before initial by aicpu or acl.");
     394            0 :             return BQS_STATUS_INNER_ERROR;
     395              :         }
     396            1 :         if ((groupInfoLen % sizeof(groupInfo.grpQueryGroupsOfProcInfo[0])) != 0U) {
     397            0 :             BQS_LOG_ERROR("[RouterServer]Group info size[%d] is invalid", groupInfoLen);
     398            0 :             return BQS_STATUS_DRIVER_ERROR;
     399              :         }
     400            1 :         const uint32_t groupNum = static_cast<uint32_t>(groupInfoLen / sizeof(groupInfo.grpQueryGroupsOfProcInfo[0]));
     401            2 :         for (uint32_t i = 0U; i < groupNum; ++i) {
     402              :             // attach and initial
     403            1 :             drvRet = halGrpAttach(groupInfo.grpQueryGroupsOfProcInfo[i].groupName, 0);
     404            1 :             if (drvRet != static_cast<int32_t>(DRV_ERROR_NONE)) {
     405            0 :                 BQS_LOG_ERROR(
     406              :                     "[RouterServer]Group[%s] attach failed for slave aicpusd[%d] ret[%d]",
     407              :                     groupInfo.grpQueryGroupsOfProcInfo[i].groupName, curPid, drvRet);
     408            0 :                 return BQS_STATUS_DRIVER_ERROR;
     409              :             }
     410            1 :             BQS_LOG_INFO(
     411              :                 "[RouterServer] halGrpAttach execute succ. group[%s] was attached by QS",
     412              :                 groupInfo.grpQueryGroupsOfProcInfo[i].groupName);
     413              :         }
     414            1 :     }
     415            2 :     attachedFlag_ = true;
     416            2 :     BuffCfg defaultCfg = {};
     417            2 :     drvRet = halBuffInit(&defaultCfg);
     418            2 :     if ((drvRet != static_cast<int32_t>(DRV_ERROR_NONE)) && (drvRet != static_cast<int32_t>(DRV_ERROR_REPEATED_INIT))) {
     419            1 :         BQS_LOG_ERROR("[RouterServer] Buffer initial failed for qs. ret[%d]", drvRet);
     420            1 :         return BQS_STATUS_DRIVER_ERROR;
     421              :     }
     422            1 :     BQS_LOG_INFO("[RouterServer] Buffer init success ret[%d]", drvRet);
     423            1 :     return BQS_STATUS_OK;
     424              : }
     425              : 
     426            2 : BqsStatus RouterServer::CreateAndGrantPipelineQueue()
     427              : {
     428            2 :     BQS_LOG_INFO("[RouterServer]Create and grant pipeline queue begin.");
     429              :     // do initial process
     430            2 :     const std::unique_lock<std::mutex> lk(mutex_);
     431            2 :     QueueAttr queAttr = {};
     432            2 :     std::string nameStr(PIPELINE_QUEUE_NAME);
     433            2 :     pid_t curPidTemp = 0;
     434            2 :     if (bqs::GetRunContext() == bqs::RunContext::HOST) {
     435            2 :         curPidTemp = getpid();
     436            2 :         queAttr.deploy_type = LOCAL_QUEUE_DEPLOY;
     437              :     } else {
     438            0 :         curPidTemp = drvDeviceGetBareTgid();
     439            0 :         queAttr.deploy_type = CLIENT_QUEUE_DEPLOY;
     440              :     }
     441            2 :     const uint32_t curPid = static_cast<uint32_t>(curPidTemp);
     442            2 :     nameStr += std::to_string(curPid);
     443              :     const auto memcpyRet =
     444            2 :         memcpy_s(queAttr.name, static_cast<uint32_t>(QUEUE_MAX_STR_LEN), nameStr.c_str(), nameStr.length() + 1UL);
     445            2 :     if (memcpyRet != EOK) {
     446            0 :         BQS_LOG_ERROR("[RouterServer]CreateAndGrantPipelineQueue memcpy_s failed, ret=%d.", memcpyRet);
     447            0 :         return BQS_STATUS_INNER_ERROR;
     448              :     }
     449            2 :     queAttr.depth = 2U;
     450            2 :     uint32_t queueId = 0U;
     451              :     // create queue
     452            2 :     auto drvRet = halQueueCreate(deviceId_, &queAttr, &queueId);
     453            2 :     if ((drvRet != DRV_ERROR_NONE) || (queueId >= MAX_QUEUE_ID_NUM)) {
     454            0 :         BQS_LOG_ERROR(
     455              :             "[RouterServer]Create queue[%s] error or qID[%u] is invalid, ret[%d]", PIPELINE_QUEUE_NAME.c_str(), queueId,
     456              :             static_cast<int32_t>(drvRet));
     457            0 :         return BQS_STATUS_DRIVER_ERROR;
     458              :     }
     459            2 :     drvRet = halQueueAttach(deviceId_, queueId, 0);
     460            2 :     if (drvRet != DRV_ERROR_NONE) {
     461            0 :         BQS_LOG_ERROR("Fail to attach queue[%ud], result[%d]", queueId, static_cast<int32_t>(drvRet));
     462            0 :         return BQS_STATUS_DRIVER_ERROR;
     463              :     }
     464            2 :     pipelineQueueId_ = queueId;
     465            2 :     if (deployMode_ == QueueSchedulerRunMode::MULTI_THREAD) {
     466            0 :         BQS_LOG_INFO("[RouterServer]Thread mode need not grant queue to other process");
     467            0 :         return BQS_STATUS_OK;
     468              :     }
     469              :     // grant pipeline queue to src process
     470              : 
     471            2 :     drvRet = halQueueGrant(deviceId_, static_cast<int32_t>(queueId), srcPid_, ADMIN_QUEUE_ATTR);
     472            2 :     if (drvRet != DRV_ERROR_NONE) {
     473            0 :         BQS_LOG_ERROR(
     474              :             "[RouterServer]Fail to add queue[%d] authority for aicpusd[%d], result[%d].", queueId, srcPid_,
     475              :             static_cast<int32_t>(drvRet));
     476            0 :         return BQS_STATUS_DRIVER_ERROR;
     477              :     }
     478            4 :     BQS_LOG_RUN_INFO("Success to init pipelineQ[%u].", pipelineQueueId_.load());
     479            2 :     return BQS_STATUS_OK;
     480            2 : }
     481              : 
     482            4 : BqsStatus RouterServer::ProcessBindInit(const event_info& info)
     483              : {
     484            4 :     if (info.priv.msg_len != sizeof(QsBindInit)) {
     485            0 :         BQS_LOG_ERROR("[RouterServer]Bind initial event message invalid, msgLen[%u]", info.priv.msg_len);
     486            0 :         return BQS_STATUS_PARAM_INVALID;
     487              :     }
     488              :     // bind initial already done
     489            4 :     if ((pipelineQueueId_ < MAX_QUEUE_ID_NUM) && (srcPid_ != -1)) {
     490            2 :         BQS_LOG_RUN_INFO("Pipeline queue already existed[%d], return pipelienQueueid", pipelineQueueId_.load());
     491            1 :         return BQS_STATUS_OK;
     492              :     }
     493            3 :     const QsBindInit* const bindInitMsg = reinterpret_cast<const QsBindInit*>(info.priv.msg);
     494            3 :     aicpuRspHead_ = bindInitMsg->syncEventHead;
     495            3 :     srcPid_ = bindInitMsg->pid;
     496            3 :     srcVersion_ = bindInitMsg->majorVersion;
     497            3 :     srcGroupId_ = static_cast<int32_t>(bindInitMsg->grpId);
     498            6 :     BQS_LOG_RUN_INFO(
     499              :         "[RouterServer]Get hostpid[%d] srcGroup[%d], srcVersion[%u]", srcPid_, srcGroupId_.load(), srcVersion_);
     500              : 
     501              :     // process mode need to attach group at first
     502            3 :     if (((deployMode_ == QueueSchedulerRunMode::SINGLE_PROCESS) ||
     503            3 :          (deployMode_ == QueueSchedulerRunMode::MULTI_PROCESS)) &&
     504            3 :         (!attachedFlag_)) {
     505            1 :         BQS_LOG_INFO("[RouterServer]start up attach and init group.");
     506            1 :         const auto attachRet = AttachAndInitGroup();
     507            1 :         if (attachRet != BQS_STATUS_OK) {
     508            0 :             BQS_LOG_ERROR("[RouterServer]AttachAndInitGroup failed, ret[%d].", static_cast<int32_t>(attachRet));
     509            0 :             return attachRet;
     510              :         }
     511              :     }
     512              : 
     513            3 :     if (qsInitGroupName_.empty()) {
     514            2 :         auto queueInitRet = QueueManager::GetInstance().InitQueue();
     515            2 :         if (queueInitRet != BQS_STATUS_OK) {
     516            1 :             BQS_LOG_ERROR("[RouterServer] Queue init failed");
     517            1 :             return queueInitRet;
     518              :         }
     519            1 :         if (numaFlag_) {
     520            0 :             queueInitRet = QueueManager::GetInstance().InitQueueExtra();
     521            0 :             if (queueInitRet != BQS_STATUS_OK) {
     522            0 :                 BQS_LOG_ERROR("[RouterServer] Queue init failed");
     523            0 :                 return queueInitRet;
     524              :             }
     525              :         }
     526              :     }
     527              : 
     528            2 :     const auto ret = CreateAndGrantPipelineQueue();
     529            2 :     if (ret != BQS_STATUS_OK) {
     530            0 :         return ret;
     531              :     }
     532            6 :     BQS_LOG_RUN_INFO(
     533              :         "First bind initial success, srcPid[%d], srvGroupId[%d], pipelineQueueId[%u]", srcPid_, srcGroupId_.load(),
     534              :         pipelineQueueId_.load());
     535            2 :     return BQS_STATUS_OK;
     536              : }
     537              : 
     538          112 : void RouterServer::FillRspContent(QsProcMsgRsp& retRsp, const int32_t resultCode)
     539              : {
     540              :     // only aicpu event need fill in aicpuRspHead
     541          112 :     retRsp.syncEventHead = isAicpuEvent_ ? aicpuRspHead_ : 0UL;
     542          112 :     retRsp.retCode = resultCode;
     543          112 :     retRsp.minorVersion = MINOR_VERSION;
     544          112 :     retRsp.majorVersion = MAJOR_VERSION;
     545          112 :     if ((subEventId_ == static_cast<uint32_t>(AICPU_BIND_QUEUE_INIT)) ||
     546          107 :         (subEventId_ == static_cast<uint32_t>(ACL_BIND_QUEUE_INIT))) {
     547              :         // init message return pipelineID
     548            5 :         retRsp.retValue =
     549            8 :             (resultCode == static_cast<int32_t>(BQS_STATUS_OK)) ? pipelineQueueId_.load() : MAX_QUEUE_ID_NUM;
     550              :     }
     551          112 :     if ((subEventId_ == static_cast<uint32_t>(AICPU_QUERY_QUEUE_NUM)) ||
     552          108 :         (subEventId_ == static_cast<uint32_t>(ACL_QUERY_QUEUE_NUM))) {
     553              :         // query num message return bind num
     554            8 :         retRsp.retValue = (resultCode != static_cast<int32_t>(BQS_STATUS_OK)) ?
     555              :                               0U :
     556            4 :                               static_cast<uint32_t>(queueRouteQueryList_.size());
     557            4 :         queueRouteQueryList_.clear();
     558              :     }
     559          112 :     if ((subEventId_ == static_cast<uint32_t>(AICPU_QUERY_QUEUE)) ||
     560          108 :         (subEventId_ == static_cast<uint32_t>(ACL_QUERY_QUEUE)) ||
     561          108 :         (subEventId_ == static_cast<uint32_t>(AICPU_BIND_QUEUE)) ||
     562          107 :         (subEventId_ == static_cast<uint32_t>(ACL_BIND_QUEUE)) ||
     563          105 :         (subEventId_ == static_cast<uint32_t>(AICPU_UNBIND_QUEUE)) ||
     564          104 :         (subEventId_ == static_cast<uint32_t>(ACL_UNBIND_QUEUE))) {
     565            8 :         retRsp.retValue = pipelineQueueId_;
     566              :     }
     567          112 :     return;
     568              : }
     569              : 
     570          112 : void RouterServer::SendRspEvent(const int32_t result)
     571              : {
     572          112 :     BQS_LOG_INFO("[RouterServer]Start to response message subeventid[%u]", subEventId_);
     573          112 :     QsProcMsgRsp retRsp = {};
     574          112 :     FillRspContent(retRsp, result);
     575          112 :     event_summary qsEvent = {};
     576          112 :     qsEvent.msg = PtrToPtr<QsProcMsgRsp, char_t>(&retRsp);
     577          112 :     qsEvent.msg_len = static_cast<uint32_t>(sizeof(retRsp));
     578          112 :     auto drvRet = DRV_ERROR_NONE;
     579          112 :     if (!isAicpuEvent_) {
     580           87 :         BQS_LOG_INFO("[RouterServer] Do ACL response");
     581           87 :         qsEvent.dst_engine = compatMsg_ ? PtrToPtr<event_sync_msg, old_event_sync_msg>(drvSyncMsg_)->dst_engine :
     582            0 :                                           drvSyncMsg_->dst_engine;
     583           87 :         qsEvent.policy = ONLY;
     584           87 :         qsEvent.pid = compatMsg_ ? PtrToPtr<event_sync_msg, old_event_sync_msg>(drvSyncMsg_)->pid : drvSyncMsg_->pid;
     585           87 :         qsEvent.grp_id = compatMsg_ ? PtrToPtr<event_sync_msg, old_event_sync_msg>(drvSyncMsg_)->gid : drvSyncMsg_->gid;
     586              :         const int32_t eventId =
     587           87 :             compatMsg_ ? PtrToPtr<event_sync_msg, old_event_sync_msg>(drvSyncMsg_)->event_id : drvSyncMsg_->event_id;
     588           87 :         qsEvent.event_id = static_cast<EVENT_ID>(eventId);
     589           87 :         qsEvent.subevent_id = compatMsg_ ? PtrToPtr<event_sync_msg, old_event_sync_msg>(drvSyncMsg_)->subevent_id :
     590            0 :                                            drvSyncMsg_->subevent_id;
     591           87 :         drvRet = halEschedSubmitEvent(deviceId_, &qsEvent); // drv interface require use 0
     592           87 :         drvSyncMsg_ = nullptr;
     593           87 :         BQS_LOG_INFO(
     594              :             "[SendRspEvent] dst_engine[%u], pid[%d], grp_id[%u], eventId[%d], subevent_id[%u], deviceId[%u]",
     595              :             qsEvent.dst_engine, qsEvent.pid, qsEvent.grp_id, qsEvent.event_id, qsEvent.subevent_id, deviceId_);
     596              :     } else {
     597           25 :         BQS_LOG_INFO("[RouterServer] Do AICPU response");
     598           25 :         qsEvent.pid = srcPid_;
     599           25 :         qsEvent.grp_id = static_cast<uint32_t>(srcGroupId_);
     600           25 :         qsEvent.event_id = EVENT_QS_MSG;
     601           25 :         qsEvent.msg_len = static_cast<uint32_t>(sizeof(QsProcMsgRspDstAicpu));
     602           25 :         const auto iter = g_reqRspMapping.find(static_cast<int32_t>(subEventId_));
     603           25 :         if (iter != g_reqRspMapping.end()) {
     604           15 :             qsEvent.subevent_id = static_cast<uint32_t>(iter->second);
     605              :         } else {
     606           10 :             BQS_LOG_ERROR("[RouterServer]ERROR Invalid subeventId[%u]", subEventId_);
     607              :         }
     608           25 :         drvRet = halEschedSubmitEvent(deviceId_, &qsEvent); // drv interface require use 0
     609           25 :         aicpuRspHead_ = 0UL;
     610              :     }
     611          112 :     BQS_LOG_ERROR_WHEN(
     612              :         drvRet != DRV_ERROR_NONE, "[RouterServer]ERROR failed to submit event[%u], result[%d].", subEventId_,
     613              :         static_cast<int32_t>(drvRet));
     614          112 :     BQS_LOG_INFO("[RouterServer]Finish response message subeventid[%u] ret[%d]", subEventId_, result);
     615          112 :     qsRouteListPtr_ = nullptr;
     616          112 :     qsRouterHeadPtr_ = nullptr;
     617          112 :     subEventId_ = 0U;
     618          224 :     return;
     619              : }
     620              : 
     621            2 : void RouterServer::ProcessBindQueue(const uint32_t index)
     622              : {
     623            2 :     BQS_LOG_INFO("[RouterServer]Bind relation [add], stage [server:process].");
     624            2 :     auto& relationInstance = BindRelation::GetInstance();
     625            2 :     retCode_ = static_cast<int32_t>(BQS_STATUS_OK);
     626            2 :     QueueRoute* queueRouteList = qsRouteListPtr_;
     627            8 :     for (uint32_t i = 0U; i < qsRouterHeadPtr_->routeNum; ++i) {
     628            6 :         if (queueRouteList->status != static_cast<int32_t>(BQS_STATUS_OK)) {
     629            3 :             retCode_ = static_cast<int32_t>(BQS_STATUS_QUEUE_AHTU_ERROR);
     630            3 :             queueRouteList->status = 0;
     631            3 :             queueRouteList = queueRouteList + 1U;
     632            3 :             continue;
     633              :         }
     634              :         // only queue
     635              :         EntityInfo srcEntity =
     636            3 :             CreateBasicEntityInfo(queueRouteList->srcId, static_cast<dgw::EntityType>(queueRouteList->srcType));
     637              :         EntityInfo dstEntity =
     638            3 :             CreateBasicEntityInfo(queueRouteList->dstId, static_cast<dgw::EntityType>(queueRouteList->dstType));
     639            3 :         const auto result = relationInstance.Bind(srcEntity, dstEntity, index);
     640            3 :         if (result == BQS_STATUS_RETRY) {
     641            0 :             queueRouteList = queueRouteList + 1U;
     642            0 :             continue;
     643            3 :         } else if (result != BQS_STATUS_OK) {
     644            0 :             retCode_ = static_cast<int32_t>(result);
     645            0 :             queueRouteList->status = 0;
     646              :         } else {
     647            3 :             queueRouteList->status = 1;
     648              :         }
     649            3 :         queueRouteList = queueRouteList + 1U;
     650            3 :         BQS_LOG_RUN_INFO(
     651              :             "Bind relation [add], stage [server:process], relation [src:%s, dst:%s, result:%d]",
     652              :             srcEntity.ToString().c_str(), dstEntity.ToString().c_str(), static_cast<int32_t>(result));
     653            3 :     }
     654            2 :     relationInstance.Order(index);
     655            2 :     return;
     656              : }
     657              : 
     658            2 : void RouterServer::ProcessUnbindQueue(const uint32_t index)
     659              : {
     660            2 :     BQS_LOG_INFO("[RouterServer]Unbind relation [del], stage [server:process].");
     661            2 :     QueueRoute* queueRouteList = qsRouteListPtr_;
     662            2 :     retCode_ = static_cast<int32_t>(BQS_STATUS_OK);
     663            2 :     auto& relationInstance = BindRelation::GetInstance();
     664            8 :     for (uint32_t i = 0U; i < qsRouterHeadPtr_->routeNum; ++i) {
     665            6 :         if (queueRouteList->status != static_cast<int32_t>(BQS_STATUS_OK)) {
     666            3 :             retCode_ = static_cast<int32_t>(BQS_STATUS_QUEUE_ID_ERROR);
     667            3 :             queueRouteList->status = 1;
     668            3 :             queueRouteList = queueRouteList + 1U;
     669            3 :             continue;
     670              :         }
     671              :         EntityInfo srcEntity =
     672            3 :             CreateBasicEntityInfo(queueRouteList->srcId, static_cast<dgw::EntityType>(queueRouteList->srcType));
     673              :         EntityInfo dstEntity =
     674            3 :             CreateBasicEntityInfo(queueRouteList->dstId, static_cast<dgw::EntityType>(queueRouteList->dstType));
     675            3 :         const auto result = relationInstance.UnBind(srcEntity, dstEntity, index);
     676            3 :         if (result == BQS_STATUS_RETRY) {
     677            0 :             queueRouteList = queueRouteList + 1U;
     678            0 :             continue;
     679            3 :         } else if (result != BQS_STATUS_OK) {
     680            0 :             retCode_ = static_cast<int32_t>(result);
     681            0 :             queueRouteList->status = 1;
     682              :         } else {
     683            3 :             queueRouteList->status = 0;
     684              :         }
     685            3 :         queueRouteList = queueRouteList + 1U;
     686            3 :         BQS_LOG_RUN_INFO(
     687              :             "Bind relation [del], stage [server:process], relation [src %s,"
     688              :             "dst %s, result:%d]",
     689              :             srcEntity.ToString().c_str(), dstEntity.ToString().c_str(), static_cast<int32_t>(result));
     690            3 :     }
     691            2 :     relationInstance.Order(index);
     692            2 :     return;
     693              : }
     694              : 
     695              : /**
     696              :  * Bqs server enqueue bind msg request process
     697              :  * @return NA
     698              :  */
     699           41 : void RouterServer::BindMsgProc(const uint32_t index)
     700              : {
     701           41 :     BQS_LOG_INFO("[RouterServer]RouterServer BindMsgProc begin.");
     702           41 :     auto& processing = (index == 0U) ? processing_ : processingExtra_;
     703           41 :     auto& done = (index == 0U) ? done_ : doneExtra_;
     704              : 
     705           41 :     const std::unique_lock<std::mutex> lk(mutex_);
     706           41 :     processing = true;
     707              : 
     708              :     // parse bind and unbind BQSMsg
     709           41 :     if ((subEventId_ == static_cast<uint32_t>(AICPU_BIND_QUEUE)) ||
     710           41 :         (subEventId_ == static_cast<uint32_t>(ACL_BIND_QUEUE))) {
     711            2 :         ProcessBindQueue(index);
     712           39 :     } else if (
     713           39 :         (subEventId_ == static_cast<uint32_t>(AICPU_UNBIND_QUEUE)) ||
     714           37 :         (subEventId_ == static_cast<uint32_t>(ACL_UNBIND_QUEUE))) {
     715            2 :         ProcessUnbindQueue(index);
     716           37 :     } else if ((subEventId_ == static_cast<uint32_t>(QueueSubEventType::UPDATE_CONFIG))) {
     717           37 :         retCode_ = (cfgInfoOperator_ == nullptr) ? static_cast<int32_t>(BQS_STATUS_INNER_ERROR) :
     718           37 :                                                    static_cast<int32_t>(cfgInfoOperator_->ProcessUpdateConfig(index));
     719           37 :         BQS_LOG_INFO("[RouterServer] Process update config ret is %d.", retCode_);
     720              :     } else {
     721            0 :         BQS_LOG_ERROR("[RouterServer]Invalid subEventId_[%d] in bind relation process.", subEventId_);
     722              :     }
     723              : 
     724           41 :     processing = false;
     725           41 :     done = true;
     726           41 :     cv_.notify_one();
     727              : 
     728           41 :     BQS_LOG_INFO("[RouterServer]RouterServer BindMsgProc end.");
     729           82 :     return;
     730           41 : }
     731              : 
     732              : /**
     733              :  * Init bqs server, including init easycomm server and bind relation
     734              :  * @return BQS_STATUS_OK:success other:failed
     735              :  */
     736            5 : BqsStatus RouterServer::InitRouterServer(const InitQsParams& params)
     737              : {
     738            5 :     BQS_LOG_INFO("[RouterServer]RouterServer Init begin");
     739            5 :     (void)signal(SIGPIPE, SIG_IGN);
     740              : 
     741            5 :     qsInitGroupName_ = params.qsInitGrpName;
     742            5 :     f2nfGroupId_ = params.f2nfGroupId;
     743            5 :     schedPolicy_ = params.schedPolicy;
     744              :     // create config info operator
     745            5 :     cfgInfoOperator_.reset(new (std::nothrow) ConfigInfoOperator(params.deviceId, qsInitGroupName_));
     746            5 :     if (cfgInfoOperator_ == nullptr) {
     747            0 :         BQS_LOG_ERROR("malloc memory for cfgInfoOperator_ failed.");
     748            0 :         return BQS_STATUS_INNER_ERROR;
     749              :     }
     750            5 :     SubscribeBufEvent();
     751              : 
     752            5 :     if (!running_) {
     753            5 :         running_ = true;
     754            5 :         deviceId_ = params.deviceId;
     755            5 :         deployMode_ = params.runMode;
     756            5 :         numaFlag_ = params.numaFlag;
     757            5 :         needAttachGroup_ = params.needAttachGroup;
     758              :         // halShrIdGetAttribute为新版驱动中才存在的接口,若没有,说明驱动为老版本,需兼容处理
     759            5 :         if ((bqs::GetRunContext() == bqs::RunContext::HOST) && (&halShrIdGetAttribute == nullptr)) {
     760            5 :             compatMsg_ = true;
     761              :         }
     762            5 :         BQS_LOG_INFO("compatMsg_ is %d", compatMsg_);
     763              : 
     764              :         try {
     765            5 :             monitorQsEvent_ = std::thread(&RouterServer::ManageQsEvent, this);
     766            0 :         } catch (std::exception& threadException) {
     767            0 :             BQS_LOG_ERROR("RouterServer Init thread failure, %s", threadException.what());
     768            0 :             return BQS_STATUS_INNER_ERROR;
     769            0 :         }
     770              : 
     771            5 :         std::unique_lock<std::mutex> lk(manageThreadMutex_);
     772           12 :         manageThreadCv_.wait(lk, [this] { return manageThreadStatus_ != ThreadStatus::NOT_INIT; });
     773            5 :         if (manageThreadStatus_ != ThreadStatus::INIT_SUCCESS) {
     774            2 :             BQS_LOG_ERROR("RouterServer thread fail to start");
     775            2 :             return BQS_STATUS_INNER_ERROR;
     776              :         }
     777              : 
     778            3 :         BQS_LOG_INFO("[RouterServer]RouterServer Init success.");
     779            5 :     } else {
     780            0 :         BQS_LOG_WARN("RouterServer is already inited");
     781              :     }
     782            3 :     return BQS_STATUS_OK;
     783              : }
     784              : 
     785            8 : BqsStatus RouterServer::ParseRelationInfo(Mbuf** mbufPtr)
     786              : {
     787            8 :     Mbuf* mBuf = nullptr;
     788            8 :     const auto drvRet = halQueueDeQueue(deviceId_, pipelineQueueId_, PtrToPtr<Mbuf*, void*>(&mBuf));
     789            8 :     if ((drvRet != DRV_ERROR_NONE) || (mBuf == nullptr)) {
     790            2 :         BQS_LOG_ERROR(
     791              :             "[RouterServer]halQueueDeQueue from queue[%u] in device[%u] failed, error[%d]", pipelineQueueId_.load(),
     792              :             deviceId_, static_cast<int32_t>(drvRet));
     793            1 :         return BQS_STATUS_DRIVER_ERROR;
     794              :     }
     795            7 :     *mbufPtr = mBuf;
     796            7 :     qsRouterHeadPtr_ = nullptr;
     797            7 :     const auto getBuffRet = halMbufGetBuffAddr(mBuf, reinterpret_cast<void**>(&qsRouterHeadPtr_));
     798            7 :     if ((getBuffRet != static_cast<int32_t>(DRV_ERROR_NONE)) || (qsRouterHeadPtr_ == nullptr)) {
     799            0 :         BQS_LOG_ERROR(
     800              :             "[RouterServer]halMbufGetBuffAddr from queue[%u] in device[%u] failed, error[%d]", pipelineQueueId_.load(),
     801              :             deviceId_, getBuffRet);
     802            0 :         return BQS_STATUS_DRIVER_ERROR;
     803              :     }
     804            7 :     if (isAicpuEvent_) {
     805            5 :         aicpuRspHead_ = qsRouterHeadPtr_->userData;
     806              :     }
     807              :     // aicpuRspHead_ is valid only in aicpu event senario, will be 0 in acl event senario
     808            7 :     subEventId_ = qsRouterHeadPtr_->subEventId;
     809            7 :     BQS_LOG_INFO("[RouterServer]Parse head[%lu] subEvnetId[%u] from mbuff success.", aicpuRspHead_, subEventId_);
     810              : 
     811              :     // query message need to get query info
     812            7 :     if ((subEventId_ == static_cast<uint32_t>(AICPU_QUERY_QUEUE)) ||
     813            3 :         (subEventId_ == static_cast<uint32_t>(ACL_QUERY_QUEUE))) {
     814            4 :         if ((((qsRouterHeadPtr_->routeNum * sizeof(QueueRoute)) + sizeof(QsRouteHead)) + sizeof(QueueRouteQuery)) !=
     815            4 :             qsRouterHeadPtr_->length) {
     816            0 :             BQS_LOG_ERROR(
     817              :                 "[RouterServer]RouteNum[%d] is inconsistence with dataLen[%d] in subEventId[%u]",
     818              :                 qsRouterHeadPtr_->routeNum, qsRouterHeadPtr_->length, subEventId_);
     819            0 :             return BQS_STATUS_PARAM_INVALID;
     820              :         }
     821            4 :         qsRouterQueryPtr_ =
     822            4 :             reinterpret_cast<QueueRouteQuery*>(reinterpret_cast<uint8_t*>(qsRouterHeadPtr_) + sizeof(QsRouteHead));
     823            4 :         BQS_LOG_INFO("[RouterServer]Get query info success. queryType[%d]", qsRouterQueryPtr_->queryType);
     824            4 :         qsRouteListPtr_ =
     825            4 :             reinterpret_cast<QueueRoute*>(reinterpret_cast<uint8_t*>(qsRouterQueryPtr_) + sizeof(QueueRouteQuery));
     826              :     } else {
     827            3 :         if (((qsRouterHeadPtr_->routeNum * sizeof(QueueRoute)) + sizeof(QsRouteHead)) != qsRouterHeadPtr_->length) {
     828            0 :             BQS_LOG_ERROR(
     829              :                 "[RouterServer]RouteNum[%d] is inconsistence with dataLen[%d] in subEventId[%u]",
     830              :                 qsRouterHeadPtr_->routeNum, qsRouterHeadPtr_->length, subEventId_);
     831            0 :             return BQS_STATUS_PARAM_INVALID;
     832              :         }
     833            3 :         qsRouterQueryPtr_ = nullptr;
     834            3 :         qsRouteListPtr_ =
     835            3 :             reinterpret_cast<QueueRoute*>(reinterpret_cast<uint8_t*>(qsRouterHeadPtr_) + sizeof(QsRouteHead));
     836              :     }
     837            7 :     BQS_LOG_INFO("[RouterServer]Get relation mbuff success, bind/unbind queue num[%d]", qsRouterHeadPtr_->routeNum);
     838            7 :     return BQS_STATUS_OK;
     839              : }
     840              : 
     841            3 : BqsStatus RouterServer::ParseBindUnbindMsg() const
     842              : {
     843            3 :     BQS_LOG_INFO("[RouterServer]Bind relation [add/del], stage [server:parse and check]");
     844            3 :     auto resultCode = BQS_STATUS_QUEUE_AHTU_ERROR;
     845            3 :     QueueRoute* queueRouteList = qsRouteListPtr_;
     846           12 :     for (uint32_t i = 0U; i < qsRouterHeadPtr_->routeNum; ++i) {
     847              :         EntityInfo srcEntity =
     848            9 :             CreateBasicEntityInfo(queueRouteList->srcId, static_cast<dgw::EntityType>(queueRouteList->srcType));
     849              :         EntityInfo dstEntity =
     850            9 :             CreateBasicEntityInfo(queueRouteList->dstId, static_cast<dgw::EntityType>(queueRouteList->dstType));
     851            9 :         BQS_LOG_INFO(
     852              :             "[RouterServer]Src[id:%u type:%d] Dst[id:%u type:%d]", srcEntity.GetId(),
     853              :             static_cast<int32_t>(srcEntity.GetType()), dstEntity.GetId(), static_cast<int32_t>(dstEntity.GetType()));
     854            9 :         if ((srcEntity.GetId() >= MAX_QUEUE_ID_NUM) || (dstEntity.GetId() >= MAX_QUEUE_ID_NUM)) {
     855            5 :             BQS_LOG_ERROR(
     856              :                 "[RouterServer]Src[%s] or Dst[%s] is invalid in this "
     857              :                 "bind/unbind relation",
     858              :                 srcEntity.ToString().c_str(), dstEntity.ToString().c_str());
     859            5 :             queueRouteList->status = static_cast<int32_t>(BQS_STATUS_QUEUE_ID_ERROR);
     860            5 :             queueRouteList = queueRouteList + 1;
     861            5 :             continue;
     862              :         }
     863              :         // preprocess: do attach queue and check src own read auth, dst own write auth
     864            4 :         if ((subEventId_ == static_cast<uint32_t>(AICPU_BIND_QUEUE)) ||
     865            4 :             (subEventId_ == static_cast<uint32_t>(ACL_BIND_QUEUE))) {
     866            4 :             const auto ret = (cfgInfoOperator_ == nullptr) ?
     867              :                                  BQS_STATUS_INNER_ERROR :
     868            0 :                                  cfgInfoOperator_->AttachAndCheckQueue(srcEntity, dstEntity);
     869            4 :             if (ret != BQS_STATUS_OK) {
     870            4 :                 BQS_LOG_ERROR(
     871              :                     "[RouterServer]Src[%s] Dst[%s] do attach queue "
     872              :                     "and check auth failed",
     873              :                     srcEntity.ToString().c_str(), dstEntity.ToString().c_str());
     874            4 :                 queueRouteList->status = static_cast<int32_t>(BQS_STATUS_QUEUE_AHTU_ERROR);
     875            4 :                 queueRouteList = queueRouteList + 1;
     876            4 :                 continue;
     877              :             }
     878              :         }
     879            0 :         resultCode = BQS_STATUS_OK;
     880            0 :         queueRouteList->status = static_cast<int32_t>(BQS_STATUS_OK);
     881            0 :         queueRouteList = queueRouteList + 1;
     882           18 :     }
     883            3 :     BQS_LOG_INFO("[RouterServer]Finish parse bind/unbind message,resultCode[%d]", static_cast<int32_t>(resultCode));
     884            3 :     return resultCode;
     885              : }
     886              : 
     887           14 : void RouterServer::FillRoutes(const EntityInfo& src, const EntityInfo& dst, const BindRelationStatus status)
     888              : {
     889           14 :     QueueRoute queueRouteInfo = {};
     890           14 :     queueRouteInfo.srcId = src.GetId();
     891           14 :     queueRouteInfo.dstId = dst.GetId();
     892           14 :     queueRouteInfo.srcType = static_cast<int16_t>(src.GetType());
     893           14 :     queueRouteInfo.dstType = static_cast<int16_t>(dst.GetType());
     894           14 :     queueRouteInfo.status = static_cast<int32_t>(status);
     895           14 :     queueRouteQueryList_.emplace_back(queueRouteInfo);
     896           14 : }
     897              : 
     898           28 : void RouterServer::SearchRelation(
     899              :     const MapEnitityInfoToInfoSet& relationMap, const EntityInfo& entityInfo, const BindRelationStatus status,
     900              :     bool bySrc)
     901              : {
     902           28 :     const auto iter = relationMap.find(entityInfo);
     903           28 :     if (iter == relationMap.end()) {
     904           21 :         return;
     905              :     }
     906           12 :     const auto dstSet = iter->second;
     907           12 :     BQS_LOG_INFO("[RouterServer]Bind relation [get], stage [server:process], relation [size:%zu].", dstSet.size());
     908           12 :     if (bySrc) {
     909           12 :         for (auto setIter = dstSet.begin(); setIter != dstSet.end(); ++setIter) {
     910            7 :             FillRoutes(entityInfo, *setIter, status);
     911              :         }
     912            5 :         return;
     913              :     }
     914              : 
     915           14 :     for (auto setIter = dstSet.begin(); setIter != dstSet.end(); ++setIter) {
     916            7 :         FillRoutes(*setIter, entityInfo, status);
     917              :     }
     918           12 : }
     919              : 
     920              : /**
     921              :  * Assembly response of get bind message according to src entity
     922              :  * @return Number of query results
     923              :  */
     924           13 : void RouterServer::GetBindRspBySingle(const EntityInfo& entityInfo, const uint32_t& queryType)
     925              : {
     926           13 :     BQS_LOG_INFO(
     927              :         "[RouterServer]RouterServer serialize get bind rsponse by entityId[%u], entityType[%d], Type[%d].",
     928              :         entityInfo.GetId(), static_cast<int32_t>(entityInfo.GetType()), queryType);
     929           13 :     auto& relationInstance = BindRelation::GetInstance();
     930           13 :     queueRouteQueryList_.clear();
     931           13 :     if ((queryType != static_cast<uint32_t>(BQS_QUERY_TYPE_SRC)) &&
     932            7 :         (queryType != static_cast<uint32_t>(BQS_QUERY_TYPE_DST))) {
     933            1 :         BQS_LOG_ERROR("[RouterServer]QueryType[%d] is not supported.", queryType);
     934            1 :         return;
     935              :     }
     936              : 
     937           12 :     if (queryType == static_cast<uint32_t>(BQS_QUERY_TYPE_SRC)) {
     938            6 :         SearchRelation(relationInstance.GetSrcToDstRelation(), entityInfo, BindRelationStatus::RelationBind, true);
     939            6 :         if (numaFlag_) {
     940            2 :             SearchRelation(
     941              :                 relationInstance.GetSrcToDstExtraRelation(), entityInfo, BindRelationStatus::RelationBind, true);
     942              :         }
     943            6 :         SearchRelation(
     944              :             relationInstance.GetAbnormalSrcToDstRelation(), entityInfo, BindRelationStatus::RelationAbnormalForQError,
     945              :             true);
     946              : 
     947            6 :         if (queueRouteQueryList_.empty()) {
     948            2 :             BQS_LOG_WARN("[RouterServer] record does not exist according to src entityId:[%u]", entityInfo.GetId());
     949              :         }
     950              :     } else {
     951            6 :         SearchRelation(relationInstance.GetDstToSrcRelation(), entityInfo, BindRelationStatus::RelationBind, false);
     952            6 :         if (numaFlag_) {
     953            2 :             SearchRelation(
     954              :                 relationInstance.GetDstToSrcExtraRelation(), entityInfo, BindRelationStatus::RelationBind, false);
     955              :         }
     956            6 :         SearchRelation(
     957              :             relationInstance.GetAbnormalDstToSrcRelation(), entityInfo, BindRelationStatus::RelationAbnormalForQError,
     958              :             false);
     959              :     }
     960           12 :     if (queueRouteQueryList_.empty()) {
     961            2 :         BQS_LOG_WARN("RouterServer get relation according to dst:%u failed, record does not exist", entityInfo.GetId());
     962              :     }
     963              : }
     964              : 
     965            5 : bool RouterServer::FindRelation(
     966              :     const MapEnitityInfoToInfoSet& relationMap, const EntityInfo& srcInfo, const EntityInfo& dstInfo) const
     967              : {
     968            5 :     const auto srcIter = relationMap.find(srcInfo);
     969            5 :     return ((srcIter != relationMap.end()) && (srcIter->second.count(dstInfo) != 0UL));
     970              : }
     971              : 
     972            4 : void RouterServer::TransRouteWithEntityInfo(
     973              :     const EntityInfo& srcInfo, const EntityInfo& dstInfo, const int32_t status, QueueRoute& routeInfo) const
     974              : {
     975            4 :     routeInfo.srcId = srcInfo.GetId();
     976            4 :     routeInfo.dstId = dstInfo.GetId();
     977            4 :     routeInfo.srcType = static_cast<int16_t>(srcInfo.GetType());
     978            4 :     routeInfo.dstType = static_cast<int16_t>(dstInfo.GetType());
     979            4 :     routeInfo.status = static_cast<int32_t>(status);
     980            4 : }
     981              : 
     982              : /**
     983              :  * Assembly response of get bind message according to dst queueId, one-to-one relation
     984              :  * @return NA
     985              :  */
     986            5 : void RouterServer::GetBindRspByDouble(const EntityInfo& src, const EntityInfo& dst, const uint32_t& queryType)
     987              : {
     988            5 :     BQS_LOG_INFO(
     989              :         "[RouterServer]RouterServer serialize get bind rsponse by srcId[%u], srcType[%d], dstId[%u], "
     990              :         "dstType[%d], Type[%u]",
     991              :         src.GetId(), static_cast<int32_t>(src.GetType()), dst.GetId(), static_cast<int32_t>(dst.GetType()), queryType);
     992            5 :     queueRouteQueryList_.clear();
     993            5 :     if (queryType == static_cast<uint32_t>(BQS_QUERY_TYPE_SRC_OR_DST)) {
     994            2 :         GetBindRspBySingle(src, static_cast<uint32_t>(BQS_QUERY_TYPE_SRC));
     995            2 :         if (queueRouteQueryList_.size() > 0UL) {
     996            0 :             return;
     997              :         }
     998            2 :         GetBindRspBySingle(dst, static_cast<uint32_t>(BQS_QUERY_TYPE_DST));
     999            2 :         return;
    1000              :     }
    1001            3 :     if (queryType == static_cast<uint32_t>(BQS_QUERY_TYPE_SRC_AND_DST)) {
    1002            3 :         auto status = BindRelationStatus::RelationUnknown;
    1003            4 :         if (FindRelation(BindRelation::GetInstance().GetSrcToDstRelation(), src, dst) ||
    1004            1 :             (numaFlag_ && (FindRelation(BindRelation::GetInstance().GetSrcToDstExtraRelation(), src, dst)))) {
    1005            2 :             status = BindRelationStatus::RelationBind;
    1006              :         } else {
    1007            1 :             if (FindRelation(BindRelation::GetInstance().GetAbnormalSrcToDstRelation(), src, dst)) {
    1008            1 :                 status = BindRelationStatus::RelationAbnormalForQError;
    1009              :             }
    1010              :         }
    1011              : 
    1012            3 :         if (status != BindRelationStatus::RelationUnknown) {
    1013            3 :             QueueRoute queueRouteInfo = {};
    1014            3 :             TransRouteWithEntityInfo(src, dst, static_cast<int32_t>(status), queueRouteInfo);
    1015            3 :             queueRouteQueryList_.emplace_back(queueRouteInfo);
    1016              :         }
    1017            3 :         return;
    1018              :     }
    1019            0 :     BQS_LOG_ERROR("[RouterServer]QueryType[%d] is not supported.", queryType);
    1020            0 :     return;
    1021              : }
    1022              : 
    1023            1 : void RouterServer::GetAllAbnormalBind()
    1024              : {
    1025            1 :     BQS_LOG_INFO("[RouterServer]RouterServer serialize get all abnormal bind rsponse");
    1026            1 :     queueRouteQueryList_.clear();
    1027            1 :     auto& relationInstance = BindRelation::GetInstance();
    1028            1 :     auto& abnormalSrcToDstRelation = relationInstance.GetAbnormalSrcToDstRelation();
    1029            2 :     for (auto iter = abnormalSrcToDstRelation.begin(); iter != abnormalSrcToDstRelation.end(); ++iter) {
    1030            1 :         const auto& src = iter->first;
    1031            1 :         const auto& dstSet = iter->second;
    1032            2 :         for (auto& dst : dstSet) {
    1033            1 :             QueueRoute queueRouteInfo = {};
    1034            1 :             TransRouteWithEntityInfo(
    1035              :                 src, dst, static_cast<int32_t>(BindRelationStatus::RelationAbnormalForQError), queueRouteInfo);
    1036            1 :             queueRouteQueryList_.emplace_back(queueRouteInfo);
    1037              :         }
    1038              :     }
    1039            1 : }
    1040              : 
    1041            8 : BqsStatus RouterServer::ProcessGetBindMsg(const uint32_t& queryType, const EntityInfo& src, const EntityInfo& dst)
    1042              : {
    1043            8 :     switch (queryType) {
    1044            2 :         case BQS_QUERY_TYPE_SRC:
    1045            2 :             GetBindRspBySingle(src, queryType);
    1046            2 :             break;
    1047            2 :         case BQS_QUERY_TYPE_DST:
    1048            2 :             GetBindRspBySingle(dst, queryType);
    1049            2 :             break;
    1050            4 :         case BQS_QUERY_TYPE_SRC_OR_DST:
    1051              :         case BQS_QUERY_TYPE_SRC_AND_DST:
    1052            4 :             GetBindRspByDouble(src, dst, queryType);
    1053            4 :             break;
    1054            0 :         case BQS_QUERY_TYPE_ABNORMAL_FOR_QUEUE_ERROR:
    1055            0 :             GetAllAbnormalBind();
    1056            0 :             break;
    1057            0 :         default:
    1058            0 :             BQS_LOG_ERROR(
    1059              :                 "[RouterServer]Unsupported query type"
    1060              :                 "{0:src, 1:dst, 2:src-or-dst, 3:src-and-dst, 100:abnormal-all}:%u",
    1061              :                 queryType);
    1062            0 :             break;
    1063              :     }
    1064              : 
    1065            8 :     if ((subEventId_ == static_cast<uint32_t>(AICPU_QUERY_QUEUE_NUM)) ||
    1066            4 :         (subEventId_ == static_cast<uint32_t>(ACL_QUERY_QUEUE_NUM))) {
    1067            4 :         return BQS_STATUS_OK;
    1068              :     }
    1069              : 
    1070            4 :     if (queueRouteQueryList_.size() != qsRouterHeadPtr_->routeNum) {
    1071            0 :         BQS_LOG_ERROR(
    1072              :             "[RouterServer]Prepare number[%d] is different with real route number[%zu].", qsRouterHeadPtr_->routeNum,
    1073              :             queueRouteQueryList_.size());
    1074            0 :         return BQS_STATUS_PARAM_INVALID;
    1075              :     }
    1076            9 :     for (size_t i = 0UL; i < qsRouterHeadPtr_->routeNum; i++) {
    1077            5 :         *qsRouteListPtr_ = queueRouteQueryList_[i];
    1078            5 :         BQS_LOG_INFO(
    1079              :             "[RouterServer]Query bind relation srcId[%u] srcType[%d] dstId[%u] dstType[%d]", qsRouteListPtr_->srcId,
    1080              :             static_cast<int32_t>(qsRouteListPtr_->srcType), qsRouteListPtr_->dstId,
    1081              :             static_cast<int32_t>(qsRouteListPtr_->dstType));
    1082            5 :         qsRouteListPtr_ = qsRouteListPtr_ + 1;
    1083              :     }
    1084            4 :     queueRouteQueryList_.clear();
    1085            4 :     return BQS_STATUS_OK;
    1086              : }
    1087              : 
    1088            4 : void RouterServer::ParseGetBindNumMsg(const event_info& info)
    1089              : {
    1090            4 :     BQS_LOG_INFO("[RouterServer]Bind relation number [get], stage [server:process], type [request]");
    1091            4 :     if (info.priv.msg_len != sizeof(QueueRouteQuery)) {
    1092            0 :         BQS_LOG_ERROR("[RouterServer]Query event[%u] message invalid, msgLen[%u]", subEventId_, info.priv.msg_len);
    1093            0 :         SendRspEvent(static_cast<int32_t>(BQS_STATUS_PARAM_INVALID));
    1094            0 :         return;
    1095              :     }
    1096            4 :     const QueueRouteQuery* const queueRouteQuery = PtrToPtr<const char_t, const QueueRouteQuery>(info.priv.msg);
    1097              :     // param syncEventHead is filled by drv when acl (call sync event interface)
    1098            4 :     aicpuRspHead_ = isAicpuEvent_ ? queueRouteQuery->syncEventHead : 0UL;
    1099            4 :     const uint32_t keyType = queueRouteQuery->queryType;
    1100              :     const EntityInfo src =
    1101            4 :         CreateBasicEntityInfo(queueRouteQuery->srcId, static_cast<dgw::EntityType>(queueRouteQuery->srcType));
    1102              :     const EntityInfo dst =
    1103            4 :         CreateBasicEntityInfo(queueRouteQuery->dstId, static_cast<dgw::EntityType>(queueRouteQuery->dstType));
    1104            4 :     queueRouteQueryList_.clear();
    1105            4 :     const auto ret = ProcessGetBindMsg(keyType, src, dst);
    1106            4 :     SendRspEvent(static_cast<int32_t>(ret));
    1107            4 :     return;
    1108            4 : }
    1109              : 
    1110            5 : BqsStatus RouterServer::ParseGetBindDetailMsg()
    1111              : {
    1112            5 :     BQS_LOG_INFO("[RouterServer]Bind relation [get], stage [server:process], type [request]");
    1113            5 :     if (qsRouterQueryPtr_ == nullptr) {
    1114            1 :         BQS_LOG_ERROR("[RouterServer]qsRouterQuery should not be null pointer");
    1115            1 :         return BQS_STATUS_INNER_ERROR;
    1116              :     }
    1117            4 :     const uint32_t keyType = qsRouterQueryPtr_->queryType;
    1118              :     const EntityInfo src =
    1119            4 :         CreateBasicEntityInfo(qsRouterQueryPtr_->srcId, static_cast<dgw::EntityType>(qsRouterQueryPtr_->srcType));
    1120              :     const EntityInfo dst =
    1121            4 :         CreateBasicEntityInfo(qsRouterQueryPtr_->dstId, static_cast<dgw::EntityType>(qsRouterQueryPtr_->dstType));
    1122            4 :     queueRouteQueryList_.clear();
    1123            4 :     const auto ret = ProcessGetBindMsg(keyType, src, dst);
    1124            4 :     return ret;
    1125            4 : }
    1126              : 
    1127           11 : ThreadStatus RouterServer::PrePareForManageThread()
    1128              : {
    1129           11 :     auto ret = halEschedAttachDevice(deviceId_);
    1130           11 :     if ((ret != DRV_ERROR_NONE) && (ret != DRV_ERROR_PROCESS_REPEAT_ADD)) {
    1131            1 :         BQS_LOG_ERROR("Failed to attach device[%u] for eSched, result[%d].", deviceId_, static_cast<int32_t>(ret));
    1132            1 :         (void)AttachGroup();
    1133            1 :         return ThreadStatus::INIT_FAIL;
    1134              :     }
    1135           10 :     ret = halEschedCreateGrp(deviceId_, bindQueueGroupId_, GRP_TYPE_BIND_CP_CPU);
    1136           10 :     if (ret != DRV_ERROR_NONE) {
    1137            1 :         (void)halEschedDettachDevice(deviceId_);
    1138            1 :         BQS_LOG_ERROR(
    1139              :             "Failed to create bindQueueGroup, groupId[%u] result[%d].", bindQueueGroupId_, static_cast<int32_t>(ret));
    1140            1 :         (void)AttachGroup();
    1141            1 :         return ThreadStatus::INIT_FAIL;
    1142              :     }
    1143              :     // subsribe bind/unbind/query event
    1144            9 :     const uint64_t eventBitmap = (1UL << static_cast<uint32_t>(EVENT_QS_MSG));
    1145            9 :     BQS_LOG_INFO("[RouterServer]BindQueue group[%u] subscribe event, eventBitmap[%lu]", bindQueueGroupId_, eventBitmap);
    1146            9 :     ret = halEschedSubscribeEvent(deviceId_, bindQueueGroupId_, 0U, eventBitmap);
    1147            9 :     if (ret != DRV_ERROR_NONE) {
    1148            1 :         BQS_LOG_ERROR(
    1149              :             "[RouterServer]halEschedSubscribeEvent failed, groupId[%u] eventBitmap[%lu] result[%d].", bindQueueGroupId_,
    1150              :             eventBitmap, static_cast<int32_t>(ret));
    1151            1 :         (void)AttachGroup();
    1152            1 :         return ThreadStatus::INIT_FAIL;
    1153              :     }
    1154              : 
    1155            8 :     if (!AttachGroup()) {
    1156            2 :         BQS_LOG_ERROR("[RouterServer] Fail to attach group");
    1157            2 :         return ThreadStatus::INIT_FAIL;
    1158              :     }
    1159            6 :     BQS_LOG_RUN_INFO("[RouterServer] ManageQsEvent of RouterServer is ready.");
    1160            6 :     return ThreadStatus::INIT_SUCCESS;
    1161              : }
    1162              : 
    1163           11 : void RouterServer::ManageQsEvent()
    1164              : {
    1165           11 :     BQS_LOG_INFO("[RouterServer] Manage QS event of router server thread start");
    1166           11 :     (void)pthread_setname_np(pthread_self(), ROUTER_SERVER_THREAD_NAME_PREFIX);
    1167              : 
    1168           11 :     if (bqs::GetRunContext() != bqs::RunContext::HOST) {
    1169            1 :         const std::vector<uint32_t>& cpuIds = QueueScheduleInterface::GetInstance().GetCtrlCpuIds();
    1170              :         // bind thread to ctrl cpu
    1171            1 :         const pthread_t threadId = pthread_self();
    1172            1 :         (void)BindCpuUtils::SetThreadAffinity(threadId, cpuIds);
    1173              :     }
    1174              : 
    1175              :     {
    1176           11 :         std::unique_lock<std::mutex> lk(manageThreadMutex_);
    1177           11 :         manageThreadStatus_ = PrePareForManageThread();
    1178           11 :         manageThreadCv_.notify_all();
    1179           11 :         if (manageThreadStatus_ != ThreadStatus::INIT_SUCCESS) {
    1180            5 :             return;
    1181              :         }
    1182           11 :     }
    1183              : 
    1184            6 :     struct event_info event = {};
    1185              :     // default wait timeout 2s
    1186            6 :     constexpr int32_t waitTimeout = 2000;
    1187           10 :     while (running_) {
    1188            5 :         const auto schedRet = halEschedWaitEvent(deviceId_, bindQueueGroupId_, 0U, waitTimeout, &event);
    1189            5 :         if (schedRet == DRV_ERROR_NONE) {
    1190            2 :             if (event.comm.event_id == EVENT_QS_MSG) {
    1191            1 :                 HandleBqsMsg(event);
    1192              :             } else {
    1193            1 :                 BQS_LOG_WARN(
    1194              :                     "Thread[%u] process unsupported eventId[%d] subEventId[%u].", 0,
    1195              :                     static_cast<int32_t>(event.comm.event_id), event.comm.subevent_id);
    1196              :             }
    1197            3 :         } else if (schedRet == DRV_ERROR_SCHED_WAIT_TIMEOUT) {
    1198            1 :             BQS_LOG_DEBUG("ManageQsEvent bind/unbind/query event waiting timeout");
    1199            1 :             continue;
    1200            2 :         } else if (schedRet == DRV_ERROR_PARA_ERROR) {
    1201            1 :             BQS_LOG_ERROR(
    1202              :                 "ManageQsEvent bind/unbind/query event failed, deviceId[%u] groupId[%u] error[%d].", deviceId_,
    1203              :                 bindQueueGroupId_, static_cast<int32_t>(schedRet));
    1204            1 :             break;
    1205              :         } else {
    1206              :             // LOG ERROR
    1207            1 :             BQS_LOG_ERROR(
    1208              :                 "ManageQsEvent bind/unbind/query event failed, deviceId[%u] groupId[%u] error[%d].", deviceId_,
    1209              :                 bindQueueGroupId_, static_cast<int32_t>(schedRet));
    1210              :         }
    1211              :     }
    1212            6 :     BQS_LOG_INFO("[RouterServer] ManageQsEvent of RouterServer thread exit.");
    1213              : }
    1214              : 
    1215            6 : BqsStatus RouterServer::SubscribeBufEvent() const
    1216              : {
    1217            6 :     BQS_LOG_INFO("[RouterServer] SubscribeBufEvent start.");
    1218            6 :     const bool needSubBufEvent =
    1219            6 :         static_cast<bool>(schedPolicy_ & static_cast<uint64_t>(SchedPolicy::POLICY_SUB_BUF_EVENT));
    1220            6 :     if (!needSubBufEvent) {
    1221            5 :         BQS_LOG_INFO("[RouterServer] needSubBufEvent is [%d]", static_cast<int32_t>(needSubBufEvent));
    1222            5 :         return BQS_STATUS_OK;
    1223              :     }
    1224              :     // load hccl so
    1225            1 :     dgw::HcclSoManager::GetInstance()->LoadSo();
    1226            1 :     return BQS_STATUS_OK;
    1227              : }
    1228              : 
    1229            3 : BqsStatus RouterServer::WaitSyncMsgProc()
    1230              : {
    1231            3 :     std::unique_lock<std::mutex> lk(mutex_);
    1232              : 
    1233              :     // waiting for aicpu thread processing
    1234            3 :     done_ = false;
    1235              :     // produce enqueue event
    1236            3 :     auto ret = QueueManager::GetInstance().EnqueueRelationEvent();
    1237            3 :     if (ret != BQS_STATUS_OK) {
    1238            1 :         return ret;
    1239              :     }
    1240              : 
    1241              :     // if numa_flag , produce extra enqueue event
    1242            2 :     if (numaFlag_) {
    1243            2 :         doneExtra_ = false;
    1244            2 :         auto result = QueueManager::GetInstance().EnqueueRelationEventExtra();
    1245            2 :         if (result != BQS_STATUS_OK) {
    1246            1 :             return result;
    1247              :         }
    1248              :     }
    1249              : 
    1250            1 :     BQS_LOG_INFO("[RouterServer] update config[add/del group, bind/unbind route], stage [server:waiting]");
    1251            3 :     (void)cv_.wait_for(lk, std::chrono::milliseconds(MAX_WAITING_NOTIFY), [this] { return done_ && doneExtra_; });
    1252            1 :     while ((!done_ || !doneExtra_) && (processing_ || processingExtra_)) {
    1253            0 :         cv_.wait(lk);
    1254              :     }
    1255            1 :     if (!done_ || !doneExtra_) {
    1256            1 :         QueueManager::GetInstance().LogErrorRelationQueueStatus();
    1257            1 :         BQS_LOG_ERROR(
    1258              :             "[RouterServer] update config[add/del group, bind/unbind route], stage [server:wait], timeout, "
    1259              :             "relation queue[enqueue cnt:%lu, dequeue cnt:%lu].",
    1260              :             StatisticManager::GetInstance().GetRelationEnqueCnt(),
    1261              :             StatisticManager::GetInstance().GetRelationDequeCnt());
    1262            1 :         return BQS_STATUS_TIMEOUT;
    1263              :     }
    1264              :     // get config process result
    1265            0 :     ret = static_cast<BqsStatus>(retCode_);
    1266              :     // rest retCode_ and updateCfgInfo_
    1267            0 :     retCode_ = static_cast<int32_t>(BQS_STATUS_OK);
    1268            0 :     return ret;
    1269            3 : }
    1270              : 
    1271           10 : void RouterServer::NotifyInitSuccess()
    1272              : {
    1273           10 :     BQS_LOG_RUN_INFO("schedule finish initing, now router is ready to handle msg");
    1274           10 :     readyToHandleMsg_ = true;
    1275           10 : }
    1276              : 
    1277           11 : bool RouterServer::AttachGroup()
    1278              : {
    1279           11 :     if (!needAttachGroup_) {
    1280            9 :         return true;
    1281              :     }
    1282            2 :     BQS_LOG_INFO("Begin to attach group");
    1283            2 :     std::stringstream grpNameStream(qsInitGroupName_);
    1284            2 :     std::string grpNameElement;
    1285            2 :     std::vector<std::string> groupNameVec;
    1286            4 :     while (getline(grpNameStream, grpNameElement, ',')) {
    1287            2 :         groupNameVec.emplace_back(grpNameElement);
    1288              :     }
    1289              : 
    1290            2 :     const int32_t halTimeOut = (RunContext::HOST == GetRunContext()) ? 3000 : -1;
    1291            2 :     for (const auto& grpName : groupNameVec) {
    1292            2 :         BQS_LOG_RUN_INFO("Begin to halGrpAttach group[%s].", grpName.c_str());
    1293            2 :         const auto drvRet = halGrpAttach(grpName.c_str(), halTimeOut);
    1294            2 :         if (drvRet != static_cast<int32_t>(DRV_ERROR_NONE)) {
    1295            2 :             BQS_LOG_ERROR("halGrpAttach group[%s] failed. ret[%d]", grpName.c_str(), drvRet);
    1296            2 :             return false;
    1297              :         }
    1298            0 :         BQS_LOG_RUN_INFO("halGrpAttach group[%s] success.", grpName.c_str());
    1299              :     }
    1300            0 :     return true;
    1301            2 : }
    1302              : 
    1303           46 : EntityInfo RouterServer::CreateBasicEntityInfo(const uint32_t id, const dgw::EntityType eType) const
    1304              : {
    1305           46 :     OptionalArg args = {};
    1306           46 :     args.eType = eType;
    1307           92 :     return EntityInfo(id, deviceId_, &args);
    1308              : }
    1309              : 
    1310              : } // namespace bqs
        

Generated by: LCOV version 2.0-1