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