Coverage Report

Created: 2026-09-27 07:10

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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(&notifyPayload[0].value, &secureChannelId,
169
0
                             &UA_TYPES[UA_TYPES_UINT32]);
170
0
        UA_Variant_setScalar(&notifyPayload[1].value, &sessionId,
171
0
                             &UA_TYPES[UA_TYPES_NODEID]);
172
0
        UA_Variant_setScalar(&notifyPayload[2].value, &ar->requestId,
173
0
                             &UA_TYPES[UA_TYPES_UINT32]);
174
0
        UA_Variant_setScalar(&notifyPayload[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