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