/src/open62541_15/src/server/ua_server_async.c
Line | Count | Source |
1 | | /* This Source Code Form is subject to the terms of the Mozilla Public |
2 | | * License, v. 2.0. If a copy of the MPL was not distributed with this |
3 | | * file, You can obtain one at http://mozilla.org/MPL/2.0/. |
4 | | * |
5 | | * Copyright 2019 (c) Fraunhofer IOSB (Author: Klaus Schick) |
6 | | * Copyright 2019, 2025 (c) Fraunhofer IOSB (Author: Julius Pfrommer) |
7 | | * Copyright 2026 (c) o6 Automation GmbH (Author: Julius Pfrommer) |
8 | | */ |
9 | | |
10 | | #include "ua_server_internal.h" |
11 | | |
12 | | /* The layout of the results array is: |
13 | | * [results-array] | padding | UA_AsyncResponse | padding | [UA_AsyncOperation] |
14 | | * |
15 | | * We need to take care about memory alignment (padding). */ |
16 | | static void * |
17 | | allocateResultsArray(const UA_DataType *resultsType, size_t resultsLen, |
18 | 253 | UA_AsyncResponse **resp, UA_AsyncOperation **ops) { |
19 | 253 | const size_t padding = sizeof(size_t) - 1; |
20 | 253 | const size_t fixedSize = sizeof(UA_AsyncResponse) + 2 * padding; |
21 | 253 | const size_t elementSize = |
22 | 253 | resultsType->memSize + sizeof(UA_AsyncOperation); |
23 | | |
24 | | /* Reserve the maximum padding at both alignment boundaries. */ |
25 | 253 | if(resultsLen > (SIZE_MAX - fixedSize) / elementSize) |
26 | 0 | return NULL; |
27 | | |
28 | 253 | size_t responseBegin = |
29 | 253 | (resultsType->memSize * resultsLen + padding) & ~padding; |
30 | 253 | size_t opsBegin = |
31 | 253 | (responseBegin + sizeof(UA_AsyncResponse) + padding) & ~padding; |
32 | 253 | size_t allocationSize = |
33 | 253 | opsBegin + sizeof(UA_AsyncOperation) * resultsLen; |
34 | | |
35 | 253 | void *arr = UA_calloc(1, allocationSize); |
36 | 253 | if(!arr) |
37 | 0 | return NULL; |
38 | 253 | uintptr_t arrMem = (uintptr_t)arr; |
39 | 253 | *resp = (UA_AsyncResponse*)(arrMem + responseBegin); |
40 | 253 | *ops = (UA_AsyncOperation*)(arrMem + opsBegin); |
41 | 253 | return arr; |
42 | 253 | } |
43 | | |
44 | | /* Cancel the operation, but don't _clear it here */ |
45 | | static void |
46 | | UA_AsyncOperation_cancel(UA_Server *server, UA_AsyncOperation *op, |
47 | 0 | UA_StatusCode opstatus) { |
48 | 0 | UA_ServerConfig *sc = &server->config; |
49 | 0 | void *cancelPtr = NULL; |
50 | | |
51 | | /* Set the status and get the pointer that identifies the operation */ |
52 | 0 | switch(op->asyncOperationType) { |
53 | 0 | case UA_ASYNCOPERATIONTYPE_READ_REQUEST: |
54 | 0 | cancelPtr = op->output.read; |
55 | 0 | op->output.read->hasStatus = true; |
56 | 0 | op->output.read->status = opstatus; |
57 | 0 | break; |
58 | 0 | case UA_ASYNCOPERATIONTYPE_READ_DIRECT: |
59 | 0 | cancelPtr = &op->output.directRead; |
60 | 0 | op->output.directRead.hasStatus = true; |
61 | 0 | op->output.directRead.status = opstatus; |
62 | 0 | break; |
63 | 0 | case UA_ASYNCOPERATIONTYPE_WRITE_REQUEST: |
64 | 0 | cancelPtr = &op->context.writeValue.value; |
65 | 0 | *op->output.write = opstatus; |
66 | 0 | break; |
67 | 0 | case UA_ASYNCOPERATIONTYPE_WRITE_DIRECT: |
68 | 0 | cancelPtr = &op->context.writeValue.value; |
69 | 0 | op->output.directWrite = opstatus; |
70 | 0 | break; |
71 | 0 | case UA_ASYNCOPERATIONTYPE_CALL_REQUEST: |
72 | | /* outputArguments is always an allocated pointer, also if the length is zero */ |
73 | 0 | cancelPtr = op->output.call->outputArguments; |
74 | 0 | op->output.call->statusCode = opstatus; |
75 | 0 | break; |
76 | 0 | case UA_ASYNCOPERATIONTYPE_CALL_DIRECT: |
77 | | /* outputArguments is always an allocated pointer, also if the length is zero */ |
78 | 0 | cancelPtr = op->output.directCall.outputArguments; |
79 | 0 | op->output.directCall.statusCode = opstatus; |
80 | 0 | break; |
81 | 0 | default: UA_assert(false); return; |
82 | 0 | } |
83 | | |
84 | | /* Notify the application that it must no longer set the async result */ |
85 | 0 | if(sc->asyncOperationCancelCallback) |
86 | 0 | sc->asyncOperationCancelCallback(server, cancelPtr); |
87 | 0 | } |
88 | | |
89 | | static void |
90 | 0 | UA_AsyncOperation_delete(UA_AsyncOperation *op) { |
91 | 0 | UA_assert(op->asyncOperationType >= UA_ASYNCOPERATIONTYPE_CALL_DIRECT); |
92 | 0 | switch(op->asyncOperationType) { |
93 | 0 | case UA_ASYNCOPERATIONTYPE_READ_DIRECT: |
94 | 0 | UA_DataValue_clear(&op->output.directRead); |
95 | 0 | break; |
96 | 0 | case UA_ASYNCOPERATIONTYPE_WRITE_DIRECT: |
97 | 0 | break; |
98 | 0 | case UA_ASYNCOPERATIONTYPE_CALL_DIRECT: |
99 | 0 | UA_CallMethodResult_clear(&op->output.directCall); |
100 | 0 | break; |
101 | 0 | default: UA_assert(false); break; |
102 | 0 | } |
103 | 0 | UA_free(op); |
104 | 0 | } |
105 | | |
106 | | static void |
107 | 0 | UA_AsyncResponse_delete(UA_AsyncResponse *ar) { |
108 | 0 | UA_NodeId_clear(&ar->sessionId); |
109 | | |
110 | | /* Clean up the results array last. Because the results array memory also |
111 | | * includes ar. */ |
112 | 0 | void *arr = NULL; |
113 | 0 | size_t arrSize = 0; |
114 | 0 | const UA_DataType *arrType; |
115 | 0 | if(ar->responseType == &UA_TYPES[UA_TYPES_CALLRESPONSE]) { |
116 | 0 | arr = ar->response.callResponse.results; |
117 | 0 | arrSize = ar->response.callResponse.resultsSize; |
118 | 0 | ar->response.callResponse.results = NULL; |
119 | 0 | ar->response.callResponse.resultsSize = 0; |
120 | 0 | arrType = &UA_TYPES[UA_TYPES_CALLMETHODRESULT]; |
121 | 0 | } else if(ar->responseType == &UA_TYPES[UA_TYPES_READRESPONSE]) { |
122 | 0 | arr = ar->response.readResponse.results; |
123 | 0 | arrSize = ar->response.readResponse.resultsSize; |
124 | 0 | ar->response.readResponse.results = NULL; |
125 | 0 | ar->response.readResponse.resultsSize = 0; |
126 | 0 | arrType = &UA_TYPES[UA_TYPES_DATAVALUE]; |
127 | 0 | } else /* if(ar->responseType == &UA_TYPES[UA_TYPES_WRITERESPONSE]) */ { |
128 | 0 | UA_assert(ar->responseType == &UA_TYPES[UA_TYPES_WRITERESPONSE]); |
129 | 0 | arr = ar->response.writeResponse.results; |
130 | 0 | arrSize = ar->response.writeResponse.resultsSize; |
131 | 0 | ar->response.writeResponse.results = NULL; |
132 | 0 | ar->response.writeResponse.resultsSize = 0; |
133 | 0 | arrType = &UA_TYPES[UA_TYPES_STATUSCODE]; |
134 | 0 | } |
135 | 0 | UA_clear(&ar->response.callResponse, ar->responseType); |
136 | 0 | UA_Array_delete(arr, arrSize, arrType); |
137 | 0 | } |
138 | | |
139 | | static void |
140 | | notifyServiceEnd(UA_Server *server, UA_AsyncResponse *ar, |
141 | 0 | UA_Session *session, UA_SecureChannel *sc) { |
142 | | /* Nothing to do? */ |
143 | 0 | UA_ServerConfig *config = UA_Server_getConfig(server); |
144 | 0 | if(!config->globalNotificationCallback && !config->serviceNotificationCallback) |
145 | 0 | return; |
146 | | |
147 | | /* Collect the payload */ |
148 | 0 | UA_NodeId sessionId = (session) ? session->sessionId : UA_NODEID_NULL; |
149 | 0 | UA_UInt32 secureChannelId = (sc) ? sc->securityToken.channelId : 0; |
150 | 0 | UA_NodeId serviceTypeId; |
151 | 0 | if(ar->responseType == &UA_TYPES[UA_TYPES_CALLRESPONSE]) { |
152 | 0 | serviceTypeId = UA_TYPES[UA_TYPES_CALLREQUEST].typeId; |
153 | 0 | } else if(ar->responseType == &UA_TYPES[UA_TYPES_READRESPONSE]) { |
154 | 0 | serviceTypeId = UA_TYPES[UA_TYPES_READREQUEST].typeId; |
155 | 0 | } else /* if(ar->responseType == &UA_TYPES[UA_TYPES_WRITERESPONSE]) */ { |
156 | 0 | serviceTypeId = UA_TYPES[UA_TYPES_WRITEREQUEST].typeId; |
157 | 0 | } |
158 | | |
159 | | /* Notify the application */ |
160 | 0 | static UA_THREAD_LOCAL UA_KeyValuePair notifyPayload[4] = { |
161 | 0 | {{0, UA_STRING_STATIC("securechannel-id")}, {0}}, |
162 | 0 | {{0, UA_STRING_STATIC("session-id")}, {0}}, |
163 | 0 | {{0, UA_STRING_STATIC("request-id")}, {0}}, |
164 | 0 | {{0, UA_STRING_STATIC("service-type")}, {0}} |
165 | 0 | }; |
166 | 0 | UA_KeyValueMap notifyPayloadMap = {4, notifyPayload}; |
167 | 0 | if(config->serviceNotificationCallback || config->globalNotificationCallback) { |
168 | 0 | UA_Variant_setScalar(¬ifyPayload[0].value, &secureChannelId, |
169 | 0 | &UA_TYPES[UA_TYPES_UINT32]); |
170 | 0 | UA_Variant_setScalar(¬ifyPayload[1].value, &sessionId, |
171 | 0 | &UA_TYPES[UA_TYPES_NODEID]); |
172 | 0 | UA_Variant_setScalar(¬ifyPayload[2].value, &ar->requestId, |
173 | 0 | &UA_TYPES[UA_TYPES_UINT32]); |
174 | 0 | UA_Variant_setScalar(¬ifyPayload[3].value, &serviceTypeId, |
175 | 0 | &UA_TYPES[UA_TYPES_NODEID]); |
176 | 0 | } |
177 | |
|
178 | 0 | UA_ApplicationNotificationType nt = UA_APPLICATIONNOTIFICATIONTYPE_SERVICE_END; |
179 | 0 | if(config->serviceNotificationCallback) |
180 | 0 | config->serviceNotificationCallback(server, nt, notifyPayloadMap); |
181 | 0 | if(config->globalNotificationCallback) |
182 | 0 | config->globalNotificationCallback(server, nt, notifyPayloadMap); |
183 | 0 | } |
184 | | |
185 | | static void |
186 | 0 | sendAsyncResponse(UA_Server *server, UA_AsyncResponse *ar) { |
187 | 0 | UA_assert(ar->opCountdown == 0); |
188 | | |
189 | | /* Get the session */ |
190 | 0 | UA_Session *session = getSessionById(server, &ar->sessionId); |
191 | 0 | UA_SecureChannel *channel = (session) ? session->channel : NULL; |
192 | | |
193 | | /* Notify that processing the service has ended */ |
194 | 0 | notifyServiceEnd(server, ar, session, channel); |
195 | | |
196 | | /* Check the session */ |
197 | 0 | if(!session) { |
198 | 0 | UA_LOG_WARNING(server->config.logging, UA_LOGCATEGORY_SERVER, |
199 | 0 | "Async Service: Session %N no longer exists", ar->sessionId); |
200 | 0 | return; |
201 | 0 | } |
202 | | |
203 | | /* Check the channel */ |
204 | 0 | if(!channel) { |
205 | 0 | UA_LOG_WARNING_SESSION(server->config.logging, session, |
206 | 0 | "Async Service Response cannot be sent. " |
207 | 0 | "No SecureChannel for the session."); |
208 | 0 | return; |
209 | 0 | } |
210 | | |
211 | | /* Set the request handle */ |
212 | 0 | UA_ResponseHeader *responseHeader = (UA_ResponseHeader*) |
213 | 0 | &ar->response.callResponse.responseHeader; |
214 | 0 | responseHeader->requestHandle = ar->requestHandle; |
215 | | |
216 | | /* Send the Response */ |
217 | 0 | UA_StatusCode res = sendResponse(server, channel, ar->requestId, |
218 | 0 | (UA_Response*)&ar->response, ar->responseType); |
219 | 0 | if(res != UA_STATUSCODE_GOOD) { |
220 | 0 | UA_LOG_WARNING_SESSION(server->config.logging, session, |
221 | 0 | "Async Response for Req# %" PRIu32 " failed " |
222 | 0 | "with StatusCode %s", ar->requestId, |
223 | 0 | UA_StatusCode_name(res)); |
224 | 0 | } |
225 | 0 | } |
226 | | |
227 | | static void |
228 | 0 | directOpCallback(UA_Server *server, UA_AsyncOperation *op) { |
229 | 0 | switch(op->asyncOperationType) { |
230 | 0 | case UA_ASYNCOPERATIONTYPE_READ_DIRECT: |
231 | 0 | op->handling.callback.method.read(server, |
232 | 0 | op->handling.callback.context, |
233 | 0 | &op->output.directRead); |
234 | 0 | break; |
235 | 0 | case UA_ASYNCOPERATIONTYPE_WRITE_DIRECT: |
236 | 0 | op->handling.callback.method.write(server, |
237 | 0 | op->handling.callback.context, |
238 | 0 | op->output.directWrite); |
239 | 0 | break; |
240 | 0 | case UA_ASYNCOPERATIONTYPE_CALL_DIRECT: |
241 | 0 | op->handling.callback.method.call(server, |
242 | 0 | op->handling.callback.context, |
243 | 0 | &op->output.directCall); |
244 | 0 | break; |
245 | 0 | default: UA_assert(false); break; |
246 | 0 | } |
247 | 0 | } |
248 | | |
249 | | /* Called from the EventLoop via a delayed callback */ |
250 | | static void |
251 | 19.4k | UA_AsyncManager_processReady(UA_Server *server, UA_AsyncManager *am) { |
252 | 19.4k | lockServer(server); |
253 | | |
254 | | /* Reset the delayed callback */ |
255 | 19.4k | UA_atomic_xchg((void**)&am->dc.callback, NULL); |
256 | | |
257 | | /* Process ready direct operations and free them */ |
258 | 19.4k | UA_AsyncOperation *op = NULL, *op_tmp = NULL; |
259 | 19.4k | TAILQ_FOREACH_SAFE(op, &am->readyOps, pointers, op_tmp) { |
260 | 0 | TAILQ_REMOVE(&am->readyOps, op, pointers); |
261 | 0 | am->opsCount--; |
262 | 0 | directOpCallback(server, op); |
263 | 0 | UA_AsyncOperation_delete(op); |
264 | 0 | } |
265 | | |
266 | | /* Send out ready responses */ |
267 | 19.4k | UA_AsyncResponse *ar, *temp; |
268 | 19.4k | TAILQ_FOREACH_SAFE(ar, &am->readyResponses, pointers, temp) { |
269 | 0 | TAILQ_REMOVE(&am->readyResponses, ar, pointers); |
270 | 0 | sendAsyncResponse(server, ar); |
271 | 0 | UA_AsyncResponse_delete(ar); |
272 | 0 | } |
273 | | |
274 | 19.4k | unlockServer(server); |
275 | 19.4k | } |
276 | | |
277 | | static void |
278 | 0 | processOperationResult(UA_Server *server, UA_AsyncOperation *op) { |
279 | 0 | UA_AsyncManager *am = &server->asyncManager; |
280 | 0 | if(op->asyncOperationType >= UA_ASYNCOPERATIONTYPE_CALL_DIRECT) { |
281 | | /* Direct operation */ |
282 | 0 | TAILQ_REMOVE(&am->waitingOps, op, pointers); |
283 | 0 | TAILQ_INSERT_TAIL(&am->readyOps, op, pointers); |
284 | 0 | } else { |
285 | | /* Part of a service request */ |
286 | 0 | TAILQ_REMOVE(&am->waitingOps, op, pointers); |
287 | 0 | am->opsCount--; |
288 | |
|
289 | 0 | UA_AsyncResponse *ar = op->handling.response; |
290 | 0 | ar->opCountdown -= 1; |
291 | 0 | if(ar->opCountdown > 0) |
292 | 0 | return; |
293 | | |
294 | | /* Enqueue ar in the readyResponses */ |
295 | 0 | TAILQ_REMOVE(&am->waitingResponses, ar, pointers); |
296 | 0 | TAILQ_INSERT_TAIL(&am->readyResponses, ar, pointers); |
297 | 0 | } |
298 | | |
299 | | /* Trigger the main server thread to handle ready operations and responses */ |
300 | 0 | if(am->dc.callback == NULL) { |
301 | 0 | UA_EventLoop *el = server->config.eventLoop; |
302 | 0 | am->dc.callback = (UA_Callback)UA_AsyncManager_processReady; |
303 | 0 | am->dc.application = server; |
304 | 0 | am->dc.context = am; |
305 | 0 | el->addDelayedCallback(el, &am->dc); |
306 | 0 | el->cancel(el); /* Wake up the EventLoop if currently waiting in select() */ |
307 | 0 | } |
308 | 0 | } |
309 | | |
310 | | /* Check if any operations have timed out */ |
311 | | static void |
312 | 0 | checkTimeouts(UA_Server *server, void *_) { |
313 | | /* Timeouts are not configured */ |
314 | 0 | if(server->config.asyncOperationTimeout <= 0.0) |
315 | 0 | return; |
316 | | |
317 | 0 | lockServer(server); |
318 | |
|
319 | 0 | UA_EventLoop *el = server->config.eventLoop; |
320 | 0 | UA_AsyncManager *am = &server->asyncManager; |
321 | 0 | const UA_DateTime tNow = el->dateTime_nowMonotonic(el); |
322 | | |
323 | | /* Loop over the waiting ops */ |
324 | 0 | UA_AsyncOperation *op = NULL, *op_tmp = NULL; |
325 | 0 | TAILQ_FOREACH_SAFE(op, &am->waitingOps, pointers, op_tmp) { |
326 | | /* Check the timeout */ |
327 | 0 | if(op->asyncOperationType <= UA_ASYNCOPERATIONTYPE_WRITE_REQUEST) { |
328 | 0 | if(tNow <= op->handling.response->timeout) |
329 | 0 | continue; |
330 | 0 | } else { |
331 | 0 | if(tNow <= op->handling.callback.timeout) |
332 | 0 | continue; |
333 | 0 | } |
334 | | |
335 | 0 | UA_LOG_WARNING(server->config.logging, UA_LOGCATEGORY_SERVER, |
336 | 0 | "Operation was removed due to a timeout"); |
337 | | |
338 | | /* Mark operation as timed out integrate */ |
339 | 0 | UA_AsyncOperation_cancel(server, op, UA_STATUSCODE_BADTIMEOUT); |
340 | 0 | processOperationResult(server, op); |
341 | 0 | } |
342 | |
|
343 | 0 | unlockServer(server); |
344 | 0 | } |
345 | | |
346 | | void |
347 | 19.4k | UA_AsyncManager_init(UA_AsyncManager *am, UA_Server *server) { |
348 | 19.4k | memset(am, 0, sizeof(UA_AsyncManager)); |
349 | 19.4k | TAILQ_INIT(&am->waitingResponses); |
350 | 19.4k | TAILQ_INIT(&am->readyResponses); |
351 | 19.4k | TAILQ_INIT(&am->waitingOps); |
352 | 19.4k | TAILQ_INIT(&am->readyOps); |
353 | 19.4k | } |
354 | | |
355 | 536 | void UA_AsyncManager_start(UA_AsyncManager *am, UA_Server *server) { |
356 | | /* Add a regular callback for cleanup and sending finished responses at a |
357 | | * 1s interval. */ |
358 | 536 | UA_StatusCode res = addRepeatedCallback(server, (UA_ServerCallback)checkTimeouts, |
359 | 536 | NULL, 1000.0, &am->checkTimeoutCallbackId); |
360 | 536 | if(res != UA_STATUSCODE_GOOD) { |
361 | 0 | UA_LOG_WARNING(server->config.logging, UA_LOGCATEGORY_SERVER, |
362 | 0 | "Failed to register async timeout callback. " |
363 | 0 | "Async operations will not be cleaned up on timeout. StatusCode: %s", |
364 | 0 | UA_StatusCode_name(res)); |
365 | 0 | am->checkTimeoutCallbackId = 0; |
366 | 0 | } |
367 | 536 | } |
368 | | |
369 | 536 | void UA_AsyncManager_stop(UA_AsyncManager *am, UA_Server *server) { |
370 | 536 | removeCallback(server, am->checkTimeoutCallbackId); |
371 | 536 | if(am->dc.callback) { |
372 | 0 | UA_EventLoop *el = server->config.eventLoop; |
373 | 0 | el->removeDelayedCallback(el, &am->dc); |
374 | 0 | } |
375 | 536 | } |
376 | | |
377 | | void |
378 | 19.4k | UA_AsyncManager_clear(UA_AsyncManager *am, UA_Server *server) { |
379 | 19.4k | UA_LOCK_ASSERT(&server->serviceMutex); |
380 | | |
381 | | /* Cancel all operations. This moves all operations and responses into the |
382 | | * ready state. */ |
383 | 19.4k | UA_AsyncOperation *op, *op_tmp; |
384 | 19.4k | TAILQ_FOREACH_SAFE(op, &am->waitingOps, pointers, op_tmp) { |
385 | 0 | UA_AsyncOperation_cancel(server, op, UA_STATUSCODE_BADSHUTDOWN); |
386 | 0 | processOperationResult(server, op); |
387 | 0 | } |
388 | | |
389 | | /* This sends out/notifies and removes all direct operations and async requests */ |
390 | 19.4k | UA_AsyncManager_processReady(server, am); |
391 | 19.4k | UA_assert(am->opsCount == 0); |
392 | 19.4k | } |
393 | | |
394 | | UA_UInt32 |
395 | 28 | UA_AsyncManager_cancel(UA_Server *server, UA_Session *session, UA_UInt32 requestHandle) { |
396 | 28 | UA_LOCK_ASSERT(&server->serviceMutex); |
397 | | |
398 | | /* Loop over all waiting operations */ |
399 | 28 | UA_UInt32 count = 0; |
400 | 28 | UA_AsyncOperation *op, *op_tmp; |
401 | 28 | UA_AsyncManager *am = &server->asyncManager; |
402 | 28 | TAILQ_FOREACH_SAFE(op, &am->waitingOps, pointers, op_tmp) { |
403 | | /* Only request operations own a handling.response. */ |
404 | 0 | if(op->asyncOperationType >= UA_ASYNCOPERATIONTYPE_CALL_DIRECT) |
405 | 0 | continue; |
406 | 0 | UA_AsyncResponse *ar = op->handling.response; |
407 | 0 | if(ar->requestHandle != requestHandle || |
408 | 0 | !UA_NodeId_equal(&session->sessionId, &ar->sessionId)) |
409 | 0 | continue; |
410 | | |
411 | 0 | count++; /* Found a matching request */ |
412 | | |
413 | | /* Set the status of the overall response */ |
414 | 0 | ar->response.callResponse.responseHeader.serviceResult = |
415 | 0 | UA_STATUSCODE_BADREQUESTCANCELLEDBYCLIENT; |
416 | | |
417 | | /* Notify, set operation status and integrate */ |
418 | 0 | UA_AsyncOperation_cancel(server, op, UA_STATUSCODE_BADOPERATIONABANDONED); |
419 | 0 | processOperationResult(server, op); |
420 | 0 | } |
421 | | |
422 | 28 | return count; |
423 | 28 | } |
424 | | |
425 | | static void |
426 | | persistAsyncResponse(UA_Server *server, UA_Session *session, |
427 | 0 | void *response, UA_AsyncResponse *ar) { |
428 | 0 | UA_LOCK_ASSERT(&server->serviceMutex); |
429 | 0 | UA_AsyncManager *am = &server->asyncManager; |
430 | | |
431 | | /* Pending results, attach the AsyncResponse to the AsyncManager. RequestId |
432 | | * and -Handle are set in the AsyncManager before processing the request. */ |
433 | 0 | ar->requestId = am->currentRequestId; |
434 | 0 | ar->requestHandle = am->currentRequestHandle; |
435 | 0 | ar->sessionId = session->sessionId; |
436 | 0 | ar->timeout = UA_INT64_MAX; |
437 | |
|
438 | 0 | UA_EventLoop *el = server->config.eventLoop; |
439 | 0 | if(server->config.asyncOperationTimeout > 0.0) |
440 | 0 | ar->timeout = el->dateTime_nowMonotonic(el) + (UA_DateTime) |
441 | 0 | (server->config.asyncOperationTimeout * (UA_DateTime)UA_DATETIME_MSEC); |
442 | | |
443 | | /* Move the response content to the AsyncResponse */ |
444 | 0 | memcpy(&ar->response, response, ar->responseType->memSize); |
445 | 0 | UA_init(response, ar->responseType); |
446 | | |
447 | | /* Enqueue the ar */ |
448 | 0 | TAILQ_INSERT_TAIL(&am->waitingResponses, ar, pointers); |
449 | 0 | } |
450 | | |
451 | | static void |
452 | | persistAsyncResponseOperation(UA_Server *server, UA_AsyncOperation *op, |
453 | | UA_AsyncOperationType opType, UA_AsyncResponse *ar, |
454 | 0 | void *outputPtr) { |
455 | | /* Set up the async operation */ |
456 | 0 | op->asyncOperationType = opType; |
457 | 0 | op->handling.response = ar; |
458 | 0 | op->output.read = (UA_DataValue*)outputPtr; |
459 | | |
460 | | /* Not enough resources to store the async operation */ |
461 | 0 | UA_AsyncManager *am = &server->asyncManager; |
462 | 0 | if(server->config.maxAsyncOperationQueueSize != 0 && |
463 | 0 | am->opsCount >= server->config.maxAsyncOperationQueueSize) { |
464 | 0 | UA_LOG_WARNING(server->config.logging, UA_LOGCATEGORY_SERVER, |
465 | 0 | "Cannot create async operation: Queue exceeds limit (%d).", |
466 | 0 | (int unsigned)server->config.maxAsyncOperationQueueSize); |
467 | | /* No need to call processOperationResult or UA_AsyncOperation_delete |
468 | | * here. The response already has the status code integrated. */ |
469 | 0 | UA_AsyncOperation_cancel(server, op, UA_STATUSCODE_BADTOOMANYOPERATIONS); |
470 | 0 | return; |
471 | 0 | } |
472 | | |
473 | | /* Enqueue the asyncop in the async manager */ |
474 | 0 | TAILQ_INSERT_TAIL(&am->waitingOps, op, pointers); |
475 | 0 | ar->opCountdown++; |
476 | 0 | am->opsCount++; |
477 | 0 | } |
478 | | |
479 | | static UA_StatusCode |
480 | | persistAsyncDirectOperation(UA_Server *server, UA_AsyncOperation *op, |
481 | | UA_AsyncOperationType opType, void *context, |
482 | 0 | uintptr_t callback, UA_DateTime timeout) { |
483 | | /* Set up the async operation */ |
484 | 0 | op->asyncOperationType = opType; |
485 | 0 | op->handling.callback.timeout = timeout; |
486 | 0 | op->handling.callback.context = context; |
487 | 0 | op->handling.callback.method.read = (UA_ServerAsyncReadResultCallback)callback; |
488 | | |
489 | | /* Not enough resources to store the async operation */ |
490 | 0 | UA_AsyncManager *am = &server->asyncManager; |
491 | 0 | if(server->config.maxAsyncOperationQueueSize != 0 && |
492 | 0 | am->opsCount >= server->config.maxAsyncOperationQueueSize) { |
493 | 0 | UA_LOG_WARNING(server->config.logging, UA_LOGCATEGORY_SERVER, |
494 | 0 | "Cannot create async operation: Queue exceeds limit (%d).", |
495 | 0 | (int unsigned)server->config.maxAsyncOperationQueueSize); |
496 | 0 | UA_AsyncOperation_cancel(server, op, UA_STATUSCODE_BADTOOMANYOPERATIONS); |
497 | 0 | UA_AsyncOperation_delete(op); |
498 | 0 | return UA_STATUSCODE_BADTOOMANYOPERATIONS; |
499 | 0 | } |
500 | | |
501 | | /* Enqueue the asyncop in the async manager */ |
502 | 0 | TAILQ_INSERT_TAIL(&am->waitingOps, op, pointers); |
503 | 0 | am->opsCount++; |
504 | 0 | return UA_STATUSCODE_GOOD; |
505 | 0 | } |
506 | | |
507 | | void |
508 | | async_cancel(UA_Server *server, void *context, UA_StatusCode opstatus, |
509 | 0 | UA_Boolean cancelSynchronous) { |
510 | 0 | UA_AsyncManager *am = &server->asyncManager; |
511 | 0 | UA_AsyncOperation *op = NULL, *op_tmp = NULL; |
512 | | |
513 | | /* Cancel operations that are still waiting for the result */ |
514 | 0 | TAILQ_FOREACH_SAFE(op, &am->waitingOps, pointers, op_tmp) { |
515 | | /* Only direct operations own a handling.callback. */ |
516 | 0 | if(op->asyncOperationType < UA_ASYNCOPERATIONTYPE_CALL_DIRECT) |
517 | 0 | continue; |
518 | 0 | if(op->handling.callback.context != context) |
519 | 0 | continue; |
520 | | |
521 | | /* Cancel the operation. This sets the StatusCode and calls the |
522 | | * asyncOperationCancelCallback. */ |
523 | 0 | UA_AsyncOperation_cancel(server, op, opstatus); |
524 | | |
525 | | /* Call the result-callback of the local async operation. |
526 | | * Right away or in the next EventLoop iteration. */ |
527 | 0 | if(cancelSynchronous) { |
528 | 0 | TAILQ_REMOVE(&am->waitingOps, op, pointers); |
529 | 0 | am->opsCount--; |
530 | 0 | directOpCallback(server, op); |
531 | 0 | UA_AsyncOperation_delete(op); |
532 | 0 | } else { |
533 | 0 | processOperationResult(server, op); |
534 | 0 | } |
535 | 0 | } |
536 | | |
537 | | /* All "ready" operations get processed in the next EventLoop iteration anyway */ |
538 | 0 | if(!cancelSynchronous) |
539 | 0 | return; |
540 | | |
541 | | /* Process matching ready operations synchronously and delete them */ |
542 | 0 | TAILQ_FOREACH_SAFE(op, &am->readyOps, pointers, op_tmp) { |
543 | 0 | if(op->handling.callback.context != context) |
544 | 0 | continue; |
545 | 0 | TAILQ_REMOVE(&am->readyOps, op, pointers); |
546 | 0 | am->opsCount--; |
547 | 0 | directOpCallback(server, op); |
548 | 0 | UA_AsyncOperation_delete(op); |
549 | 0 | } |
550 | 0 | } |
551 | | |
552 | | void |
553 | | UA_Server_cancelAsync(UA_Server *server, void *context, UA_StatusCode opstatus, |
554 | 0 | UA_Boolean synchronousResultCallback) { |
555 | 0 | lockServer(server); |
556 | 0 | async_cancel(server, context, opstatus, synchronousResultCallback); |
557 | 0 | unlockServer(server); |
558 | 0 | } |
559 | | |
560 | | /********/ |
561 | | /* Read */ |
562 | | /********/ |
563 | | |
564 | | UA_Boolean |
565 | | Service_Read(UA_Server *server, UA_Session *session, const UA_ReadRequest *request, |
566 | 310 | UA_ReadResponse *response) { |
567 | 310 | UA_LOG_DEBUG_SESSION(server->config.logging, session, "Processing ReadRequest"); |
568 | 310 | UA_LOCK_ASSERT(&server->serviceMutex); |
569 | | |
570 | | /* Check if the timestampstoreturn is valid */ |
571 | 310 | if(request->timestampsToReturn > UA_TIMESTAMPSTORETURN_NEITHER) { |
572 | 62 | response->responseHeader.serviceResult = UA_STATUSCODE_BADTIMESTAMPSTORETURNINVALID; |
573 | 62 | return true; |
574 | 62 | } |
575 | | |
576 | | /* Check if maxAge is valid */ |
577 | 248 | if(request->maxAge < 0) { |
578 | 13 | response->responseHeader.serviceResult = UA_STATUSCODE_BADMAXAGEINVALID; |
579 | 13 | return true; |
580 | 13 | } |
581 | | |
582 | | /* Check if there are too many operations */ |
583 | 235 | if(server->config.maxNodesPerRead != 0 && |
584 | 0 | request->nodesToReadSize > server->config.maxNodesPerRead) { |
585 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADTOOMANYOPERATIONS; |
586 | 0 | return true; |
587 | 0 | } |
588 | | |
589 | | /* Check if there are no operations */ |
590 | 235 | if(request->nodesToReadSize == 0) { |
591 | 19 | response->responseHeader.serviceResult = UA_STATUSCODE_BADNOTHINGTODO; |
592 | 19 | return true; |
593 | 19 | } |
594 | | |
595 | | /* Allocate the results array */ |
596 | 216 | UA_AsyncResponse *ar = NULL; |
597 | 216 | UA_AsyncOperation *aopArray = NULL; |
598 | 216 | response->results = (UA_DataValue*) |
599 | 216 | allocateResultsArray(&UA_TYPES[UA_TYPES_DATAVALUE], |
600 | 216 | request->nodesToReadSize, &ar, &aopArray); |
601 | 216 | if(!response->results) { |
602 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADOUTOFMEMORY; |
603 | 0 | return true; |
604 | 0 | } |
605 | 216 | response->resultsSize = request->nodesToReadSize; |
606 | | |
607 | | /* Execute the operations */ |
608 | 500 | for(size_t i = 0; i < request->nodesToReadSize; i++) { |
609 | 284 | UA_Boolean done = Operation_Read(server, session, request->timestampsToReturn, |
610 | 284 | &request->nodesToRead[i], &response->results[i]); |
611 | 284 | if(!done) |
612 | 0 | persistAsyncResponseOperation(server, &aopArray[i], |
613 | 0 | UA_ASYNCOPERATIONTYPE_READ_REQUEST, |
614 | 0 | ar, &response->results[i]); |
615 | 284 | } |
616 | | |
617 | | /* If async operations are pending, persist them and signal the service is |
618 | | * not done */ |
619 | 216 | if(ar->opCountdown > 0) { |
620 | 0 | ar->responseType = &UA_TYPES[UA_TYPES_READRESPONSE]; |
621 | 0 | persistAsyncResponse(server, session, response, ar); |
622 | 0 | } |
623 | 216 | return (ar->opCountdown == 0); |
624 | 216 | } |
625 | | |
626 | | UA_StatusCode |
627 | | read_async(UA_Server *server, UA_Session *session, const UA_ReadValueId *operation, |
628 | | UA_TimestampsToReturn ttr, UA_ServerAsyncReadResultCallback callback, |
629 | 0 | void *context, UA_UInt32 timeout) { |
630 | | /* Allocate the async operation. Do this first as we need the pointer to the |
631 | | * datavalue to be stable.*/ |
632 | 0 | UA_AsyncOperation *op = (UA_AsyncOperation*)UA_calloc(1, sizeof(UA_AsyncOperation)); |
633 | 0 | if(!op) |
634 | 0 | return UA_STATUSCODE_BADOUTOFMEMORY; |
635 | | |
636 | 0 | UA_AsyncManager *am = &server->asyncManager; |
637 | 0 | if(server->config.maxAsyncOperationQueueSize != 0 && |
638 | 0 | am->opsCount >= server->config.maxAsyncOperationQueueSize) { |
639 | 0 | UA_free(op); |
640 | 0 | return UA_STATUSCODE_BADTOOMANYOPERATIONS; |
641 | 0 | } |
642 | | |
643 | 0 | UA_DateTime timeoutDate = UA_INT64_MAX; |
644 | 0 | if(timeout > 0) { |
645 | 0 | UA_EventLoop *el = server->config.eventLoop; |
646 | 0 | const UA_DateTime tNow = el->dateTime_nowMonotonic(el); |
647 | 0 | timeoutDate = tNow + (timeout * UA_DATETIME_MSEC); |
648 | 0 | } |
649 | | |
650 | | /* Call the operation */ |
651 | 0 | UA_Boolean done = Operation_Read(server, session, ttr, operation, &op->output.directRead); |
652 | 0 | if(!done) |
653 | 0 | return persistAsyncDirectOperation(server, op, UA_ASYNCOPERATIONTYPE_READ_DIRECT, |
654 | 0 | context, (uintptr_t)callback, timeoutDate); |
655 | | |
656 | 0 | callback(server, context, &op->output.directRead); |
657 | 0 | UA_DataValue_clear(&op->output.directRead); |
658 | 0 | UA_free(op); |
659 | 0 | return UA_STATUSCODE_GOOD; |
660 | 0 | } |
661 | | |
662 | | UA_StatusCode |
663 | | UA_Server_read_async(UA_Server *server, const UA_ReadValueId *operation, |
664 | | UA_TimestampsToReturn ttr, UA_ServerAsyncReadResultCallback callback, |
665 | 0 | void *context, UA_UInt32 timeout) { |
666 | 0 | lockServer(server); |
667 | 0 | UA_StatusCode res = read_async(server, &server->adminSession, operation, |
668 | 0 | ttr, callback, context, timeout); |
669 | 0 | unlockServer(server); |
670 | 0 | return res; |
671 | 0 | } |
672 | | |
673 | | UA_StatusCode |
674 | 0 | UA_Server_setAsyncReadResult(UA_Server *server, UA_DataValue *result) { |
675 | 0 | lockServer(server); |
676 | 0 | UA_AsyncManager *am = &server->asyncManager; |
677 | 0 | UA_AsyncOperation *op = NULL; |
678 | 0 | TAILQ_FOREACH(op, &am->waitingOps, pointers) { |
679 | 0 | if(op->output.read == result || &op->output.directRead == result) { |
680 | 0 | processOperationResult(server, op); |
681 | 0 | break; |
682 | 0 | } |
683 | 0 | } |
684 | 0 | unlockServer(server); |
685 | 0 | return (op) ? UA_STATUSCODE_GOOD : UA_STATUSCODE_BADNOTFOUND; |
686 | 0 | } |
687 | | |
688 | | /*********/ |
689 | | /* Write */ |
690 | | /*********/ |
691 | | |
692 | | UA_Boolean |
693 | | Service_Write(UA_Server *server, UA_Session *session, |
694 | 1 | const UA_WriteRequest *request, UA_WriteResponse *response) { |
695 | 1 | UA_assert(session != NULL); |
696 | 1 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
697 | 1 | "Processing WriteRequest"); |
698 | 1 | UA_LOCK_ASSERT(&server->serviceMutex); |
699 | | |
700 | 1 | if(server->config.maxNodesPerWrite != 0 && |
701 | 0 | request->nodesToWriteSize > server->config.maxNodesPerWrite) { |
702 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADTOOMANYOPERATIONS; |
703 | 0 | return true; |
704 | 0 | } |
705 | | |
706 | 1 | if(request->nodesToWriteSize == 0) { |
707 | 1 | response->responseHeader.serviceResult = UA_STATUSCODE_BADNOTHINGTODO; |
708 | 1 | return true; |
709 | 1 | } |
710 | | |
711 | | /* Allocate the results array */ |
712 | 0 | UA_AsyncResponse *ar = NULL; |
713 | 0 | UA_AsyncOperation *aopArray = NULL; |
714 | 0 | response->results = (UA_StatusCode*) |
715 | 0 | allocateResultsArray(&UA_TYPES[UA_TYPES_STATUSCODE], |
716 | 0 | request->nodesToWriteSize, &ar, &aopArray); |
717 | 0 | if(!response->results) { |
718 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADOUTOFMEMORY; |
719 | 0 | return true; |
720 | 0 | } |
721 | 0 | response->resultsSize = request->nodesToWriteSize; |
722 | | |
723 | | /* Execute the operations */ |
724 | 0 | for(size_t i = 0; i < request->nodesToWriteSize; i++) { |
725 | | /* Ensure a stable pointer for the writevalue. Doesn't get written to, |
726 | | * just used for the lookup of the async operation later on. |
727 | | * The original writeValue might be _clear'ed before the lookup. */ |
728 | 0 | UA_AsyncOperation *aop = &aopArray[i]; |
729 | 0 | aop->context.writeValue = request->nodesToWrite[i]; |
730 | 0 | UA_Boolean done = Operation_Write(server, session, &aop->context.writeValue, |
731 | 0 | &response->results[i]); |
732 | 0 | if(!done) |
733 | 0 | persistAsyncResponseOperation(server, aop, UA_ASYNCOPERATIONTYPE_WRITE_REQUEST, |
734 | 0 | ar, &response->results[i]); |
735 | 0 | } |
736 | | |
737 | | /* If async operations are pending, persist them and signal the service is |
738 | | * not done */ |
739 | 0 | if(ar->opCountdown > 0) { |
740 | 0 | ar->responseType = &UA_TYPES[UA_TYPES_WRITERESPONSE]; |
741 | 0 | persistAsyncResponse(server, session, response, ar); |
742 | 0 | } |
743 | 0 | return (ar->opCountdown == 0); |
744 | 0 | } |
745 | | |
746 | | UA_StatusCode |
747 | | write_async(UA_Server *server, UA_Session *session, const UA_WriteValue *operation, |
748 | | UA_ServerAsyncWriteResultCallback callback, void *context, |
749 | 0 | UA_UInt32 timeout) { |
750 | | /* Allocate the async operation. Do this first as we need the pointer to the |
751 | | * datavalue to be stable.*/ |
752 | 0 | UA_AsyncOperation *op = (UA_AsyncOperation*)UA_calloc(1, sizeof(UA_AsyncOperation)); |
753 | 0 | if(!op) |
754 | 0 | return UA_STATUSCODE_BADOUTOFMEMORY; |
755 | | |
756 | 0 | UA_AsyncManager *am = &server->asyncManager; |
757 | 0 | if(server->config.maxAsyncOperationQueueSize != 0 && |
758 | 0 | am->opsCount >= server->config.maxAsyncOperationQueueSize) { |
759 | 0 | UA_free(op); |
760 | 0 | return UA_STATUSCODE_BADTOOMANYOPERATIONS; |
761 | 0 | } |
762 | | |
763 | 0 | UA_DateTime timeoutDate = UA_INT64_MAX; |
764 | 0 | if(timeout > 0) { |
765 | 0 | UA_EventLoop *el = server->config.eventLoop; |
766 | 0 | const UA_DateTime tNow = el->dateTime_nowMonotonic(el); |
767 | 0 | timeoutDate = tNow + (timeout * UA_DATETIME_MSEC); |
768 | 0 | } |
769 | | |
770 | | /* Call the operation */ |
771 | 0 | op->context.writeValue = *operation; /* Stable pointer */ |
772 | 0 | UA_Boolean done = Operation_Write(server, session, &op->context.writeValue, |
773 | 0 | &op->output.directWrite); |
774 | 0 | if(!done) |
775 | 0 | return persistAsyncDirectOperation(server, op, UA_ASYNCOPERATIONTYPE_WRITE_DIRECT, |
776 | 0 | context, (uintptr_t)callback, timeoutDate); |
777 | | |
778 | | /* Done, return right away */ |
779 | 0 | callback(server, context, op->output.directWrite); |
780 | 0 | UA_free(op); |
781 | 0 | return UA_STATUSCODE_GOOD; |
782 | 0 | } |
783 | | |
784 | | UA_StatusCode |
785 | | UA_Server_write_async(UA_Server *server, const UA_WriteValue *operation, |
786 | | UA_ServerAsyncWriteResultCallback callback, |
787 | 0 | void *context, UA_UInt32 timeout) { |
788 | 0 | lockServer(server); |
789 | 0 | UA_StatusCode res = write_async(server, &server->adminSession, operation, |
790 | 0 | callback, context, timeout); |
791 | 0 | unlockServer(server); |
792 | 0 | return res; |
793 | 0 | } |
794 | | |
795 | | UA_StatusCode |
796 | | UA_Server_setAsyncWriteResult(UA_Server *server, |
797 | | const UA_DataValue *value, |
798 | 0 | UA_StatusCode result) { |
799 | 0 | lockServer(server); |
800 | 0 | UA_AsyncManager *am = &server->asyncManager; |
801 | 0 | UA_AsyncOperation *op = NULL; |
802 | 0 | TAILQ_FOREACH(op, &am->waitingOps, pointers) { |
803 | 0 | if(&op->context.writeValue.value == value) { |
804 | 0 | if(op->asyncOperationType == UA_ASYNCOPERATIONTYPE_WRITE_REQUEST) |
805 | 0 | *op->output.write = result; |
806 | 0 | else |
807 | 0 | op->output.directWrite = result; |
808 | 0 | processOperationResult(server, op); |
809 | 0 | break; |
810 | 0 | } |
811 | 0 | } |
812 | 0 | unlockServer(server); |
813 | 0 | return (op) ? UA_STATUSCODE_GOOD : UA_STATUSCODE_BADNOTFOUND; |
814 | 0 | } |
815 | | |
816 | | /********/ |
817 | | /* Call */ |
818 | | /********/ |
819 | | |
820 | | #ifdef UA_ENABLE_METHODCALLS |
821 | | UA_Boolean |
822 | | Service_Call(UA_Server *server, UA_Session *session, |
823 | 48 | const UA_CallRequest *request, UA_CallResponse *response) { |
824 | 48 | UA_LOG_DEBUG_SESSION(server->config.logging, session, "Processing CallRequest"); |
825 | 48 | UA_LOCK_ASSERT(&server->serviceMutex); |
826 | | |
827 | 48 | if(server->config.maxNodesPerMethodCall != 0 && |
828 | 0 | request->methodsToCallSize > server->config.maxNodesPerMethodCall) { |
829 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADTOOMANYOPERATIONS; |
830 | 0 | return true; |
831 | 0 | } |
832 | | |
833 | 48 | if(request->methodsToCallSize == 0) { |
834 | 11 | response->responseHeader.serviceResult = UA_STATUSCODE_BADNOTHINGTODO; |
835 | 11 | return true; |
836 | 11 | } |
837 | | |
838 | | /* Allocate the results array */ |
839 | 37 | UA_AsyncResponse *ar = NULL; |
840 | 37 | UA_AsyncOperation *aopArray = NULL; |
841 | 37 | response->results = (UA_CallMethodResult*) |
842 | 37 | allocateResultsArray(&UA_TYPES[UA_TYPES_CALLMETHODRESULT], |
843 | 37 | request->methodsToCallSize, &ar, &aopArray); |
844 | 37 | if(!response->results) { |
845 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADOUTOFMEMORY; |
846 | 0 | return true; |
847 | 0 | } |
848 | 37 | response->resultsSize = request->methodsToCallSize; |
849 | | |
850 | | /* Execute the operations */ |
851 | 158 | for(size_t i = 0; i < request->methodsToCallSize; i++) { |
852 | 121 | UA_Boolean done = Operation_CallMethod(server, session, &request->methodsToCall[i], |
853 | 121 | &response->results[i]); |
854 | 121 | if(!done) |
855 | 0 | persistAsyncResponseOperation(server, &aopArray[i], |
856 | 0 | UA_ASYNCOPERATIONTYPE_CALL_REQUEST, |
857 | 0 | ar, &response->results[i]); |
858 | 121 | } |
859 | | |
860 | | /* If async operations are pending, persist them and signal the service is |
861 | | * not done */ |
862 | 37 | if(ar->opCountdown > 0) { |
863 | 0 | ar->responseType = &UA_TYPES[UA_TYPES_CALLRESPONSE]; |
864 | 0 | persistAsyncResponse(server, session, response, ar); |
865 | 0 | } |
866 | 37 | return (ar->opCountdown == 0); |
867 | 37 | } |
868 | | |
869 | | UA_StatusCode |
870 | | call_async(UA_Server *server, UA_Session *session, const UA_CallMethodRequest *operation, |
871 | | UA_ServerAsyncMethodResultCallback callback, void *context, |
872 | 0 | UA_UInt32 timeout) { |
873 | | /* Allocate the async operation. Do this first as we need the pointer to the |
874 | | * datavalue to be stable.*/ |
875 | 0 | UA_AsyncOperation *op = (UA_AsyncOperation*)UA_calloc(1, sizeof(UA_AsyncOperation)); |
876 | 0 | if(!op) |
877 | 0 | return UA_STATUSCODE_BADOUTOFMEMORY; |
878 | | |
879 | 0 | UA_AsyncManager *am = &server->asyncManager; |
880 | 0 | if(server->config.maxAsyncOperationQueueSize != 0 && |
881 | 0 | am->opsCount >= server->config.maxAsyncOperationQueueSize) { |
882 | 0 | UA_free(op); |
883 | 0 | return UA_STATUSCODE_BADTOOMANYOPERATIONS; |
884 | 0 | } |
885 | | |
886 | 0 | UA_DateTime timeoutDate = UA_INT64_MAX; |
887 | 0 | if(timeout > 0) { |
888 | 0 | UA_EventLoop *el = server->config.eventLoop; |
889 | 0 | const UA_DateTime tNow = el->dateTime_nowMonotonic(el); |
890 | 0 | timeoutDate = tNow + (timeout * UA_DATETIME_MSEC); |
891 | 0 | } |
892 | | |
893 | | /* Call the operation */ |
894 | 0 | UA_Boolean done = Operation_CallMethod(server, session, operation, |
895 | 0 | &op->output.directCall); |
896 | 0 | if(!done) |
897 | 0 | return persistAsyncDirectOperation(server, op, UA_ASYNCOPERATIONTYPE_CALL_DIRECT, |
898 | 0 | context, (uintptr_t)callback, timeoutDate); |
899 | | |
900 | | /* Done, return right away */ |
901 | 0 | callback(server, context, &op->output.directCall); |
902 | 0 | UA_CallMethodResult_clear(&op->output.directCall); |
903 | 0 | UA_free(op); |
904 | 0 | return UA_STATUSCODE_GOOD; |
905 | 0 | } |
906 | | |
907 | | UA_StatusCode |
908 | | UA_Server_call_async(UA_Server *server, const UA_CallMethodRequest *operation, |
909 | | UA_ServerAsyncMethodResultCallback callback, |
910 | 0 | void *context, UA_UInt32 timeout) { |
911 | 0 | lockServer(server); |
912 | 0 | UA_StatusCode res = |
913 | 0 | call_async(server, &server->adminSession, operation, callback, context, timeout); |
914 | 0 | unlockServer(server); |
915 | 0 | return res; |
916 | 0 | } |
917 | | |
918 | | UA_StatusCode |
919 | | UA_Server_setAsyncCallMethodResult(UA_Server *server, UA_Variant *output, |
920 | 0 | UA_StatusCode result) { |
921 | 0 | lockServer(server); |
922 | 0 | UA_AsyncManager *am = &server->asyncManager; |
923 | 0 | UA_AsyncOperation *op = NULL; |
924 | 0 | TAILQ_FOREACH(op, &am->waitingOps, pointers) { |
925 | 0 | if(op->asyncOperationType == UA_ASYNCOPERATIONTYPE_CALL_REQUEST) { |
926 | 0 | if(op->output.call->outputArguments == output) { |
927 | 0 | op->output.call->statusCode = result; |
928 | 0 | processOperationResult(server, op); |
929 | 0 | break; |
930 | 0 | } |
931 | 0 | } else if(op->asyncOperationType == UA_ASYNCOPERATIONTYPE_CALL_DIRECT) { |
932 | 0 | if(op->output.directCall.outputArguments == output) { |
933 | 0 | op->output.directCall.statusCode = result; |
934 | 0 | processOperationResult(server, op); |
935 | 0 | break; |
936 | 0 | } |
937 | 0 | } |
938 | 0 | } |
939 | 0 | unlockServer(server); |
940 | 0 | return (op) ? UA_STATUSCODE_GOOD : UA_STATUSCODE_BADNOTFOUND; |
941 | 0 | } |
942 | | #endif |