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
|