Line data Source code
1 : /**
2 : * Copyright (c) 2025 Huawei Technologies Co., Ltd.
3 : * This program is free software, you can redistribute it and/or modify it under the terms and conditions of
4 : * CANN Open Software License Agreement Version 2.0 (the "License").
5 : * Please refer to the License for details. You may not use this file except in compliance with the License.
6 : * THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
7 : * INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE.
8 : * See LICENSE in the root of the software repository for the full text of the License.
9 : */
10 :
11 : #include "operator_kernel_model_batch_dequeue_buff.h"
12 :
13 : #include "aicpusd_status.h"
14 : #include "aicpusd_monitor.h"
15 : #include "aicpusd_model_execute.h"
16 : #include "aicpusd_resource_manager.h"
17 : #include "operator_kernel_common.h"
18 :
19 : namespace AicpuSchedule {
20 : namespace {
21 : const std::string KERNEL_MODEL_BATCH_DEQUEUE_BUFF = "modelBatchDequeueBuff";
22 : } // namespace
23 :
24 4 : int32_t OperatorKernelModelBatchDequeueBuff::Compute(const AicpuTaskInfo& kernelTaskInfo, const RunContext& taskContext)
25 : {
26 4 : aicpusd_info("Begin to batch dequeue buff. modelId[%u].", taskContext.modelId);
27 : const BatchDequeueBuffDesc* const batchDeqBufDesc =
28 4 : PtrToPtr<void, BatchDequeueBuffDesc>(ValueToPtr(kernelTaskInfo.paraBase));
29 4 : if (batchDeqBufDesc == nullptr) {
30 1 : aicpusd_err(
31 : "KernelTaskInfo paramBase is null, modelId[%u], streamId[%u], taskId[%u].", taskContext.modelId,
32 : taskContext.streamId, kernelTaskInfo.taskID);
33 1 : return AICPU_SCHEDULE_ERROR_PARAMETER_NOT_VALID;
34 : }
35 3 : BatchDequeueBuffInfo batchDeqBufInfo = {};
36 3 : auto ret = CheckAndParseBatchDeqBufParams(batchDeqBufDesc, kernelTaskInfo, taskContext, batchDeqBufInfo);
37 3 : if (ret != AICPU_SCHEDULE_OK) {
38 1 : return ret;
39 : }
40 :
41 2 : aicpusd_info("batch dequeue for %u queues", batchDeqBufInfo.inputNums);
42 3 : for (uint32_t i = 0U; i < batchDeqBufInfo.inputNums; ++i) {
43 : BufEnQueueBuffInfo queueInfo = {
44 2 : batchDeqBufInfo.queueIds[i], batchDeqBufInfo.deviceIds[i], batchDeqBufInfo.mbufAddrs[i]};
45 2 : const auto res = ModelAttachAndDequeueBuff(queueInfo, taskContext);
46 2 : if (res == AICPU_SCHEDULE_ERROR_MODEL_UNLOAD) {
47 1 : aicpusd_warn("model is destroy");
48 1 : bool* const pending = const_cast<bool*>(&taskContext.pending);
49 1 : *pending = true;
50 1 : return AICPU_SCHEDULE_OK;
51 : }
52 : }
53 :
54 1 : if (batchDeqBufInfo.alignOffsets != nullptr) {
55 0 : ret = AlignBatchDequeueBuff(batchDeqBufInfo, taskContext);
56 : }
57 1 : return ret;
58 : }
59 :
60 2 : int32_t OperatorKernelModelBatchDequeueBuff::CheckAndParseBatchDeqBufParams(
61 : const BatchDequeueBuffDesc* const batchDeqBufDesc, const AicpuTaskInfo& kernelTaskInfo,
62 : const RunContext& taskContext, BatchDequeueBuffInfo& batchDeqBufInfo) const
63 : {
64 2 : batchDeqBufInfo.inputNums = batchDeqBufDesc->inputNums;
65 2 : batchDeqBufInfo.alignInterval = batchDeqBufDesc->alignInterval;
66 2 : batchDeqBufInfo.alignOffsets = PtrToPtr<void, uint32_t>(ValueToPtr(batchDeqBufDesc->alignOffsetsAddr));
67 2 : batchDeqBufInfo.queueIds = PtrToPtr<void, uint32_t>(ValueToPtr(batchDeqBufDesc->queueIdsAddr));
68 2 : if (batchDeqBufInfo.queueIds == nullptr) {
69 0 : aicpusd_err(
70 : "KernelTaskInfo queueIds is null, modelId[%u], streamId[%u], taskId[%u]", taskContext.modelId,
71 : taskContext.streamId, kernelTaskInfo.taskID);
72 0 : return AICPU_SCHEDULE_ERROR_PARAMETER_NOT_VALID;
73 : }
74 2 : batchDeqBufInfo.mbufAddrs = PtrToPtr<void, uint64_t>(ValueToPtr(batchDeqBufDesc->mbufAddrsAddr));
75 2 : if (batchDeqBufInfo.mbufAddrs == nullptr) {
76 0 : aicpusd_err(
77 : "KernelTaskInfo mbufAddrs is null, modelId[%u], streamId[%u], taskId[%u]", taskContext.modelId,
78 : taskContext.streamId, kernelTaskInfo.taskID);
79 0 : return AICPU_SCHEDULE_ERROR_PARAMETER_NOT_VALID;
80 : }
81 2 : batchDeqBufInfo.deviceIds = PtrToPtr<void, int32_t>(ValueToPtr(batchDeqBufDesc->deviceIdAddr));
82 2 : if (batchDeqBufInfo.deviceIds == nullptr) {
83 0 : aicpusd_err(
84 : "KernelTaskInfo deviceIds is null, modelId[%u], streamId[%u], taskId[%u]", taskContext.modelId,
85 : taskContext.streamId, kernelTaskInfo.taskID);
86 0 : return AICPU_SCHEDULE_ERROR_PARAMETER_NOT_VALID;
87 : }
88 2 : return AICPU_SCHEDULE_OK;
89 : }
90 :
91 3 : int32_t OperatorKernelModelBatchDequeueBuff::ModelAttachAndDequeueBuff(
92 : BufEnQueueBuffInfo& queueInfo, const RunContext& taskContext, bool tryOnce) const
93 : {
94 3 : const auto drvRet = halQueueAttach(static_cast<uint32_t>(queueInfo.deviceId), queueInfo.queueID, 0);
95 3 : if ((drvRet != DRV_ERROR_NONE) && (drvRet != DRV_ERROR_REPEATED_INIT)) {
96 0 : aicpusd_err("Aicpusd attached queue[%u] failed ret[%d]", queueInfo.queueID, drvRet);
97 0 : return AICPU_SCHEDULE_ERROR_FROM_DRV;
98 : }
99 3 : aicpusd_info("queue id %u device id %d.", queueInfo.queueID, queueInfo.deviceId);
100 3 : return ModelDequeueBuffTaskKernel(queueInfo, taskContext, tryOnce);
101 : }
102 :
103 6 : int32_t OperatorKernelModelBatchDequeueBuff::ModelDequeueBuffTaskKernel(
104 : BufEnQueueBuffInfo& bufInfo, const RunContext& taskContext, bool tryOnce) const
105 : {
106 6 : Mbuf** const mBufPptr = PtrToPtr<void, Mbuf*>(ValueToPtr(bufInfo.mBufPtr));
107 6 : if (mBufPptr == nullptr) {
108 1 : aicpusd_err("param mBufPptr is null.");
109 1 : return AICPU_SCHEDULE_ERROR_PARAMETER_NOT_VALID;
110 : }
111 5 : const auto model = AicpuModelManager::GetInstance().GetModel(taskContext.modelId);
112 5 : if (model == nullptr) {
113 0 : aicpusd_err("Cannot get model by modelId:[%u], streamId[%u].", taskContext.modelId, taskContext.streamId);
114 0 : return AICPU_SCHEDULE_ERROR_PARAMETER_NOT_VALID;
115 : }
116 5 : const uint32_t queueId = bufInfo.queueID;
117 5 : const int32_t deviceId = bufInfo.deviceId;
118 5 : uint64_t respLen = 0U;
119 : while (true) {
120 5 : if (model->GetModelDestroyStatus()) {
121 2 : aicpusd_info("model is unloading, exit");
122 2 : return AICPU_SCHEDULE_ERROR_MODEL_UNLOAD;
123 : }
124 3 : const auto ret = halQueuePeek(static_cast<uint32_t>(deviceId), queueId, &respLen, 0);
125 3 : if ((ret == DRV_ERROR_NONE) && (respLen != 0U)) {
126 2 : aicpusd_info("halQueuePeek success");
127 2 : break;
128 : }
129 1 : if (tryOnce) {
130 1 : return AICPU_SCHEDULE_OK;
131 : }
132 0 : }
133 :
134 2 : const int32_t deqRet = DequeueBuff(bufInfo, taskContext, respLen, queueId, deviceId);
135 2 : if (deqRet != AICPU_SCHEDULE_OK) {
136 1 : aicpusd_warn(
137 : "DequeueBuff failed from queue[%u] in device[%u] failed, deqRet[%d], respLen[%u]", queueId, deviceId,
138 : static_cast<int32_t>(deqRet), respLen);
139 : }
140 2 : return deqRet;
141 : }
142 :
143 2 : int32_t OperatorKernelModelBatchDequeueBuff::DequeueBuff(
144 : BufEnQueueBuffInfo& bufInfo, const RunContext& taskContext, const uint64_t respLen, const uint32_t queueId,
145 : const int32_t deviceId) const
146 : {
147 : // mBufPptr在之前已经做过判空处理,此处不在进行校验
148 2 : Mbuf** const mBufPptr = reinterpret_cast<Mbuf**>(static_cast<uintptr_t>(bufInfo.mBufPtr));
149 2 : Mbuf* outMBuf = BufManager::GetInstance().MallocAndGuardBuf(respLen, taskContext.modelId);
150 2 : if (outMBuf == nullptr) {
151 1 : aicpusd_err("Failed to alloc mbuf, respLen[%u], modelId[%u].", respLen, taskContext.modelId);
152 1 : AicpuMonitor::GetInstance().SendKillMsgToTsd();
153 1 : return AICPU_SCHEDULE_ERROR_FROM_DRV;
154 : }
155 1 : auto drvRet = halMbufSetDataLen(outMBuf, static_cast<uint64_t>(respLen));
156 1 : if (drvRet != DRV_ERROR_NONE) {
157 0 : aicpusd_err("halMbufSetDataLen error, queue id:%u, drvRet=%d", queueId, drvRet);
158 0 : return AICPU_SCHEDULE_ERROR_FROM_DRV;
159 : }
160 :
161 1 : uint32_t headSize = 0U;
162 1 : void* headBuf = nullptr;
163 1 : drvRet = halMbufGetPrivInfo(outMBuf, &headBuf, &headSize);
164 1 : if (drvRet != DRV_ERROR_NONE) {
165 0 : aicpusd_err("Failed to get head info in input information, ret[%d].", drvRet);
166 0 : return AICPU_SCHEDULE_ERROR_FROM_DRV;
167 : }
168 :
169 1 : void* dataAddrPtr = nullptr;
170 1 : drvRet = OperatorKernelCommon::GetMbufDataPtr(
171 : static_cast<uint64_t>(reinterpret_cast<uintptr_t>(&outMBuf)), &dataAddrPtr);
172 1 : if (drvRet != AICPU_SCHEDULE_OK) {
173 0 : aicpusd_err("Failed to get mbuf data addr. ret is [%d]", drvRet);
174 0 : return drvRet;
175 : }
176 :
177 1 : constexpr size_t totalLen = sizeof(struct buff_iovec) + sizeof(struct iovec_info);
178 1 : std::unique_ptr<char_t[]> vecUniquePtr(new (std::nothrow) char_t[totalLen], std::default_delete<char_t[]>());
179 1 : if (vecUniquePtr == nullptr) {
180 0 : aicpusd_err("failed to alloc memory for vecUniquePtr, size[%zu].", totalLen);
181 0 : return AICPU_SCHEDULE_ERROR_INNER_ERROR;
182 : }
183 1 : buff_iovec* const buffvec = reinterpret_cast<buff_iovec*>(vecUniquePtr.get());
184 1 : buffvec->context_base = headBuf;
185 1 : buffvec->context_len = headSize;
186 1 : buffvec->count = 1U;
187 1 : buffvec->ptr[0U].iovec_base = dataAddrPtr;
188 1 : buffvec->ptr[0U].len = respLen;
189 1 : drvRet = halQueueDeQueueBuff(static_cast<uint32_t>(deviceId), queueId, buffvec, -1);
190 1 : if (drvRet != DRV_ERROR_NONE) {
191 0 : aicpusd_err(
192 : "halQueueDeQueueBuff to queue[%u] in device[%d] failed, error[%d]", queueId, deviceId,
193 : static_cast<int32_t>(drvRet));
194 0 : return AICPU_SCHEDULE_ERROR_FROM_DRV;
195 : }
196 :
197 1 : *mBufPptr = PtrToPtr<void, Mbuf>(outMBuf);
198 :
199 1 : (void)ProcessMbufHeadInDequeueTask(taskContext.modelId, headBuf, headSize);
200 1 : (void)SetModelEndOfSequence(taskContext.modelId, headBuf, headSize);
201 :
202 1 : OperatorKernelCommon::TraceQueueData(taskContext, headBuf, headSize, "DequeuedBuff");
203 1 : return AICPU_SCHEDULE_OK;
204 1 : }
205 :
206 : // if max(inputs timestamp-alignOffset)-min(inputs timestamp-alignOffset) < alignInterval, return ok
207 : // else delete oldest data, then re-dequeue data until inputs timestamp alignment is satisfied or the queue is empty.
208 3 : int32_t OperatorKernelModelBatchDequeueBuff::AlignBatchDequeueBuff(
209 : BatchDequeueBuffInfo& batchDeqBufInfo, const RunContext& taskContext) const
210 : {
211 : BatchDequeueInfo batchDeqInfo = {
212 3 : batchDeqBufInfo.inputNums, batchDeqBufInfo.alignInterval, batchDeqBufInfo.alignOffsets,
213 3 : batchDeqBufInfo.queueIds, batchDeqBufInfo.mbufAddrs};
214 : while (true) {
215 3 : uint32_t maxAlignTimestamp = 0U;
216 3 : uint32_t minAlignTimestamp = UINT32_MAX;
217 3 : uint32_t minTimestampIndex = 0U;
218 3 : auto ret = AlignTimestamp(batchDeqInfo, taskContext, maxAlignTimestamp, minAlignTimestamp, minTimestampIndex);
219 3 : if (ret != AICPU_SCHEDULE_OK) {
220 0 : aicpusd_err("AlignTimestamp failed! ret[%ld]", ret);
221 3 : return ret;
222 : }
223 3 : if ((maxAlignTimestamp - minAlignTimestamp) <= batchDeqInfo.alignInterval) {
224 1 : return AICPU_SCHEDULE_OK;
225 : }
226 : BufEnQueueBuffInfo queBuffInfo = {
227 2 : batchDeqBufInfo.queueIds[minTimestampIndex], batchDeqBufInfo.deviceIds[minTimestampIndex],
228 2 : batchDeqBufInfo.mbufAddrs[minTimestampIndex]};
229 2 : ret = ModelDequeueBuffTaskKernel(queBuffInfo, taskContext);
230 2 : if (ret == AICPU_SCHEDULE_ERROR_MODEL_UNLOAD) {
231 1 : bool* const pending = const_cast<bool*>(&taskContext.pending);
232 1 : *pending = true;
233 1 : return AICPU_SCHEDULE_OK;
234 : }
235 1 : if (ret != AICPU_SCHEDULE_OK) {
236 1 : aicpusd_err("ModelDequeueBuffTaskKernel failed! ret[%ld]", ret);
237 1 : return ret;
238 : }
239 0 : }
240 : return AICPU_SCHEDULE_OK;
241 : }
242 :
243 6 : REGISTER_OPERATOR_KERNEL(KERNEL_MODEL_BATCH_DEQUEUE_BUFF, OperatorKernelModelBatchDequeueBuff);
244 : } // namespace AicpuSchedule
|