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 : #ifndef QUEUE_PROCESSOR_H
12 : #define QUEUE_PROCESSOR_H
13 :
14 : #include <mutex>
15 : #include <memory>
16 : #include <map>
17 : #include "queue.h"
18 : #include "runtime/rt_mem_queue.h"
19 : #include "mmpa/mmpa_api.h"
20 :
21 : namespace acl {
22 :
23 : using QueueDataMutex = struct TagQueueDataMutex {
24 : std::mutex muForEnqueue;
25 : std::mutex muForDequeue;
26 : };
27 :
28 : enum PID_QUERY_TYPE : int32_t {
29 : CP_PID,
30 : QS_PID
31 : };
32 :
33 : constexpr int32_t MSEC_TO_USEC = 1000;
34 :
35 : constexpr size_t QUERY_BUFF_GRP_MAX_NUM = 1024U;
36 :
37 : using QueueDataMutexPtr = std::shared_ptr<QueueDataMutex>;
38 :
39 : class QueueProcessor {
40 : public:
41 : virtual aclError acltdtCreateQueue(const acltdtQueueAttr *const attr, uint32_t *const qid) = 0;
42 :
43 : virtual aclError acltdtDestroyQueue(const uint32_t qid) = 0;
44 :
45 : virtual aclError acltdtEnqueue(const uint32_t qid, const acltdtBuf buf, const int32_t timeout);
46 :
47 : virtual aclError acltdtDequeue(const uint32_t qid, acltdtBuf *const buf, const int32_t timeout);
48 :
49 : virtual aclError acltdtGrantQueue(const uint32_t qid, const int32_t pid, const uint32_t permission,
50 : const int32_t timeout);
51 :
52 : virtual aclError acltdtAttachQueue(const uint32_t qid, const int32_t timeout,
53 : uint32_t *const permission);
54 :
55 : virtual aclError acltdtBindQueueRoutes(acltdtQueueRouteList *const qRouteList) = 0;
56 :
57 : virtual aclError acltdtUnbindQueueRoutes(acltdtQueueRouteList *const qRouteList) = 0;
58 :
59 : virtual aclError acltdtQueryQueueRoutes(const acltdtQueueRouteQueryInfo *const queryInfo,
60 : acltdtQueueRouteList *const qRouteList) = 0;
61 :
62 : virtual aclError acltdtAllocBuf(const size_t size, const uint32_t type, acltdtBuf *const buf) = 0;
63 :
64 : virtual aclError acltdtAllocBufData(const size_t size, const uint32_t type, acltdtBuf *const buf);
65 :
66 : virtual aclError acltdtFreeBuf(acltdtBuf buf);
67 :
68 : virtual aclError acltdtSetBufDataLen(const acltdtBuf buf, const size_t len);
69 :
70 : virtual aclError acltdtGetBufDataLen(const acltdtBuf buf, size_t *const len);
71 :
72 : virtual aclError acltdtAppendBufChain(const acltdtBuf headBuf, const acltdtBuf buf);
73 :
74 : virtual aclError acltdtGetBufChainNum(const acltdtBuf headBuf, uint32_t *const num);
75 :
76 : virtual aclError acltdtGetBufFromChain(const acltdtBuf headBuf, const uint32_t index, acltdtBuf *const buf);
77 :
78 : virtual aclError acltdtGetBufData(const acltdtBuf buf, void **const dataPtr, size_t *const size);
79 :
80 : virtual aclError acltdtGetBufUserData(const acltdtBuf buf, void *dataPtr,
81 : const size_t size, const size_t offset);
82 :
83 : virtual aclError acltdtSetBufUserData(acltdtBuf buf, const void *dataPtr,
84 : const size_t size, const size_t offset);
85 :
86 : virtual aclError acltdtCopyBufRef(const acltdtBuf buf, acltdtBuf *const newBuf);
87 :
88 : virtual aclError QueryAllocGroup();
89 :
90 : virtual aclError QueryGroupId(const std::string &grpName);
91 :
92 : aclError InitQueueSchedule(const int32_t devId) const;
93 :
94 : aclError acltdtDestroyQueueOndevice(const uint32_t qid, const bool isThreadMode = false);
95 :
96 : aclError SendBindUnbindMsgOnDevice(acltdtQueueRouteList *const qRouteList,
97 : const bool isBind, rtEschedEventSummary_t &eventSum, rtEschedEventReply_t &ack) const;
98 :
99 : aclError SendConnectQsMsg(const int32_t deviceId, rtEschedEventSummary_t &eventSum, rtEschedEventReply_t &ack);
100 : aclError GetDstInfo(const int32_t deviceId, const PID_QUERY_TYPE type,
101 : int32_t &dstPid, const bool isThreadMode = false) const;
102 : aclError GetQueuePermission(const int32_t deviceId, uint32_t qid, rtMemQueueShareAttr_t &permission) const;
103 : aclError GetQueueRouteNum(const acltdtQueueRouteQueryInfo *const queryInfo,
104 : const int32_t deviceId,
105 : rtEschedEventSummary_t &eventSum,
106 : rtEschedEventReply_t &ack,
107 : size_t &routeNum) const;
108 :
109 : aclError QueryQueueRoutesOnDevice(const acltdtQueueRouteQueryInfo *const queryInfo, const size_t routeNum,
110 : rtEschedEventSummary_t &eventSum, rtEschedEventReply_t &ack, acltdtQueueRouteList *const qRouteList) const;
111 :
112 : QueueDataMutexPtr GetMutexForData(const uint32_t qid);
113 : void DeleteMutexForData(const uint32_t qid);
114 :
115 : uint64_t GetTimestamp() const;
116 :
117 : aclError GetDeviceId(int32_t& deviceId) const;
118 :
119 : virtual aclError acltdtEnqueueData(const uint32_t qid, const void *const data, const size_t dataSize,
120 : const void *const userData, const size_t userDataSize, const int32_t timeout, const uint32_t rsv);
121 :
122 : aclError acltdtDequeueData(const uint32_t qid, void *const data, const size_t dataSize, size_t *const retDataSize,
123 : void *const userData, const size_t userDataSize, const int32_t timeout);
124 :
125 : // set queue attr to default,depth is 8,name is empty
126 : static void acltdtSetDefaultQueueAttr(acltdtQueueAttr &attr);
127 :
128 : aclError acltdtCreateQueueWithAttr(const int32_t deviceId, const acltdtQueueAttr *const attr,
129 : uint32_t *const qid) const;
130 :
131 35 : QueueProcessor() = default;
132 35 : virtual ~QueueProcessor() = default;
133 :
134 : // not allow copy constructor and assignment operators
135 : QueueProcessor(const QueueProcessor &) = delete;
136 :
137 : QueueProcessor &operator=(const QueueProcessor &) = delete;
138 :
139 : QueueProcessor(QueueProcessor &&) = delete;
140 :
141 : QueueProcessor &&operator=(QueueProcessor &&) = delete;
142 :
143 : protected:
144 : std::recursive_mutex muForQueueCtrl_;
145 : std::mutex muForQueryGroup_;
146 : std::mutex muForQueueMap_;
147 : std::mutex muForCreateGroup_;
148 : std::map<uint32_t, QueueDataMutexPtr> muForQueue_;
149 : bool isQsInit_ = false;
150 : uint32_t qsContactId_ = 0U;
151 : int32_t qsGroupId_ = 0;
152 : static bool isInitQs_;
153 : static bool isMbufInit_;
154 : // new mbuf version is enhanced, mbuf can not be operated after enqueue
155 : bool isMbufEnhanced_ = false;
156 : };
157 : }
158 : #endif // QUEUE_PROCESS_H
|