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