Coverage Report

Created: 2026-09-28 07:11

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/open62541_15/src/client/ua_client_subscriptions.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 2015-2018 (c) Fraunhofer IOSB (Author: Julius Pfrommer)
6
 *    Copyright 2015 (c) Oleksiy Vasylyev
7
 *    Copyright 2016 (c) Sten Grüner
8
 *    Copyright 2017-2018 (c) Thomas Stalder, Blue Time Concept SA
9
 *    Copyright 2016-2017 (c) Florian Palm
10
 *    Copyright 2017 (c) Frank Meerkötter
11
 *    Copyright 2017 (c) Stefan Profanter, fortiss GmbH
12
 *    Copyright 2026 (c) o6 Automation GmbH (Author: Julius Pfrommer)
13
 */
14
15
#include <open62541/client_highlevel.h>
16
#include <open62541/client_highlevel_async.h>
17
18
#include "ua_client_internal.h"
19
20
struct UA_Client_MonitoredItem_ForDelete {
21
    UA_Client *client;
22
    UA_Client_Subscription *sub;
23
    UA_UInt32 *monitoredItemId;
24
};
25
26
/*****************/
27
/* Subscriptions */
28
/*****************/
29
30
/* For ZIP_TREE we use clientHandle comparison */
31
static enum ZIP_CMP
32
0
UA_ClientHandle_cmp(const void *a, const void *b) {
33
0
    const UA_Client_MonitoredItem *aa = (const UA_Client_MonitoredItem *)a;
34
0
    const UA_Client_MonitoredItem *bb = (const UA_Client_MonitoredItem *)b;
35
0
    if(aa->parameters.clientHandle < bb->parameters.clientHandle)
36
0
        return ZIP_CMP_LESS;
37
0
    if(aa->parameters.clientHandle > bb->parameters.clientHandle)
38
0
        return ZIP_CMP_MORE;
39
0
    return ZIP_CMP_EQ;
40
0
}
41
42
0
ZIP_FUNCTIONS(MonitorItemsTree, UA_Client_MonitoredItem, zipfields,
Unexecuted instantiation: ua_client_subscriptions.c:MonitorItemsTree_ZIP_INSERT
Unexecuted instantiation: ua_client_subscriptions.c:MonitorItemsTree_ZIP_REMOVE
Unexecuted instantiation: ua_client_subscriptions.c:MonitorItemsTree_ZIP_ITER
43
              UA_Client_MonitoredItem, zipfields, UA_ClientHandle_cmp)
44
45
static UA_Client_MonitoredItem *
46
0
findMonitoredItemByHandle(UA_Client_Subscription *sub, UA_UInt32 clientHandle) {
47
0
    UA_Client_MonitoredItem dummy;
48
0
    memset(&dummy, 0, sizeof(dummy));
49
0
    dummy.parameters.clientHandle = clientHandle;
50
0
    return ZIP_FIND(MonitorItemsTree, &sub->monitoredItems, &dummy);
51
0
}
52
53
static void *
54
0
MonitoredItem_findPendingByHandle(void *data, UA_Client_MonitoredItem *mon) {
55
0
    UA_UInt32 clientHandle = *(UA_UInt32*)data;
56
0
    if(mon->pendingParameters.clientHandle == clientHandle)
57
0
        return mon;
58
0
    return NULL;
59
0
}
60
61
static void
62
MonitoredItem_delete(UA_Client *client, UA_Client_Subscription *sub,
63
                     UA_Client_MonitoredItem *mon);
64
65
static void
66
Subscription_create(UA_Client *client, UA_Client_Subscription *newSub,
67
0
                    UA_CreateSubscriptionResponse *response) {
68
0
    UA_LOCK_ASSERT(&client->clientMutex);
69
70
0
    UA_EventLoop *el = client->config.eventLoop;
71
72
0
    newSub->subscriptionId = response->subscriptionId;
73
0
    newSub->sequenceNumber = 0;
74
0
    newSub->lastActivity = el->dateTime_nowMonotonic(el);
75
0
    newSub->publishingInterval = response->revisedPublishingInterval;
76
0
    newSub->maxKeepAliveCount = response->revisedMaxKeepAliveCount;
77
0
    ZIP_INIT(&newSub->monitoredItems);
78
0
    newSub->pendingRekeys = 0;
79
0
    LIST_INSERT_HEAD(&client->subscriptions, newSub, listEntry);
80
81
    /* Immediately send the first publish requests if there are none
82
     * outstanding */
83
0
    __Client_Subscriptions_backgroundPublish(client);
84
0
}
85
86
static void
87
Subscriptions_create_handler(UA_Client *client, void *data,
88
0
                             UA_UInt32 requestId, void *r) {
89
0
    UA_LOCK_ASSERT(&client->clientMutex);
90
91
0
    UA_CreateSubscriptionResponse *response = (UA_CreateSubscriptionResponse *)r;
92
0
    CustomCallback *cc = (CustomCallback *)data;
93
0
    UA_Client_Subscription *newSub = (UA_Client_Subscription *)cc->clientData;
94
0
    if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD) {
95
0
        UA_free(newSub);
96
0
        goto cleanup;
97
0
    }
98
99
    /* Prepare the internal representation */
100
0
    Subscription_create(client, newSub, response);
101
102
0
cleanup:
103
0
    if(cc->userCallback)
104
0
        cc->userCallback(client, cc->userData, requestId, response);
105
0
    UA_free(cc);
106
0
}
107
108
UA_CreateSubscriptionResponse
109
UA_Client_Subscriptions_create(UA_Client *client,
110
                               const UA_CreateSubscriptionRequest request,
111
                               void *subscriptionContext,
112
                               UA_Client_StatusChangeNotificationCallback statusChangeCallback,
113
0
                               UA_Client_DeleteSubscriptionCallback deleteCallback) {
114
0
    lockClient(client);
115
116
0
    UA_CreateSubscriptionResponse response;
117
0
    UA_Client_Subscription *sub = (UA_Client_Subscription *)
118
0
        UA_malloc(sizeof(UA_Client_Subscription));
119
0
    if(!sub) {
120
0
        UA_CreateSubscriptionResponse_init(&response);
121
0
        response.responseHeader.serviceResult = UA_STATUSCODE_BADOUTOFMEMORY;
122
0
        unlockClient(client);
123
0
        return response;
124
0
    }
125
0
    sub->context = subscriptionContext;
126
0
    sub->statusChangeCallback = statusChangeCallback;
127
0
    sub->deleteCallback = deleteCallback;
128
129
    /* Send the request as a synchronous service call */
130
0
    __Client_Service(client, &request, &UA_TYPES[UA_TYPES_CREATESUBSCRIPTIONREQUEST],
131
0
                     &response, &UA_TYPES[UA_TYPES_CREATESUBSCRIPTIONRESPONSE]);
132
0
    if(response.responseHeader.serviceResult != UA_STATUSCODE_GOOD) {
133
0
        UA_free(sub);
134
0
        unlockClient(client);
135
0
        return response;
136
0
    }
137
138
0
    Subscription_create(client, sub, &response);
139
140
0
    unlockClient(client);
141
0
    return response;
142
0
}
143
144
UA_StatusCode
145
UA_Client_Subscriptions_create_async(UA_Client *client, const UA_CreateSubscriptionRequest request,
146
                                     void *subscriptionContext,
147
                                     UA_Client_StatusChangeNotificationCallback statusChangeCallback,
148
                                     UA_Client_DeleteSubscriptionCallback deleteCallback,
149
                                     UA_ClientAsyncCreateSubscriptionCallback createCallback,
150
0
                                     void *userdata, UA_UInt32 *requestId) {
151
0
    CustomCallback *cc = (CustomCallback *)UA_calloc(1, sizeof(CustomCallback));
152
0
    if(!cc)
153
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
154
155
0
    UA_Client_Subscription *sub = (UA_Client_Subscription *)
156
0
        UA_malloc(sizeof(UA_Client_Subscription));
157
0
    if(!sub) {
158
0
        UA_free(cc);
159
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
160
0
    }
161
0
    sub->context = subscriptionContext;
162
0
    sub->statusChangeCallback = statusChangeCallback;
163
0
    sub->deleteCallback = deleteCallback;
164
165
0
    cc->userCallback = (UA_ClientAsyncServiceCallback)createCallback;
166
0
    cc->userData = userdata;
167
0
    cc->clientData = sub;
168
169
    /* Send the request as asynchronous service call */
170
0
    UA_StatusCode res =
171
0
        __UA_Client_AsyncService(client, &request,
172
0
                                 &UA_TYPES[UA_TYPES_CREATESUBSCRIPTIONREQUEST],
173
0
                                 Subscriptions_create_handler,
174
0
                                 &UA_TYPES[UA_TYPES_CREATESUBSCRIPTIONRESPONSE],
175
0
                                 cc, requestId);
176
0
    if(res != UA_STATUSCODE_GOOD) {
177
0
        UA_free(cc);
178
0
        UA_free(sub);
179
0
    }
180
0
    return res;
181
0
}
182
183
static UA_Client_Subscription *
184
0
findSubscriptionById(const UA_Client *client, UA_UInt32 subscriptionId) {
185
0
    UA_Client_Subscription *sub = NULL;
186
0
    LIST_FOREACH(sub, &client->subscriptions, listEntry) {
187
0
        if(sub->subscriptionId == subscriptionId)
188
0
            break;
189
0
    }
190
0
    return sub;
191
0
}
192
193
static void
194
Subscription_modify(UA_Client *client, UA_Client_Subscription *sub,
195
0
                    const UA_ModifySubscriptionResponse *response) {
196
0
    sub->publishingInterval = response->revisedPublishingInterval;
197
0
    sub->maxKeepAliveCount = response->revisedMaxKeepAliveCount;
198
0
}
199
200
static void
201
Subscription_modify_handler(UA_Client *client, void *data,
202
0
                            UA_UInt32 requestId, void *r) {
203
0
    UA_LOCK_ASSERT(&client->clientMutex);
204
205
0
    UA_ModifySubscriptionResponse *response = (UA_ModifySubscriptionResponse *)r;
206
0
    CustomCallback *cc = (CustomCallback *)data;
207
0
    UA_Client_Subscription *sub =
208
0
        findSubscriptionById(client, (UA_UInt32)(uintptr_t)cc->clientData);
209
0
    if(sub) {
210
0
        Subscription_modify(client, sub, response);
211
0
    } else {
212
0
        UA_LOG_INFO(client->config.logging, UA_LOGCATEGORY_CLIENT,
213
0
                    "No internal representation of subscription %" PRIu32,
214
0
                    (UA_UInt32)(uintptr_t)cc->clientData);
215
0
    }
216
217
0
    if(cc->userCallback)
218
0
        cc->userCallback(client, cc->userData, requestId, response);
219
0
    UA_free(cc);
220
0
}
221
222
UA_StatusCode
223
UA_Client_Subscriptions_getContext(UA_Client *client, UA_UInt32 subscriptionId,
224
0
                                   void **subContext) {
225
0
    if(!client || !subContext)
226
0
        return UA_STATUSCODE_BADINVALIDARGUMENT;
227
228
0
    lockClient(client);
229
0
    UA_Client_Subscription *sub = findSubscriptionById(client, subscriptionId);
230
0
    if(!sub) {
231
0
        unlockClient(client);
232
0
        return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
233
0
    }
234
235
0
    *subContext = sub->context;
236
0
    unlockClient(client);
237
0
    return UA_STATUSCODE_GOOD;
238
0
}
239
240
UA_StatusCode
241
UA_Client_Subscriptions_setContext(UA_Client *client, UA_UInt32 subscriptionId,
242
0
                                   void *subContext) {
243
0
    if(!client)
244
0
        return UA_STATUSCODE_BADINVALIDARGUMENT;
245
246
0
    lockClient(client);
247
0
    UA_Client_Subscription *sub = findSubscriptionById(client, subscriptionId);
248
0
    if(!sub) {
249
0
        unlockClient(client);
250
0
        return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
251
0
    }
252
253
0
    sub->context = subContext;
254
0
    unlockClient(client);
255
0
    return UA_STATUSCODE_GOOD;
256
0
}
257
258
UA_ModifySubscriptionResponse
259
UA_Client_Subscriptions_modify(UA_Client *client,
260
0
                               const UA_ModifySubscriptionRequest request) {
261
0
    UA_ModifySubscriptionResponse response;
262
0
    UA_ModifySubscriptionResponse_init(&response);
263
264
    /* Find the internal representation */
265
0
    lockClient(client);
266
0
    UA_Client_Subscription *sub = findSubscriptionById(client, request.subscriptionId);
267
0
    if(!sub) {
268
0
        unlockClient(client);
269
0
        response.responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
270
0
        return response;
271
0
    }
272
273
    /* Call the service */
274
0
    __Client_Service(client,
275
0
                     &request, &UA_TYPES[UA_TYPES_MODIFYSUBSCRIPTIONREQUEST],
276
0
                     &response, &UA_TYPES[UA_TYPES_MODIFYSUBSCRIPTIONRESPONSE]);
277
278
    /* Adjust the internal representation. Lookup again for thread-safety. */
279
0
    sub = findSubscriptionById(client, request.subscriptionId);
280
0
    if(!sub) {
281
0
        response.responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
282
0
        unlockClient(client);
283
0
        return response;
284
0
    }
285
0
    Subscription_modify(client, sub, &response);
286
0
    unlockClient(client);
287
0
    return response;
288
0
}
289
290
UA_StatusCode
291
UA_Client_Subscriptions_modify_async(UA_Client *client,
292
                                     const UA_ModifySubscriptionRequest request,
293
                                     UA_ClientAsyncModifySubscriptionCallback callback,
294
0
                                     void *userdata, UA_UInt32 *requestId) {
295
0
    lockClient(client);
296
297
0
    UA_Client_Subscription *sub = findSubscriptionById(client, request.subscriptionId);
298
0
    if(!sub) {
299
0
        unlockClient(client);
300
0
        return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
301
0
    }
302
303
0
    CustomCallback *cc = (CustomCallback *)UA_calloc(1, sizeof(CustomCallback));
304
0
    if(!cc) {
305
0
        unlockClient(client);
306
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
307
0
    }
308
309
0
    cc->clientData = (void *)(uintptr_t)request.subscriptionId;
310
0
    cc->userData = userdata;
311
0
    cc->userCallback = (UA_ClientAsyncServiceCallback)callback;
312
313
0
    UA_StatusCode res =
314
0
        __Client_AsyncService(client, &request,
315
0
                              &UA_TYPES[UA_TYPES_MODIFYSUBSCRIPTIONREQUEST],
316
0
                              Subscription_modify_handler,
317
0
                              &UA_TYPES[UA_TYPES_MODIFYSUBSCRIPTIONRESPONSE],
318
0
                              cc, requestId);
319
320
0
    unlockClient(client);
321
0
    return res;
322
0
}
323
324
static void *
325
0
MonitoredItem_delete_wrapper(void *data, UA_Client_MonitoredItem *mon) {
326
0
    struct UA_Client_MonitoredItem_ForDelete *deleteMonitoredItem =
327
0
        (struct UA_Client_MonitoredItem_ForDelete *)data;
328
0
    if(deleteMonitoredItem != NULL) {
329
0
        if(deleteMonitoredItem->monitoredItemId != NULL &&
330
0
           (mon->monitoredItemId != *deleteMonitoredItem->monitoredItemId)) {
331
0
            return NULL;
332
0
        }
333
0
        MonitoredItem_delete(deleteMonitoredItem->client, deleteMonitoredItem->sub, mon);
334
0
    }
335
0
    return NULL;
336
0
}
337
338
static void
339
__Client_Subscription_deleteInternal(UA_Client *client,
340
0
                                     UA_Client_Subscription *sub) {
341
    /* Remove the MonitoredItems */
342
0
    struct UA_Client_MonitoredItem_ForDelete deleteMonitoredItem;
343
0
    memset(&deleteMonitoredItem, 0, sizeof(struct UA_Client_MonitoredItem_ForDelete));
344
0
    deleteMonitoredItem.client = client;
345
0
    deleteMonitoredItem.sub = sub;
346
0
    ZIP_ITER(MonitorItemsTree, &sub->monitoredItems,
347
0
             MonitoredItem_delete_wrapper, &deleteMonitoredItem);
348
349
    /* Call the delete callback */
350
0
    if(sub->deleteCallback) {
351
0
        void *subC = sub->context;
352
0
        UA_UInt32 subId = sub->subscriptionId;
353
0
        sub->deleteCallback(client, subId, subC);
354
0
    }
355
356
    /* Remove */
357
0
    LIST_REMOVE(sub, listEntry);
358
0
    UA_free(sub);
359
0
}
360
361
static void
362
__Client_Subscription_processDelete(UA_Client *client,
363
                                    const UA_DeleteSubscriptionsRequest *request,
364
0
                                    const UA_DeleteSubscriptionsResponse *response)  {
365
0
    if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD)
366
0
        return;
367
368
    /* Check that the request and response size -- use the same index for both */
369
0
    if(request->subscriptionIdsSize != response->resultsSize)
370
0
        return;
371
372
0
    for(size_t i = 0; i < request->subscriptionIdsSize; i++) {
373
0
        if(response->results[i] != UA_STATUSCODE_GOOD &&
374
0
           response->results[i] != UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID)
375
0
            continue;
376
377
        /* Get the Subscription */
378
0
        UA_Client_Subscription *sub =
379
0
            findSubscriptionById(client, request->subscriptionIds[i]);
380
0
        if(!sub) {
381
0
            UA_LOG_INFO(client->config.logging, UA_LOGCATEGORY_CLIENT,
382
0
                        "No internal representation of subscription %" PRIu32,
383
0
                        request->subscriptionIds[i]);
384
0
            continue;
385
0
        }
386
387
        /* Delete the Subscription */
388
0
        __Client_Subscription_deleteInternal(client, sub);
389
0
    }
390
0
}
391
392
typedef struct {
393
    UA_DeleteSubscriptionsRequest request;
394
    UA_ClientAsyncServiceCallback userCallback;
395
    void *userData;
396
} DeleteSubscriptionCallback;
397
398
static void
399
Subscriptions_delete_handler(UA_Client *client, void *data,
400
0
                             UA_UInt32 requestId, void *r) {
401
0
    UA_DeleteSubscriptionsResponse *response =
402
0
        (UA_DeleteSubscriptionsResponse *)r;
403
0
    DeleteSubscriptionCallback *dsc =
404
0
        (DeleteSubscriptionCallback*)data;
405
406
0
    lockClient(client);
407
408
    /* Delete */
409
0
    __Client_Subscription_processDelete(client, &dsc->request, response);
410
411
    /* Userland Callback */
412
0
    if(dsc->userCallback)
413
0
        dsc->userCallback(client, dsc->userData, requestId, response);
414
415
    /* Cleanup */
416
0
    UA_DeleteSubscriptionsRequest_clear(&dsc->request);
417
0
    UA_free(dsc);
418
419
0
    unlockClient(client);
420
0
}
421
422
UA_StatusCode
423
UA_Client_Subscriptions_delete_async(UA_Client *client,
424
                                     const UA_DeleteSubscriptionsRequest request,
425
                                     UA_ClientAsyncDeleteSubscriptionsCallback callback,
426
0
                                     void *userdata, UA_UInt32 *requestId) {
427
    /* Make a copy of the request that persists into the async callback */
428
0
    DeleteSubscriptionCallback *dsc = (DeleteSubscriptionCallback*)
429
0
        UA_malloc(sizeof(DeleteSubscriptionCallback));
430
0
    if(!dsc)
431
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
432
0
    dsc->userCallback = (UA_ClientAsyncServiceCallback)callback;
433
0
    dsc->userData = userdata;
434
0
    UA_StatusCode res = UA_DeleteSubscriptionsRequest_copy(&request, &dsc->request);
435
0
    if(res != UA_STATUSCODE_GOOD) {
436
0
        UA_free(dsc);
437
0
        return res;
438
0
    }
439
440
    /* Make the async call */
441
0
    res = __UA_Client_AsyncService(client, &request,
442
0
                                   &UA_TYPES[UA_TYPES_DELETESUBSCRIPTIONSREQUEST],
443
0
                                   Subscriptions_delete_handler,
444
0
                                   &UA_TYPES[UA_TYPES_DELETESUBSCRIPTIONSRESPONSE],
445
0
                                   dsc, requestId);
446
0
    if(res != UA_STATUSCODE_GOOD) {
447
0
        UA_DeleteSubscriptionsRequest_clear(&dsc->request);
448
0
        UA_free(dsc);
449
0
    }
450
0
    return res;
451
0
}
452
453
UA_DeleteSubscriptionsResponse
454
UA_Client_Subscriptions_delete(UA_Client *client,
455
0
                               const UA_DeleteSubscriptionsRequest request) {
456
0
    lockClient(client);
457
458
    /* Send the request */
459
0
    UA_DeleteSubscriptionsResponse response;
460
0
    __Client_Service(client, &request,
461
0
                     &UA_TYPES[UA_TYPES_DELETESUBSCRIPTIONSREQUEST],
462
0
                     &response, &UA_TYPES[UA_TYPES_DELETESUBSCRIPTIONSRESPONSE]);
463
464
    /* Process */
465
0
    __Client_Subscription_processDelete(client, &request, &response);
466
467
0
    unlockClient(client);
468
0
    return response;
469
0
}
470
471
UA_StatusCode
472
0
UA_Client_Subscriptions_deleteSingle(UA_Client *client, UA_UInt32 subscriptionId) {
473
0
    UA_DeleteSubscriptionsRequest request;
474
0
    UA_DeleteSubscriptionsRequest_init(&request);
475
0
    request.subscriptionIds = &subscriptionId;
476
0
    request.subscriptionIdsSize = 1;
477
478
0
    UA_DeleteSubscriptionsResponse response =
479
0
        UA_Client_Subscriptions_delete(client, request);
480
481
0
    UA_StatusCode retval = response.responseHeader.serviceResult;
482
0
    if(retval != UA_STATUSCODE_GOOD) {
483
0
        UA_DeleteSubscriptionsResponse_clear(&response);
484
0
        return retval;
485
0
    }
486
487
0
    if(response.resultsSize != 1) {
488
0
        UA_DeleteSubscriptionsResponse_clear(&response);
489
0
        return UA_STATUSCODE_BADINTERNALERROR;
490
0
    }
491
492
0
    retval = response.results[0];
493
0
    UA_DeleteSubscriptionsResponse_clear(&response);
494
0
    return retval;
495
0
}
496
497
/******************/
498
/* MonitoredItems */
499
/******************/
500
501
static void
502
0
EventFields_clear(UA_KeyValueMap *eventFields) {
503
    /* The values are borrowed from the PublishResponse. Detach them before
504
     * clearing the map-owned field-name keys. */
505
0
    for(size_t i = 0; i < eventFields->mapSize; i++)
506
0
        UA_Variant_init(&eventFields->map[i].value);
507
0
    UA_KeyValueMap_clear(eventFields);
508
0
}
509
510
static void
511
MonitoredItem_delete(UA_Client *client, UA_Client_Subscription *sub,
512
0
                     UA_Client_MonitoredItem *mon) {
513
0
    UA_LOCK_ASSERT(&client->clientMutex);
514
515
0
    ZIP_REMOVE(MonitorItemsTree, &sub->monitoredItems, mon);
516
0
    if(mon->deleteCallback)
517
0
        mon->deleteCallback(client, sub->subscriptionId, sub->context,
518
0
                            mon->monitoredItemId, mon->context);
519
0
    EventFields_clear(&mon->eventFields);
520
0
    UA_MonitoringParameters_clear(&mon->parameters);
521
0
    if(mon->pendingParameters.clientHandle != 0)
522
0
        sub->pendingRekeys--;
523
0
    UA_MonitoringParameters_clear(&mon->pendingParameters);
524
0
    UA_free(mon);
525
0
}
526
527
static UA_StatusCode
528
prepareEventFieldsMap(UA_KeyValueMap *eventFields,
529
0
                      UA_MonitoringParameters *params) {
530
    /* Get the EventFilter */
531
0
    UA_ExtensionObject *eo = &params->filter;
532
0
    if(eo->content.decoded.type != &UA_TYPES[UA_TYPES_EVENTFILTER])
533
0
        return UA_STATUSCODE_GOOD;
534
0
    UA_EventFilter *ef = (UA_EventFilter*)eo->content.decoded.data;
535
0
    UA_StatusCode res = UA_STATUSCODE_GOOD;
536
537
    /* Check whether there are fields */
538
0
    if(ef->selectClausesSize == 0)
539
0
        return UA_STATUSCODE_GOOD;
540
541
    /* Allocate the map */
542
0
    eventFields->map = (UA_KeyValuePair*)
543
0
        UA_calloc(ef->selectClausesSize, sizeof(UA_KeyValuePair));
544
0
    if(!eventFields->map)
545
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
546
0
    eventFields->mapSize = ef->selectClausesSize;
547
548
    /* Create the key-strings for the fields */
549
0
    for(size_t i = 0; i < eventFields->mapSize; i++) {
550
0
        res |= UA_SimpleAttributeOperand_print(&ef->selectClauses[i],
551
0
                                               &eventFields->map[i].key.name);
552
0
    }
553
554
0
    return res;
555
0
}
556
557
static UA_StatusCode
558
MonitoredItem_createBegin(UA_Client *client, UA_Client_Subscription *sub,
559
                          UA_MonitoredItemCreateRequest *item,
560
                          UA_Client_DeleteMonitoredItemCallback deleteCallback,
561
                          void *context, void *handlingCallback,
562
0
                          UA_Client_MonitoredItem **outMon) {
563
    /* Allocate MonitoredItem */
564
0
    UA_Client_MonitoredItem *mon = (UA_Client_MonitoredItem *)
565
0
        UA_calloc(1, sizeof(UA_Client_MonitoredItem));
566
0
    if(!mon)
567
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
568
569
    /* Set a unique ClientHandle and retain the active parameters. */
570
0
    item->requestedParameters.clientHandle = ++client->monitoredItemHandles;
571
0
    UA_StatusCode res =
572
0
        UA_MonitoringParameters_copy(&item->requestedParameters,
573
0
                                     &mon->parameters);
574
0
    if(res != UA_STATUSCODE_GOOD) {
575
0
        UA_free(mon);
576
0
        return res;
577
0
    }
578
579
0
    mon->isEventMonitoredItem =
580
0
        (item->itemToMonitor.attributeId == UA_ATTRIBUTEID_EVENTNOTIFIER);
581
582
    /* Fill in members and add to the client  */
583
0
    mon->context = context;
584
0
    mon->deleteCallback = deleteCallback;
585
0
    mon->handler.dataChangeCallback =
586
0
        (UA_Client_DataChangeNotificationCallback)(uintptr_t)handlingCallback;
587
0
    ZIP_INSERT(MonitorItemsTree, &sub->monitoredItems, mon);
588
589
0
    UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
590
0
                 "Subscription %" PRIu32 " | Added a MonitoredItem with handle %" PRIu32,
591
0
                 sub->subscriptionId, mon->parameters.clientHandle);
592
593
0
    *outMon = mon;
594
0
    return UA_STATUSCODE_GOOD;
595
0
}
596
597
static void
598
MonitoredItem_createFinish(UA_Client *client, UA_Client_Subscription *sub,
599
                           UA_Client_MonitoredItem *mon,
600
0
                           UA_MonitoredItemCreateResult *result) {
601
0
    UA_assert(result->statusCode == UA_STATUSCODE_GOOD);
602
603
0
    mon->monitoredItemId = result->monitoredItemId;
604
    /* revisedSamplingInterval; */
605
    /* revisedQueueSize; */
606
    /* filterResult; */
607
0
}
608
609
/************************************/
610
/* CreateMonitoredItems Synchronous */
611
/************************************/
612
613
static void
614
Client_MonitoredItems_create(UA_Client *client,
615
                             const UA_CreateMonitoredItemsRequest *constRequest,
616
                             void **contexts, void **handlingCallbacks,
617
                             UA_Client_DeleteMonitoredItemCallback *deleteCallbacks,
618
0
                             UA_CreateMonitoredItemsResponse *response) {
619
0
    UA_LOCK_ASSERT(&client->clientMutex);
620
0
    UA_CreateMonitoredItemsResponse_init(response);
621
622
    /* Any items? */
623
0
    if(constRequest->itemsToCreateSize == 0) {
624
0
        response->responseHeader.serviceResult = UA_STATUSCODE_BADNOTHINGTODO;
625
0
        return;
626
0
    }
627
628
    /* Get the subscription */
629
0
    UA_Client_Subscription *sub =
630
0
        findSubscriptionById(client, constRequest->subscriptionId);
631
0
    if(!sub) {
632
0
        response->responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
633
0
        return;
634
0
    }
635
636
    /* Make a mutable copy. We modify the request to set the internal
637
     * clientHandle. */
638
0
    UA_CreateMonitoredItemsRequest request;
639
0
    UA_StatusCode res =
640
0
        UA_CreateMonitoredItemsRequest_copy(constRequest, &request);
641
0
    if(res != UA_STATUSCODE_GOOD) {
642
0
        response->responseHeader.serviceResult = res;
643
0
        return;
644
0
    }
645
646
    /* Create the MonitoredItems */
647
0
    UA_STACKARRAY(UA_Client_MonitoredItem*, mons, request.itemsToCreateSize);
648
0
    memset(mons, 0, sizeof(UA_Client_MonitoredItem*) * request.itemsToCreateSize);
649
0
    for(size_t i = 0; i < request.itemsToCreateSize; i++) {
650
0
        void *context = (contexts) ? contexts[i] : NULL;
651
0
        void *handlingCallback = (handlingCallbacks) ? handlingCallbacks[i] : NULL;
652
0
        UA_Client_DeleteMonitoredItemCallback deleteCallback =
653
0
            (deleteCallbacks) ? deleteCallbacks[i] : NULL;
654
0
        res |= MonitoredItem_createBegin(client, sub, &request.itemsToCreate[i],
655
0
                                         deleteCallback, context,
656
0
                                         handlingCallback, &mons[i]);
657
0
    }
658
659
    /* Failure -> Delete created MonitoredItems. Directly call deleteCallback if
660
     * creation failed. The MonitoredItemId is not yet known, use zero. */
661
0
    if(res != UA_STATUSCODE_GOOD) {
662
0
        for(size_t i = 0; i < request.itemsToCreateSize; i++) {
663
0
            if(mons[i])
664
0
                MonitoredItem_delete(client, sub, mons[i]);
665
0
            else if(deleteCallbacks && contexts)
666
0
                deleteCallbacks[i](client, request.subscriptionId,
667
0
                                   sub->context, 0, contexts[i]);
668
0
        }
669
0
        UA_CreateMonitoredItemsRequest_clear(&request);
670
0
        response->responseHeader.serviceResult = res;
671
0
        return;
672
0
    }
673
674
    /* Call the service */
675
0
    __Client_Service(client, &request,
676
0
                     &UA_TYPES[UA_TYPES_CREATEMONITOREDITEMSREQUEST],
677
0
                     response, &UA_TYPES[UA_TYPES_CREATEMONITOREDITEMSRESPONSE]);
678
679
    /* Check that the response size is good */
680
0
    if(response->responseHeader.serviceResult == UA_STATUSCODE_GOOD &&
681
0
       response->resultsSize != request.itemsToCreateSize)
682
0
        response->responseHeader.serviceResult = UA_STATUSCODE_BADINTERNALERROR;
683
684
    /* Update the MonitoredItems */
685
0
    for(size_t i = 0; i < request.itemsToCreateSize; i++) {
686
0
        UA_assert(mons[i]);
687
0
        UA_MonitoredItemCreateResult *item = &response->results[i];
688
0
        if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD ||
689
0
           item->statusCode != UA_STATUSCODE_GOOD) {
690
0
            MonitoredItem_delete(client, sub, mons[i]);
691
0
            continue;
692
0
        }
693
0
        MonitoredItem_createFinish(client, sub, mons[i], item);
694
0
    }
695
696
0
    UA_CreateMonitoredItemsRequest_clear(&request);
697
0
}
698
699
UA_CreateMonitoredItemsResponse
700
UA_Client_MonitoredItems_createDataChanges(UA_Client *client,
701
                                           const UA_CreateMonitoredItemsRequest request,
702
                                           void **contexts,
703
                                           UA_Client_DataChangeNotificationCallback *callbacks,
704
0
                                           UA_Client_DeleteMonitoredItemCallback *deleteCallbacks) {
705
0
    UA_CreateMonitoredItemsResponse response;
706
0
    lockClient(client);
707
0
    Client_MonitoredItems_create(client, &request, contexts, (void **)callbacks,
708
0
                                 deleteCallbacks, &response);
709
0
    unlockClient(client);
710
0
    return response;
711
0
}
712
713
UA_MonitoredItemCreateResult
714
UA_Client_MonitoredItems_createDataChange(UA_Client *client, UA_UInt32 subscriptionId,
715
                                          UA_TimestampsToReturn timestampsToReturn,
716
                                          const UA_MonitoredItemCreateRequest item,
717
                                          void *context,
718
                                          UA_Client_DataChangeNotificationCallback callback,
719
0
                                          UA_Client_DeleteMonitoredItemCallback deleteCallback) {
720
0
    UA_CreateMonitoredItemsRequest request;
721
0
    UA_CreateMonitoredItemsRequest_init(&request);
722
0
    request.subscriptionId = subscriptionId;
723
0
    request.timestampsToReturn = timestampsToReturn;
724
0
    request.itemsToCreate = (UA_MonitoredItemCreateRequest*)(uintptr_t)&item;
725
0
    request.itemsToCreateSize = 1;
726
0
    UA_CreateMonitoredItemsResponse response =
727
0
       UA_Client_MonitoredItems_createDataChanges(client, request, &context,
728
0
                                                   &callback, &deleteCallback);
729
0
    UA_MonitoredItemCreateResult result;
730
0
    UA_MonitoredItemCreateResult_init(&result);
731
0
    if(response.responseHeader.serviceResult != UA_STATUSCODE_GOOD)
732
0
        result.statusCode = response.responseHeader.serviceResult;
733
0
    if(result.statusCode == UA_STATUSCODE_GOOD &&
734
0
       response.resultsSize != 1)
735
0
        result.statusCode = UA_STATUSCODE_BADINTERNALERROR;
736
0
    if(result.statusCode == UA_STATUSCODE_GOOD) {
737
0
        result = response.results[0];
738
0
        UA_MonitoredItemCreateResult_init(&response.results[0]);
739
0
    }
740
0
    UA_CreateMonitoredItemsResponse_clear(&response);
741
0
    return result;
742
0
}
743
744
UA_CreateMonitoredItemsResponse
745
UA_Client_MonitoredItems_createEvents(UA_Client *client,
746
                                      const UA_CreateMonitoredItemsRequest request,
747
                                      void **contexts,
748
                                      UA_Client_EventNotificationCallback *callback,
749
0
                                      UA_Client_DeleteMonitoredItemCallback *deleteCallback) {
750
0
    UA_CreateMonitoredItemsResponse response;
751
0
    lockClient(client);
752
0
    Client_MonitoredItems_create(client, &request, contexts, (void **)callback,
753
0
                                 deleteCallback, &response);
754
0
    unlockClient(client);
755
0
    return response;
756
0
}
757
758
UA_MonitoredItemCreateResult
759
UA_Client_MonitoredItems_createEvent(UA_Client *client, UA_UInt32 subscriptionId,
760
                                     UA_TimestampsToReturn timestampsToReturn,
761
                                     const UA_MonitoredItemCreateRequest item, void *context,
762
                                     UA_Client_EventNotificationCallback callback,
763
0
                                     UA_Client_DeleteMonitoredItemCallback deleteCallback) {
764
0
    UA_CreateMonitoredItemsRequest request;
765
0
    UA_CreateMonitoredItemsRequest_init(&request);
766
0
    request.subscriptionId = subscriptionId;
767
0
    request.timestampsToReturn = timestampsToReturn;
768
0
    request.itemsToCreate = (UA_MonitoredItemCreateRequest*)(uintptr_t)&item;
769
0
    request.itemsToCreateSize = 1;
770
0
    UA_CreateMonitoredItemsResponse response =
771
0
       UA_Client_MonitoredItems_createEvents(client, request, &context,
772
0
                                             &callback, &deleteCallback);
773
0
    UA_MonitoredItemCreateResult result;
774
0
    UA_MonitoredItemCreateResult_init(&result);
775
0
    if(response.responseHeader.serviceResult != UA_STATUSCODE_GOOD)
776
0
        result.statusCode = response.responseHeader.serviceResult;
777
0
    if(result.statusCode == UA_STATUSCODE_GOOD &&
778
0
       response.resultsSize != 1)
779
0
        result.statusCode = UA_STATUSCODE_BADINTERNALERROR;
780
0
    if(result.statusCode == UA_STATUSCODE_GOOD) {
781
0
        result = response.results[0];
782
0
        UA_MonitoredItemCreateResult_init(&response.results[0]);
783
0
    }
784
0
    UA_CreateMonitoredItemsResponse_clear(&response);
785
0
    return result;
786
0
}
787
788
/*************************************/
789
/* CreateMonitoredItems Asynchronous */
790
/*************************************/
791
792
/* The handles are an array of [subId, monSize, monHandleId1, monHandle2, ...] */
793
static void
794
MonitoredItems_create_async_handler(UA_Client *client, void *data,
795
0
                                    UA_UInt32 requestId, void *resp) {
796
0
    CustomCallback *cc = (CustomCallback*)data;
797
0
    UA_UInt32 *handles = (UA_UInt32*)cc->clientData;
798
0
    UA_CreateMonitoredItemsResponse *response =
799
0
        (UA_CreateMonitoredItemsResponse *)resp;
800
801
0
    lockClient(client);
802
803
    /* Extract the first elements from the handles */
804
0
    UA_UInt32 subId = handles[0];
805
0
    UA_UInt32 monSize = handles[1];
806
0
    UA_UInt32 *monHandles = handles + 2;
807
808
    /* Check that the response size is good */
809
0
    if(response->responseHeader.serviceResult == UA_STATUSCODE_GOOD &&
810
0
       response->resultsSize != monSize)
811
0
        response->responseHeader.serviceResult = UA_STATUSCODE_BADINTERNALERROR;
812
813
    /* Get the Subscription from the SubscriptionId */
814
0
    UA_Client_Subscription *sub = findSubscriptionById(client, subId);
815
0
    if(!sub)
816
0
        response->responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
817
818
    /* Update the MonitoredItems */
819
0
    for(size_t i = 0; sub && i < monSize; i++) {
820
        /* Get the MonitoredItem from the ClientHandle */
821
0
        UA_Client_MonitoredItem *mon =
822
0
            findMonitoredItemByHandle(sub, monHandles[i]);
823
0
        if(!mon)
824
0
            continue;
825
826
        /* Delete MonitoredItem if the creation failed */
827
0
        UA_MonitoredItemCreateResult *item = &response->results[i];
828
0
        if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD ||
829
0
           item->statusCode != UA_STATUSCODE_GOOD) {
830
0
            MonitoredItem_delete(client, sub, mon);
831
0
            continue;
832
0
        }
833
834
        /* Update the MonitoredItem with the server's response  */
835
0
        MonitoredItem_createFinish(client, sub, mon, item);
836
0
    }
837
838
    /* Notify the application */
839
0
    if(cc->userCallback) {
840
0
        UA_ClientAsyncCreateMonitoredItemsCallback cb =
841
0
            (UA_ClientAsyncCreateMonitoredItemsCallback)cc->userCallback;
842
0
        cb(client, cc->userData, requestId, response);
843
0
    }
844
845
    /* Clean up */
846
0
    UA_free(handles);
847
0
    UA_free(cc);
848
849
0
    unlockClient(client);
850
0
}
851
852
static UA_StatusCode
853
Client_MonitoredItems_createAsync(UA_Client *client,
854
                                  const UA_CreateMonitoredItemsRequest *constRequest,
855
                                  void **contexts, void **handlingCallbacks,
856
                                  UA_Client_DeleteMonitoredItemCallback *deleteCallbacks,
857
                                  UA_ClientAsyncServiceCallback createCallback,
858
0
                                  void *userdata, UA_UInt32 *requestId) {
859
0
    UA_LOCK_ASSERT(&client->clientMutex);
860
861
    /* Get the Subscription */
862
0
    UA_Client_Subscription *sub =
863
0
        findSubscriptionById(client, constRequest->subscriptionId);
864
0
    if(!sub)
865
0
        return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
866
867
    /* Any items? */
868
0
    if(constRequest->itemsToCreateSize == 0)
869
0
        return UA_STATUSCODE_BADNOTHINGTODO;
870
871
    /* Make a mutable copy. We modify the request to set the internal
872
     * clientHandle. */
873
0
    UA_CreateMonitoredItemsRequest request;
874
0
    UA_StatusCode res =
875
0
        UA_CreateMonitoredItemsRequest_copy(constRequest, &request);
876
0
    if(res != UA_STATUSCODE_GOOD)
877
0
        return res;
878
879
    /* Allocate context for the async handling */
880
0
    CustomCallback *cc = (CustomCallback*)UA_calloc(1, sizeof(CustomCallback));
881
0
    if(!cc) {
882
0
        UA_CreateMonitoredItemsRequest_clear(&request);
883
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
884
0
    }
885
0
    UA_UInt32 *handles = (UA_UInt32*)
886
0
        UA_malloc(sizeof(UA_UInt32) * (request.itemsToCreateSize + 2));
887
0
    if(!handles) {
888
0
        UA_free(cc);
889
0
        UA_CreateMonitoredItemsRequest_clear(&request);
890
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
891
0
    }
892
893
    /* Create the MonitoredItems locally */
894
0
    UA_STACKARRAY(UA_Client_MonitoredItem*, mons, request.itemsToCreateSize);
895
0
    memset(mons, 0, sizeof(UA_Client_MonitoredItem*) * request.itemsToCreateSize);
896
0
    for(size_t i = 0; i < request.itemsToCreateSize; i++) {
897
0
        void *context = (contexts) ? contexts[i] : NULL;
898
0
        void *handlingCallback = (handlingCallbacks) ? handlingCallbacks[i] : NULL;
899
0
        UA_Client_DeleteMonitoredItemCallback deleteCallback =
900
0
            (deleteCallbacks) ? deleteCallbacks[i] : NULL;
901
0
        res |= MonitoredItem_createBegin(client, sub, &request.itemsToCreate[i],
902
0
                                         deleteCallback, context,
903
0
                                         handlingCallback, &mons[i]);
904
0
    }
905
906
    /* Failure -> Delete created MonitoredItems. Directly call deleteCallback if
907
     * creation failed. The MonitoredItemId is not yet known, use zero. */
908
0
    if(res != UA_STATUSCODE_GOOD) {
909
0
        for(size_t i = 0; i < request.itemsToCreateSize; i++) {
910
0
            if(mons[i])
911
0
                MonitoredItem_delete(client, sub, mons[i]);
912
0
            else if(deleteCallbacks && contexts)
913
0
                deleteCallbacks[i](client, request.subscriptionId,
914
0
                                   sub->context, 0, contexts[i]);
915
0
        }
916
0
        UA_free(handles);
917
0
        UA_free(cc);
918
0
        UA_CreateMonitoredItemsRequest_clear(&request);
919
0
        return res;
920
0
    }
921
922
    /* Set the async handler context */
923
0
    handles[0] = sub->subscriptionId;
924
0
    handles[1] = (UA_UInt32)request.itemsToCreateSize;
925
0
    for(size_t i = 0; i < request.itemsToCreateSize; i++) {
926
0
        handles[i+2] = mons[i]->parameters.clientHandle;
927
0
    }
928
0
    cc->clientData = handles;
929
0
    cc->userCallback = createCallback;
930
0
    cc->userData = userdata;
931
932
    /* Call the service asynchronously */
933
0
    res = __Client_AsyncService(client, &request,
934
0
                                &UA_TYPES[UA_TYPES_CREATEMONITOREDITEMSREQUEST],
935
0
                                MonitoredItems_create_async_handler,
936
0
                                &UA_TYPES[UA_TYPES_CREATEMONITOREDITEMSRESPONSE],
937
0
                                cc, requestId);
938
939
    /* Manually clean up the context in the failure case */
940
0
    if(res != UA_STATUSCODE_GOOD) {
941
0
        for(size_t i = 0; i < request.itemsToCreateSize; i++) {
942
0
            UA_assert(mons[i]);
943
0
            MonitoredItem_delete(client, sub, mons[i]);
944
0
        }
945
0
        UA_free(handles);
946
0
        UA_free(cc);
947
0
    }
948
949
    /* Clean up */
950
0
    UA_CreateMonitoredItemsRequest_clear(&request);
951
0
    return res;
952
0
}
953
954
UA_StatusCode
955
UA_Client_MonitoredItems_createDataChanges_async(UA_Client *client,
956
                                                 const UA_CreateMonitoredItemsRequest request,
957
                                                 void **contexts,
958
                                                 UA_Client_DataChangeNotificationCallback *callbacks,
959
                                                 UA_Client_DeleteMonitoredItemCallback *deleteCallbacks,
960
                                                 UA_ClientAsyncCreateMonitoredItemsCallback createCallback,
961
0
                                                 void *userdata, UA_UInt32 *requestId) {
962
0
    lockClient(client);
963
0
    UA_StatusCode res =
964
0
        Client_MonitoredItems_createAsync(client, &request, contexts,
965
0
                                          (void **)callbacks, deleteCallbacks,
966
0
                                          (UA_ClientAsyncServiceCallback)createCallback,
967
0
                                          userdata, requestId);
968
0
    unlockClient(client);
969
0
    return res;
970
0
}
971
972
UA_StatusCode
973
UA_Client_MonitoredItems_createEvents_async(UA_Client *client,
974
                                            const UA_CreateMonitoredItemsRequest request,
975
                                            void **contexts,
976
                                            UA_Client_EventNotificationCallback *callbacks,
977
                                            UA_Client_DeleteMonitoredItemCallback *deleteCallbacks,
978
                                            UA_ClientAsyncCreateMonitoredItemsCallback createCallback,
979
0
                                            void *userdata, UA_UInt32 *requestId) {
980
0
    lockClient(client);
981
0
    UA_StatusCode res =
982
0
        Client_MonitoredItems_createAsync(client, &request, contexts,
983
0
                                          (void **)callbacks, deleteCallbacks,
984
0
                                          (UA_ClientAsyncServiceCallback)createCallback,
985
0
                                          userdata, requestId);
986
0
    unlockClient(client);
987
0
    return res;
988
0
}
989
990
static void
991
MonitoredItems_delete(UA_Client *client, UA_Client_Subscription *sub,
992
                      const UA_DeleteMonitoredItemsRequest *request,
993
0
                      const UA_DeleteMonitoredItemsResponse *response) {
994
#ifdef __clang_analyzer__
995
    return;
996
#endif
997
998
    /* Loop over deleted MonitoredItems */
999
0
    struct UA_Client_MonitoredItem_ForDelete deleteMonitoredItem;
1000
0
    memset(&deleteMonitoredItem, 0, sizeof(struct UA_Client_MonitoredItem_ForDelete));
1001
0
    deleteMonitoredItem.client = client;
1002
0
    deleteMonitoredItem.sub = sub;
1003
1004
0
    for(size_t i = 0; i < response->resultsSize; i++) {
1005
0
        if(response->results[i] != UA_STATUSCODE_GOOD &&
1006
0
           response->results[i] != UA_STATUSCODE_BADMONITOREDITEMIDINVALID) {
1007
0
            continue;
1008
0
        }
1009
0
        deleteMonitoredItem.monitoredItemId = &request->monitoredItemIds[i];
1010
        /* Delete the internal representation */
1011
0
        ZIP_ITER(MonitorItemsTree,&sub->monitoredItems,
1012
0
                 MonitoredItem_delete_wrapper, &deleteMonitoredItem);
1013
0
    }
1014
0
}
1015
1016
static void
1017
0
MonitoredItems_delete_handler(UA_Client *client, void *d, UA_UInt32 requestId, void *r) {
1018
0
    UA_Client_Subscription *sub = NULL;
1019
0
    CustomCallback *cc = (CustomCallback *)d;
1020
0
    UA_DeleteMonitoredItemsResponse *response = (UA_DeleteMonitoredItemsResponse *)r;
1021
0
    UA_DeleteMonitoredItemsRequest *request =
1022
0
        (UA_DeleteMonitoredItemsRequest *)cc->clientData;
1023
1024
0
    lockClient(client);
1025
1026
0
    if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD)
1027
0
        goto cleanup;
1028
1029
0
    sub = findSubscriptionById(client, request->subscriptionId);
1030
0
    if(!sub) {
1031
0
        UA_LOG_INFO(client->config.logging, UA_LOGCATEGORY_CLIENT,
1032
0
                    "No internal representation of subscription %" PRIu32,
1033
0
                    request->subscriptionId);
1034
0
        goto cleanup;
1035
0
    }
1036
1037
    /* Delete MonitoredItems from the internal representation */
1038
0
    MonitoredItems_delete(client, sub, request, response);
1039
1040
0
cleanup:
1041
0
    if(cc->userCallback)
1042
0
        cc->userCallback(client, cc->userData, requestId, response);
1043
0
    UA_DeleteMonitoredItemsRequest_delete(request);
1044
0
    UA_free(cc);
1045
1046
0
    unlockClient(client);
1047
0
}
1048
1049
UA_DeleteMonitoredItemsResponse
1050
UA_Client_MonitoredItems_delete(UA_Client *client,
1051
0
                                const UA_DeleteMonitoredItemsRequest request) {
1052
    /* Send the request */
1053
0
    UA_DeleteMonitoredItemsResponse response;
1054
0
    __UA_Client_Service(client, &request, &UA_TYPES[UA_TYPES_DELETEMONITOREDITEMSREQUEST],
1055
0
                        &response, &UA_TYPES[UA_TYPES_DELETEMONITOREDITEMSRESPONSE]);
1056
1057
    /* A problem occured remote? */
1058
0
    if(response.responseHeader.serviceResult != UA_STATUSCODE_GOOD)
1059
0
        return response;
1060
1061
0
    lockClient(client);
1062
1063
    /* Find the internal subscription representation */
1064
0
    UA_Client_Subscription *sub = findSubscriptionById(client, request.subscriptionId);
1065
0
    if(!sub) {
1066
0
        UA_LOG_INFO(client->config.logging, UA_LOGCATEGORY_CLIENT,
1067
0
                    "No internal representation of subscription %" PRIu32,
1068
0
                    request.subscriptionId);
1069
0
        unlockClient(client);
1070
0
        return response;
1071
0
    }
1072
1073
    /* Remove MonitoredItems in the internal representation */
1074
0
    MonitoredItems_delete(client, sub, &request, &response);
1075
1076
0
    unlockClient(client);
1077
1078
0
    return response;
1079
0
}
1080
1081
UA_StatusCode
1082
UA_Client_MonitoredItems_delete_async(UA_Client *client,
1083
                                      const UA_DeleteMonitoredItemsRequest request,
1084
                                      UA_ClientAsyncDeleteMonitoredItemsCallback callback,
1085
0
                                      void *userdata, UA_UInt32 *requestId) {
1086
    /* Send the request */
1087
0
    CustomCallback *cc = (CustomCallback *)UA_calloc(1, sizeof(CustomCallback));
1088
0
    if(!cc)
1089
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
1090
1091
0
    UA_DeleteMonitoredItemsRequest *req_copy = UA_DeleteMonitoredItemsRequest_new();
1092
0
    if(!req_copy) {
1093
0
        UA_free(cc);
1094
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
1095
0
    }
1096
1097
0
    UA_DeleteMonitoredItemsRequest_copy(&request, req_copy);
1098
0
    cc->clientData = req_copy;
1099
0
    cc->userCallback = (UA_ClientAsyncServiceCallback)callback;
1100
0
    cc->userData = userdata;
1101
1102
0
    UA_StatusCode res =
1103
0
        __UA_Client_AsyncService(client, &request,
1104
0
                                 &UA_TYPES[UA_TYPES_DELETEMONITOREDITEMSREQUEST],
1105
0
                                 MonitoredItems_delete_handler,
1106
0
                                 &UA_TYPES[UA_TYPES_DELETEMONITOREDITEMSRESPONSE],
1107
0
                                 cc, requestId);
1108
0
    if(res != UA_STATUSCODE_GOOD) {
1109
0
        UA_DeleteMonitoredItemsRequest_delete(req_copy);
1110
0
        UA_free(cc);
1111
0
    }
1112
0
    return res;
1113
0
}
1114
1115
UA_StatusCode
1116
UA_Client_MonitoredItems_deleteSingle(UA_Client *client, UA_UInt32 subscriptionId,
1117
0
                                      UA_UInt32 monitoredItemId) {
1118
0
    UA_DeleteMonitoredItemsRequest request;
1119
0
    UA_DeleteMonitoredItemsRequest_init(&request);
1120
0
    request.subscriptionId = subscriptionId;
1121
0
    request.monitoredItemIds = &monitoredItemId;
1122
0
    request.monitoredItemIdsSize = 1;
1123
1124
0
    UA_DeleteMonitoredItemsResponse response =
1125
0
        UA_Client_MonitoredItems_delete(client, request);
1126
1127
0
    UA_StatusCode retval = response.responseHeader.serviceResult;
1128
0
    if(retval != UA_STATUSCODE_GOOD) {
1129
0
        UA_DeleteMonitoredItemsResponse_clear(&response);
1130
0
        return retval;
1131
0
    }
1132
1133
0
    if(response.resultsSize != 1) {
1134
0
        UA_DeleteMonitoredItemsResponse_clear(&response);
1135
0
        return UA_STATUSCODE_BADINTERNALERROR;
1136
0
    }
1137
1138
0
    retval = response.results[0];
1139
0
    UA_DeleteMonitoredItemsResponse_clear(&response);
1140
0
    return retval;
1141
0
}
1142
1143
static void *
1144
0
MonitoredItem_findByID(void *data, UA_Client_MonitoredItem *mon) {
1145
0
    UA_UInt32 monitorId = *(UA_UInt32*)data;
1146
0
    if(monitorId && (mon->monitoredItemId == monitorId))
1147
0
        return mon;
1148
0
    return NULL;
1149
0
}
1150
1151
static UA_Client_MonitoredItem *
1152
0
findMonitoredItemById(UA_Client_Subscription *sub, UA_UInt32 monitoredItemId) {
1153
0
    return (UA_Client_MonitoredItem *)
1154
0
        ZIP_ITER(MonitorItemsTree, &sub->monitoredItems,
1155
0
                 MonitoredItem_findByID, &monitoredItemId);
1156
0
}
1157
1158
static UA_StatusCode
1159
MonitoredItems_prepareModify(UA_Client *client, UA_Client_Subscription *sub,
1160
0
                             UA_ModifyMonitoredItemsRequest *request) {
1161
0
    UA_STACKARRAY(UA_Client_MonitoredItem*, mons, request->itemsToModifySize);
1162
0
    UA_STACKARRAY(UA_MonitoringParameters, preparedParameters,
1163
0
                  request->itemsToModifySize);
1164
0
    memset(mons, 0, sizeof(UA_Client_MonitoredItem*) *
1165
0
           request->itemsToModifySize);
1166
0
    memset(preparedParameters, 0, sizeof(UA_MonitoringParameters) *
1167
0
           request->itemsToModifySize);
1168
1169
    /* Prepare every settings slot before changing the MonitoredItems. So an
1170
     * allocation failure leaves all active and pending settings untouched. */
1171
0
    for(size_t i = 0; i < request->itemsToModifySize; ++i) {
1172
0
        UA_MonitoredItemModifyRequest *mimr = &request->itemsToModify[i];
1173
0
        mimr->requestedParameters.clientHandle = 0;
1174
0
        mons[i] = findMonitoredItemById(sub, mimr->monitoredItemId);
1175
0
        if(!mons[i])
1176
0
            continue;
1177
1178
0
        mimr->requestedParameters.clientHandle = ++client->monitoredItemHandles;
1179
0
        UA_StatusCode res =
1180
0
            UA_MonitoringParameters_copy(&mimr->requestedParameters,
1181
0
                                         &preparedParameters[i]);
1182
0
        if(res != UA_STATUSCODE_GOOD) {
1183
0
            for(size_t j = 0; j <= i; j++)
1184
0
                UA_MonitoringParameters_clear(&preparedParameters[j]);
1185
0
            return res;
1186
0
        }
1187
0
    }
1188
1189
    /* Commit the prepared settings. A newer modification supersedes a pending
1190
     * generation that has not produced a notification yet. */
1191
0
    for(size_t i = 0; i < request->itemsToModifySize; ++i) {
1192
0
        UA_Client_MonitoredItem *mon = mons[i];
1193
0
        if(!mon)
1194
0
            continue;
1195
        /* Newly entering the pending re-key state (a fresh handle is
1196
         * always non-zero); a supersede keeps it pending, no double count. */
1197
0
        if(mon->pendingParameters.clientHandle == 0)
1198
0
            sub->pendingRekeys++;
1199
0
        UA_MonitoringParameters_clear(&mon->pendingParameters);
1200
0
        mon->pendingParameters = preparedParameters[i];
1201
0
        memset(&preparedParameters[i], 0, sizeof(UA_MonitoringParameters));
1202
0
    }
1203
0
    return UA_STATUSCODE_GOOD;
1204
0
}
1205
1206
static void
1207
MonitoredItems_reconcileModify(UA_Client_Subscription *sub,
1208
                               const UA_ModifyMonitoredItemsRequest *request,
1209
0
                               UA_ModifyMonitoredItemsResponse *response) {
1210
0
    UA_Boolean validResponse = response &&
1211
0
        response->responseHeader.serviceResult == UA_STATUSCODE_GOOD &&
1212
0
        response->resultsSize == request->itemsToModifySize;
1213
0
    if(response && response->responseHeader.serviceResult == UA_STATUSCODE_GOOD &&
1214
0
       !validResponse)
1215
0
        response->responseHeader.serviceResult = UA_STATUSCODE_BADINTERNALERROR;
1216
1217
0
    for(size_t i = 0; i < request->itemsToModifySize; ++i) {
1218
0
        if(validResponse && response->results[i].statusCode == UA_STATUSCODE_GOOD)
1219
0
            continue;
1220
0
        const UA_MonitoredItemModifyRequest *mimr = &request->itemsToModify[i];
1221
0
        UA_Client_MonitoredItem *mon =
1222
0
            findMonitoredItemById(sub, mimr->monitoredItemId);
1223
0
        if(mon && mon->pendingParameters.clientHandle ==
1224
0
                  mimr->requestedParameters.clientHandle) {
1225
0
            UA_MonitoringParameters_clear(&mon->pendingParameters);
1226
0
            sub->pendingRekeys--;
1227
0
        }
1228
0
    }
1229
0
}
1230
1231
UA_ModifyMonitoredItemsResponse
1232
UA_Client_MonitoredItems_modify(UA_Client *client,
1233
0
                                const UA_ModifyMonitoredItemsRequest request) {
1234
0
    UA_ModifyMonitoredItemsResponse response;
1235
0
    UA_ModifyMonitoredItemsResponse_init(&response);
1236
1237
    /* Make a modifiable copy of the request */
1238
0
    UA_ModifyMonitoredItemsRequest modifiedRequest;
1239
0
    UA_StatusCode res = UA_ModifyMonitoredItemsRequest_copy(&request, &modifiedRequest);
1240
0
    if(res != UA_STATUSCODE_GOOD) {
1241
0
        response.responseHeader.serviceResult = res;
1242
0
        return response;
1243
0
    }
1244
1245
0
    lockClient(client);
1246
1247
    /* Get the subscription */
1248
0
    UA_Client_Subscription *sub = findSubscriptionById(client, request.subscriptionId);
1249
0
    if(!sub) {
1250
0
        unlockClient(client);
1251
0
        UA_ModifyMonitoredItemsRequest_clear(&modifiedRequest);
1252
0
        response.responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
1253
0
        return response;
1254
0
    }
1255
1256
0
    res = MonitoredItems_prepareModify(client, sub, &modifiedRequest);
1257
0
    if(res != UA_STATUSCODE_GOOD) {
1258
0
        response.responseHeader.serviceResult = res;
1259
0
        unlockClient(client);
1260
0
        UA_ModifyMonitoredItemsRequest_clear(&modifiedRequest);
1261
0
        return response;
1262
0
    }
1263
1264
    /* Call the service */
1265
0
    __Client_Service(client, &modifiedRequest,
1266
0
                     &UA_TYPES[UA_TYPES_MODIFYMONITOREDITEMSREQUEST], &response,
1267
0
                     &UA_TYPES[UA_TYPES_MODIFYMONITOREDITEMSRESPONSE]);
1268
1269
0
    MonitoredItems_reconcileModify(sub, &modifiedRequest, &response);
1270
1271
0
    unlockClient(client);
1272
1273
    /* Clean up */
1274
0
    UA_ModifyMonitoredItemsRequest_clear(&modifiedRequest);
1275
0
    return response;
1276
0
}
1277
1278
static void
1279
MonitoredItems_modify_async_handler(UA_Client *client, void *data,
1280
0
                                    UA_UInt32 requestId, void *resp) {
1281
0
    CustomCallback *cc = (CustomCallback*)data;
1282
0
    UA_ModifyMonitoredItemsRequest *request =
1283
0
        (UA_ModifyMonitoredItemsRequest*)cc->clientData;
1284
0
    UA_ModifyMonitoredItemsResponse *response =
1285
0
        (UA_ModifyMonitoredItemsResponse*)resp;
1286
1287
0
    lockClient(client);
1288
0
    UA_Client_Subscription *sub =
1289
0
        findSubscriptionById(client, request->subscriptionId);
1290
0
    if(sub)
1291
0
        MonitoredItems_reconcileModify(sub, request, response);
1292
1293
0
    if(cc->userCallback) {
1294
0
        UA_ClientAsyncModifyMonitoredItemsCallback cb =
1295
0
            (UA_ClientAsyncModifyMonitoredItemsCallback)cc->userCallback;
1296
0
        cb(client, cc->userData, requestId, response);
1297
0
    }
1298
1299
0
    UA_ModifyMonitoredItemsRequest_delete(request);
1300
0
    UA_free(cc);
1301
0
    unlockClient(client);
1302
0
}
1303
1304
UA_StatusCode
1305
UA_Client_MonitoredItems_modify_async(UA_Client *client,
1306
                                      const UA_ModifyMonitoredItemsRequest request,
1307
                                      UA_ClientAsyncModifyMonitoredItemsCallback callback,
1308
0
                                      void *userdata, UA_UInt32 *requestId) {
1309
0
    UA_ModifyMonitoredItemsRequest *requestCopy =
1310
0
        UA_ModifyMonitoredItemsRequest_new();
1311
0
    if(!requestCopy)
1312
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
1313
0
    UA_StatusCode res =
1314
0
        UA_ModifyMonitoredItemsRequest_copy(&request, requestCopy);
1315
0
    if(res != UA_STATUSCODE_GOOD) {
1316
0
        UA_ModifyMonitoredItemsRequest_delete(requestCopy);
1317
0
        return res;
1318
0
    }
1319
1320
0
    lockClient(client);
1321
1322
    /* Get the subscription */
1323
0
    UA_Client_Subscription *sub = findSubscriptionById(client, request.subscriptionId);
1324
0
    if(!sub) {
1325
0
        unlockClient(client);
1326
0
        UA_ModifyMonitoredItemsRequest_delete(requestCopy);
1327
0
        return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
1328
0
    }
1329
1330
0
    res = MonitoredItems_prepareModify(client, sub, requestCopy);
1331
0
    if(res != UA_STATUSCODE_GOOD) {
1332
0
        unlockClient(client);
1333
0
        UA_ModifyMonitoredItemsRequest_delete(requestCopy);
1334
0
        return res;
1335
0
    }
1336
1337
0
    CustomCallback *cc = (CustomCallback*)UA_calloc(1, sizeof(CustomCallback));
1338
0
    if(!cc) {
1339
0
        MonitoredItems_reconcileModify(sub, requestCopy, NULL);
1340
0
        unlockClient(client);
1341
0
        UA_ModifyMonitoredItemsRequest_delete(requestCopy);
1342
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
1343
0
    }
1344
0
    cc->clientData = requestCopy;
1345
0
    cc->userCallback = (UA_ClientAsyncServiceCallback)callback;
1346
0
    cc->userData = userdata;
1347
1348
    /* Call the service */
1349
0
    UA_StatusCode statusCode = __Client_AsyncService(
1350
0
        client, requestCopy, &UA_TYPES[UA_TYPES_MODIFYMONITOREDITEMSREQUEST],
1351
0
        MonitoredItems_modify_async_handler,
1352
0
        &UA_TYPES[UA_TYPES_MODIFYMONITOREDITEMSRESPONSE], cc, requestId);
1353
1354
0
    if(statusCode != UA_STATUSCODE_GOOD) {
1355
0
        MonitoredItems_reconcileModify(sub, requestCopy, NULL);
1356
0
        UA_ModifyMonitoredItemsRequest_delete(requestCopy);
1357
0
        UA_free(cc);
1358
0
    }
1359
1360
0
    unlockClient(client);
1361
0
    return statusCode;
1362
0
}
1363
1364
UA_StatusCode
1365
UA_Client_MonitoredItem_getContext(UA_Client *client, UA_UInt32 subscriptionId,
1366
0
                                   UA_UInt32 monitoredItemId, void **monContext) {
1367
0
    if(!client || !monContext)
1368
0
        return UA_STATUSCODE_BADINVALIDARGUMENT;
1369
1370
0
    *monContext = NULL;
1371
1372
0
    lockClient(client);
1373
0
    UA_Client_Subscription *sub = findSubscriptionById(client, subscriptionId);
1374
0
    if(!sub) {
1375
0
        unlockClient(client);
1376
0
        return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
1377
0
    }
1378
1379
0
    UA_StatusCode itemstatus = UA_STATUSCODE_BADMONITOREDITEMIDINVALID;
1380
0
    UA_Client_MonitoredItem *monItem = findMonitoredItemById(sub, monitoredItemId);
1381
0
    if(monItem) {
1382
0
        *monContext = monItem->context;
1383
0
        itemstatus = UA_STATUSCODE_GOOD;
1384
0
    }
1385
0
    unlockClient(client);
1386
0
    return itemstatus;
1387
0
}
1388
1389
UA_StatusCode
1390
UA_Client_MonitoredItem_setContext(UA_Client *client, UA_UInt32 subscriptionId,
1391
0
                                   UA_UInt32 monitoredItemId, void *monContext) {
1392
0
    if(!client)
1393
0
        return UA_STATUSCODE_BADINVALIDARGUMENT;
1394
1395
0
    lockClient(client);
1396
0
    UA_Client_Subscription *sub = findSubscriptionById(client, subscriptionId);
1397
0
    if(!sub) {
1398
0
        unlockClient(client);
1399
0
        return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
1400
0
    }
1401
1402
0
    UA_StatusCode itemstatus = UA_STATUSCODE_BADMONITOREDITEMIDINVALID;
1403
0
    UA_Client_MonitoredItem *monItem = findMonitoredItemById(sub, monitoredItemId);
1404
0
    if(monItem) {
1405
0
        monItem->context = monContext;
1406
0
        itemstatus = UA_STATUSCODE_GOOD;
1407
0
    }
1408
0
    unlockClient(client);
1409
0
    return itemstatus;
1410
0
}
1411
1412
/*************************************/
1413
/* Async Processing of Notifications */
1414
/*************************************/
1415
1416
/* Assume the request is already initialized */
1417
UA_StatusCode
1418
0
__Client_preparePublishRequest(UA_Client *client, UA_PublishRequest *request) {
1419
0
    UA_LOCK_ASSERT(&client->clientMutex);
1420
1421
    /* Count acks */
1422
0
    UA_Client_NotificationsAckNumber *ack;
1423
0
    LIST_FOREACH(ack, &client->pendingNotificationsAcks, listEntry)
1424
0
        ++request->subscriptionAcknowledgementsSize;
1425
1426
    /* Create the array. Returns a sentinel pointer if the length is zero. */
1427
0
    request->subscriptionAcknowledgements = (UA_SubscriptionAcknowledgement*)
1428
0
        UA_Array_new(request->subscriptionAcknowledgementsSize,
1429
0
                     &UA_TYPES[UA_TYPES_SUBSCRIPTIONACKNOWLEDGEMENT]);
1430
0
    if(!request->subscriptionAcknowledgements) {
1431
0
        request->subscriptionAcknowledgementsSize = 0;
1432
0
        return UA_STATUSCODE_BADOUTOFMEMORY;
1433
0
    }
1434
1435
0
    size_t i = 0;
1436
0
    UA_Client_NotificationsAckNumber *ack_tmp;
1437
0
    LIST_FOREACH_SAFE(ack, &client->pendingNotificationsAcks, listEntry, ack_tmp) {
1438
0
        LIST_REMOVE(ack, listEntry);
1439
0
        UA_SubscriptionAcknowledgement *reqAck = &request->subscriptionAcknowledgements[i];
1440
0
        reqAck->sequenceNumber = ack->subAck.sequenceNumber;
1441
0
        reqAck->subscriptionId = ack->subAck.subscriptionId;
1442
0
        UA_free(ack);
1443
0
        i++;
1444
0
    }
1445
0
    return UA_STATUSCODE_GOOD;
1446
0
}
1447
1448
/* According to specification, Part 4 5.13.1, the value 0 is never used for the
1449
 * sequence number */
1450
static UA_UInt32
1451
0
__nextSequenceNumber(UA_UInt32 sequenceNumber) {
1452
0
    UA_UInt32 nextSequenceNumber = sequenceNumber + 1;
1453
0
    if(nextSequenceNumber == 0)
1454
0
        nextSequenceNumber = 1;
1455
0
    return nextSequenceNumber;
1456
0
}
1457
1458
static UA_Client_MonitoredItem *
1459
findMonitoredItemForNotification(UA_Client_Subscription *sub,
1460
0
                                 UA_UInt32 clientHandle) {
1461
0
    UA_Client_MonitoredItem *mon =
1462
0
        findMonitoredItemByHandle(sub, clientHandle);
1463
0
    if(mon)
1464
0
        return mon;
1465
1466
    /* A notification is the authoritative signal that the server has started
1467
     * to use the modified settings. Promote the embedded pending slot and
1468
     * re-key the existing tree node. */
1469
    /* Skip the O(n) scan entirely when no re-key is pending: stops a
1470
     * malicious server from burning CPU on the single client thread by
1471
     * flooding notifications with unknown clientHandles (O(N*n)). */
1472
0
    if(sub->pendingRekeys == 0)
1473
0
        return NULL;
1474
0
    mon = (UA_Client_MonitoredItem*)
1475
0
        ZIP_ITER(MonitorItemsTree, &sub->monitoredItems,
1476
0
                 MonitoredItem_findPendingByHandle, &clientHandle);
1477
0
    if(!mon)
1478
0
        return NULL;
1479
1480
0
    ZIP_REMOVE(MonitorItemsTree, &sub->monitoredItems, mon);
1481
0
    UA_MonitoringParameters_clear(&mon->parameters);
1482
0
    mon->parameters = mon->pendingParameters;
1483
0
    UA_MonitoringParameters_init(&mon->pendingParameters);
1484
0
    sub->pendingRekeys--;
1485
0
    EventFields_clear(&mon->eventFields);
1486
0
    ZIP_INSERT(MonitorItemsTree, &sub->monitoredItems, mon);
1487
0
    return mon;
1488
0
}
1489
1490
static void
1491
processDataChangeNotification(UA_Client *client, UA_Client_Subscription *sub,
1492
0
                              UA_DataChangeNotification *dataChangeNotification) {
1493
0
    UA_LOCK_ASSERT(&client->clientMutex);
1494
1495
0
    for(size_t j = 0; j < dataChangeNotification->monitoredItemsSize; ++j) {
1496
0
        UA_MonitoredItemNotification *min = &dataChangeNotification->monitoredItems[j];
1497
1498
        /* Find the MonitoredItem */
1499
0
        UA_Client_MonitoredItem *mon =
1500
0
            findMonitoredItemForNotification(sub, min->clientHandle);
1501
1502
0
        if(!mon) {
1503
0
            UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1504
0
                           "Could not process a notification with clienthandle %" PRIu32
1505
0
                           " on subscription %" PRIu32, min->clientHandle, sub->subscriptionId);
1506
0
            continue;
1507
0
        }
1508
1509
0
        if(mon->isEventMonitoredItem) {
1510
0
            UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1511
0
                           "MonitoredItem is configured for Events. But received a "
1512
0
                           "DataChangeNotification.");
1513
0
            continue;
1514
0
        }
1515
1516
0
        if(mon->handler.dataChangeCallback) {
1517
0
            void *subC = sub->context;
1518
0
            void *monC = mon->context;
1519
0
            UA_UInt32 subId = sub->subscriptionId;
1520
0
            UA_UInt32 monId = mon->monitoredItemId;
1521
0
            mon->handler.dataChangeCallback(client, subId, subC, monId, monC, &min->value);
1522
0
        }
1523
0
    }
1524
0
}
1525
1526
static void
1527
processEventNotification(UA_Client *client, UA_Client_Subscription *sub,
1528
0
                         UA_EventNotificationList *eventNotificationList) {
1529
0
    UA_LOCK_ASSERT(&client->clientMutex);
1530
1531
0
    for(size_t j = 0; j < eventNotificationList->eventsSize; ++j) {
1532
0
        UA_EventFieldList *efl = &eventNotificationList->events[j];
1533
1534
        /* Find the MonitoredItem */
1535
0
        UA_Client_MonitoredItem *mon =
1536
0
            findMonitoredItemForNotification(sub, efl->clientHandle);
1537
1538
0
        if(!mon) {
1539
0
            UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
1540
0
                         "Could not process a notification with clienthandle %" PRIu32
1541
0
                         " on subscription %" PRIu32, efl->clientHandle,
1542
0
                         sub->subscriptionId);
1543
0
            continue;
1544
0
        }
1545
1546
0
        if(!mon->isEventMonitoredItem) {
1547
0
            UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
1548
0
                         "MonitoredItem is configured for DataChanges. But received a "
1549
0
                         "EventNotification");
1550
0
            continue;
1551
0
        }
1552
1553
        /* Build the callback metadata from the active filter only when it is
1554
         * first needed. After a promotion the cache is empty and rebuilt from
1555
         * the newly active parameters. */
1556
0
        if(mon->eventFields.mapSize == 0) {
1557
0
            UA_StatusCode res =
1558
0
                prepareEventFieldsMap(&mon->eventFields, &mon->parameters);
1559
0
            if(res != UA_STATUSCODE_GOOD) {
1560
0
                EventFields_clear(&mon->eventFields);
1561
0
                UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1562
0
                               "Could not prepare event fields: %s",
1563
0
                               UA_StatusCode_name(res));
1564
0
                continue;
1565
0
            }
1566
0
        }
1567
1568
0
        if(mon->eventFields.mapSize != efl->eventFieldsSize) {
1569
0
            UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
1570
0
                         "MonitoredItem received a EventNotification with the "
1571
0
                          "wrong number of event fields");
1572
0
            continue;
1573
0
        }
1574
1575
        /* Add the borrowed notification values to the cached field names. */
1576
0
        for(size_t i = 0; i < mon->eventFields.mapSize; i++)
1577
0
            mon->eventFields.map[i].value = efl->eventFields[i];
1578
0
        mon->handler.eventCallback(client, sub->subscriptionId, sub->context,
1579
0
                                   mon->monitoredItemId, mon->context,
1580
0
                                   mon->eventFields);
1581
0
    }
1582
0
}
1583
1584
static void
1585
processNotificationMessage(UA_Client *client, UA_Client_Subscription *sub,
1586
0
                           UA_ExtensionObject *msg) {
1587
0
    UA_LOCK_ASSERT(&client->clientMutex);
1588
1589
0
    if(msg->encoding != UA_EXTENSIONOBJECT_DECODED)
1590
0
        return;
1591
1592
    /* Handle DataChangeNotification */
1593
0
    if(msg->content.decoded.type == &UA_TYPES[UA_TYPES_DATACHANGENOTIFICATION]) {
1594
0
        UA_DataChangeNotification *dataChangeNotification =
1595
0
            (UA_DataChangeNotification *)msg->content.decoded.data;
1596
0
        processDataChangeNotification(client, sub, dataChangeNotification);
1597
0
        return;
1598
0
    }
1599
1600
    /* Handle EventNotification */
1601
0
    if(msg->content.decoded.type == &UA_TYPES[UA_TYPES_EVENTNOTIFICATIONLIST]) {
1602
0
        UA_EventNotificationList *eventNotificationList =
1603
0
            (UA_EventNotificationList *)msg->content.decoded.data;
1604
0
        processEventNotification(client, sub, eventNotificationList);
1605
0
        return;
1606
0
    }
1607
1608
    /* Handle StatusChangeNotification */
1609
0
    if(msg->content.decoded.type == &UA_TYPES[UA_TYPES_STATUSCHANGENOTIFICATION]) {
1610
0
        if(sub->statusChangeCallback) {
1611
0
            void *subC = sub->context;
1612
0
            UA_UInt32 subId = sub->subscriptionId;
1613
0
            sub->statusChangeCallback(client, subId, subC,
1614
0
                                      (UA_StatusChangeNotification*)msg->content.decoded.data);
1615
0
        } else {
1616
0
            UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1617
0
                           "Dropped a StatusChangeNotification since no "
1618
0
                           "callback is registered");
1619
0
        }
1620
0
        return;
1621
0
    }
1622
1623
0
    UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1624
0
                   "Unknown notification message type");
1625
0
}
1626
1627
void
1628
__Client_Subscriptions_processPublishResponse(UA_Client *client, UA_PublishRequest *request,
1629
0
                                              UA_PublishResponse *response) {
1630
0
    UA_LOCK_ASSERT(&client->clientMutex);
1631
1632
    /* Reduce the number of "in-flight" PublishRequests */
1633
0
    client->currentlyOutStandingPublishRequests--;
1634
1635
    /* Process ServiceResult for bad StatusCodes without referring to a
1636
     * SubscriptionId */
1637
0
    switch(response->responseHeader.serviceResult) {
1638
0
    case UA_STATUSCODE_BADTOOMANYPUBLISHREQUESTS:
1639
        /* Correct the assumed number of max outstanding requests */
1640
0
        if(client->config.outStandingPublishRequests > 1) {
1641
0
            client->config.outStandingPublishRequests--;
1642
0
            UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1643
0
                           "PublishResponse: Too many PublishRequest, reduce "
1644
0
                           "outStandingPublishRequests to %" PRId16,
1645
0
                           client->config.outStandingPublishRequests);
1646
0
        } else {
1647
0
            UA_LOG_ERROR(client->config.logging, UA_LOGCATEGORY_CLIENT,
1648
0
                         "PublishResponse: Too many PublishRequests when "
1649
0
                         "outStandingPublishRequests = 1");
1650
0
            UA_Client_Subscriptions_deleteSingle(client, response->subscriptionId);
1651
0
        }
1652
0
        return;
1653
1654
0
    case UA_STATUSCODE_BADSHUTDOWN:
1655
        /* If the remote server shuts down, DEBUG-log to avoid a warning-storm
1656
         * for normal operations */
1657
0
        UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
1658
0
                     "PublishResponse: Received BadShutdown status");
1659
0
        return;
1660
1661
0
    case UA_STATUSCODE_BADNOSUBSCRIPTION:
1662
        /* There is no Subscription configured, the server expects no
1663
         * PublishRequests. We demote this to debug-logging, as it can occur
1664
         * during regular shutdown when the Subscriptions are removed. */
1665
0
        UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
1666
0
                     "PublishResponse: Received BadNoSubscription status");
1667
0
        return;
1668
1669
0
    default:
1670
0
        break;
1671
0
    }
1672
1673
    /* Get the Subscription */
1674
0
    UA_Client_Subscription *sub = findSubscriptionById(client, response->subscriptionId);
1675
0
    if(!sub) {
1676
0
        response->responseHeader.serviceResult = UA_STATUSCODE_BADNOSUBSCRIPTION;
1677
0
        UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1678
0
                       "PublishResponse: Received response for an unknown Subscription");
1679
0
        return;
1680
0
    }
1681
1682
    /* Process ServiceResult */
1683
0
    switch(response->responseHeader.serviceResult) {
1684
0
    case UA_STATUSCODE_BADSESSIONCLOSED:
1685
        /* The Session no longer exists on the server - remove the Subscription */
1686
0
        __Client_Subscription_deleteInternal(client, sub);
1687
0
        return;
1688
1689
0
    case UA_STATUSCODE_BADTIMEOUT:
1690
0
        UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1691
0
                       "PublishResponse: Aborted with BadTimeout status");
1692
0
        if(client->config.subscriptionInactivityCallback) {
1693
0
            void *subC = sub->context;
1694
0
            UA_UInt32 subId = sub->subscriptionId;
1695
0
            client->config.subscriptionInactivityCallback(client, subId, subC);
1696
0
        }
1697
0
        return;
1698
1699
0
    case UA_STATUSCODE_GOOD:
1700
0
        break; /* Continue below */
1701
1702
0
    default:
1703
        /* Catch-all for other bad StatusCodes */
1704
0
        UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1705
0
                       "PublishResponse: Received %s status",
1706
0
                       UA_StatusCode_name(response->responseHeader.serviceResult));
1707
0
        return;
1708
0
    }
1709
1710
    /* Update the LastActivity for the Subscription */
1711
0
    UA_EventLoop *el = client->config.eventLoop;
1712
0
    sub->lastActivity = el->dateTime_nowMonotonic(el);
1713
1714
    /* Detect missing message - OPC Unified Architecture, Part 4 5.13.1.1 e) */
1715
0
    UA_NotificationMessage *msg = &response->notificationMessage;
1716
0
    if(__nextSequenceNumber(sub->sequenceNumber) != msg->sequenceNumber) {
1717
0
        UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1718
0
                       "PublishResponse: Invalid subscription sequence number: "
1719
0
                       "Expected %" PRIu32 " but got %" PRIu32,
1720
0
                       __nextSequenceNumber(sub->sequenceNumber), msg->sequenceNumber);
1721
        /* This is an error. But we do not abort the connection. Some server
1722
         * SDKs misbehave from time to time and send out-of-order sequence
1723
         * numbers. (Probably some multi-threading synchronization issue.) */
1724
        /* UA_Client_disconnect(client);
1725
           return; */
1726
0
    }
1727
1728
    /* According to f), a keep-alive message contains no notifications and has
1729
     * the sequence number of the next NotificationMessage that is to be sent =>
1730
     * More than one consecutive keep-alive message or a NotificationMessage
1731
     * following a keep-alive message will share the same sequence number. */
1732
0
    if(msg->notificationDataSize)
1733
0
        sub->sequenceNumber = msg->sequenceNumber;
1734
1735
    /* Process the notification messages */
1736
0
    for(size_t k = 0; k < msg->notificationDataSize; ++k)
1737
0
        processNotificationMessage(client, sub, &msg->notificationData[k]);
1738
1739
    /* Add the current NotificationMessage (SequenceNumber) to the list of
1740
     * pending acks to be acknowledged. But only if it is in the list of
1741
     * sequence numbers the server has available. */
1742
0
    for(size_t i = 0; i < response->availableSequenceNumbersSize; i++) {
1743
0
        if(response->availableSequenceNumbers[i] != msg->sequenceNumber)
1744
0
            continue;
1745
0
        UA_Client_NotificationsAckNumber *tmpAck = (UA_Client_NotificationsAckNumber*)
1746
0
            UA_malloc(sizeof(UA_Client_NotificationsAckNumber));
1747
0
        if(!tmpAck) {
1748
0
            UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1749
0
                           "PublishResponse: Not enough memory to store the pending "
1750
0
                           "acknowledgement for Subscription %" PRIu32, sub->subscriptionId);
1751
0
            break;
1752
0
        }
1753
0
        tmpAck->subAck.sequenceNumber = msg->sequenceNumber;
1754
0
        tmpAck->subAck.subscriptionId = sub->subscriptionId;
1755
0
        LIST_INSERT_HEAD(&client->pendingNotificationsAcks, tmpAck, listEntry);
1756
0
        break;
1757
0
    }
1758
0
}
1759
1760
static void
1761
processPublishResponseAsync(UA_Client *client, void *userdata,
1762
0
                            UA_UInt32 requestId, void *response) {
1763
0
    UA_PublishRequest *req = (UA_PublishRequest*)userdata;
1764
0
    UA_PublishResponse *res = (UA_PublishResponse*)response;
1765
1766
0
    lockClient(client);
1767
1768
    /* Process the response */
1769
0
    __Client_Subscriptions_processPublishResponse(client, req, res);
1770
1771
    /* Delete the cached request */
1772
0
    UA_PublishRequest_delete(req);
1773
1774
    /* Fill up the outstanding publish requests */
1775
0
    __Client_Subscriptions_backgroundPublish(client);
1776
1777
0
    unlockClient(client);
1778
0
}
1779
1780
void
1781
120
__Client_Subscriptions_clear(UA_Client *client) {
1782
120
    UA_Client_NotificationsAckNumber *n;
1783
120
    UA_Client_NotificationsAckNumber *tmp;
1784
120
    LIST_FOREACH_SAFE(n, &client->pendingNotificationsAcks, listEntry, tmp) {
1785
0
        LIST_REMOVE(n, listEntry);
1786
0
        UA_free(n);
1787
0
    }
1788
1789
120
    UA_Client_Subscription *sub;
1790
120
    UA_Client_Subscription *tmps;
1791
120
    LIST_FOREACH_SAFE(sub, &client->subscriptions, listEntry, tmps)
1792
0
        __Client_Subscription_deleteInternal(client, sub); /* force local removal */
1793
1794
120
    client->monitoredItemHandles = 0;
1795
120
}
1796
1797
void
1798
0
__Client_Subscriptions_backgroundPublishInactivityCheck(UA_Client *client) {
1799
0
    UA_LOCK_ASSERT(&client->clientMutex);
1800
1801
0
    UA_EventLoop *el = client->config.eventLoop;
1802
0
    UA_DateTime nowm = el->dateTime_nowMonotonic(el);
1803
1804
0
    UA_Client_Subscription *sub;
1805
0
    LIST_FOREACH(sub, &client->subscriptions, listEntry) {
1806
0
        UA_DateTime maxSilence = (UA_DateTime)
1807
0
            ((sub->publishingInterval * sub->maxKeepAliveCount) +
1808
0
             client->config.timeout) * UA_DATETIME_MSEC;
1809
0
        if(maxSilence + sub->lastActivity < nowm) {
1810
            /* Reset activity */
1811
0
            sub->lastActivity = nowm;
1812
1813
0
            if(client->config.subscriptionInactivityCallback) {
1814
0
                void *subC = sub->context;
1815
0
                UA_UInt32 subId = sub->subscriptionId;
1816
0
                client->config.subscriptionInactivityCallback(client, subId, subC);
1817
0
            }
1818
0
            UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
1819
0
                           "Inactivity for Subscription %" PRIu32 ".",
1820
0
                           sub->subscriptionId);
1821
0
        }
1822
0
    }
1823
0
}
1824
1825
void
1826
0
__Client_Subscriptions_backgroundPublish(UA_Client *client) {
1827
0
    UA_LOCK_ASSERT(&client->clientMutex);
1828
1829
0
    if(client->sessionState != UA_SESSIONSTATE_ACTIVATED)
1830
0
        return;
1831
1832
    /* The session must have at least one subscription */
1833
0
    if(!LIST_FIRST(&client->subscriptions))
1834
0
        return;
1835
1836
0
    while(client->currentlyOutStandingPublishRequests < client->config.outStandingPublishRequests) {
1837
0
        UA_PublishRequest *request = UA_PublishRequest_new();
1838
0
        if(!request)
1839
0
            return;
1840
1841
        /* Publish requests are valid for 10 minutes */
1842
0
        request->requestHeader.timeoutHint = 10 * 60 * 1000;
1843
1844
0
        UA_StatusCode retval = __Client_preparePublishRequest(client, request);
1845
0
        if(retval != UA_STATUSCODE_GOOD) {
1846
0
            UA_PublishRequest_delete(request);
1847
0
            return;
1848
0
        }
1849
1850
0
        retval = __Client_AsyncService(client, request,
1851
0
                                         &UA_TYPES[UA_TYPES_PUBLISHREQUEST],
1852
0
                                         processPublishResponseAsync,
1853
0
                                         &UA_TYPES[UA_TYPES_PUBLISHRESPONSE],
1854
0
                                         (void*)request, NULL);
1855
0
        if(retval != UA_STATUSCODE_GOOD) {
1856
0
            UA_PublishRequest_delete(request);
1857
0
            return;
1858
0
        }
1859
1860
0
        client->currentlyOutStandingPublishRequests++;
1861
0
    }
1862
0
}
1863
1864
UA_SetPublishingModeResponse
1865
UA_Client_Subscriptions_setPublishingMode(UA_Client *client,
1866
0
                                          const UA_SetPublishingModeRequest request) {
1867
0
    UA_SetPublishingModeResponse response;
1868
0
    __UA_Client_Service(client, &request, &UA_TYPES[UA_TYPES_SETPUBLISHINGMODEREQUEST],
1869
0
                        &response, &UA_TYPES[UA_TYPES_SETPUBLISHINGMODERESPONSE]);
1870
0
    return response;
1871
0
}
1872
1873
UA_SetMonitoringModeResponse
1874
UA_Client_MonitoredItems_setMonitoringMode(UA_Client *client,
1875
0
                                           const UA_SetMonitoringModeRequest request) {
1876
0
    UA_SetMonitoringModeResponse response;
1877
0
    __UA_Client_Service(client, &request, &UA_TYPES[UA_TYPES_SETMONITORINGMODEREQUEST],
1878
0
                        &response, &UA_TYPES[UA_TYPES_SETMONITORINGMODERESPONSE]);
1879
0
    return response;
1880
0
}
1881
1882
UA_StatusCode
1883
UA_Client_MonitoredItems_setMonitoringMode_async(UA_Client *client,
1884
                                                 const UA_SetMonitoringModeRequest request,
1885
                                                 UA_ClientAsyncSetMonitoringModeCallback callback,
1886
0
                                                 void *userdata, UA_UInt32 *requestId) {
1887
0
    return __UA_Client_AsyncService(client, &request,
1888
0
                                    &UA_TYPES[UA_TYPES_SETMONITORINGMODEREQUEST],
1889
0
                                    (UA_ClientAsyncServiceCallback)callback,
1890
0
                                    &UA_TYPES[UA_TYPES_SETMONITORINGMODERESPONSE],
1891
0
                                    userdata, requestId);
1892
0
}
1893
1894
UA_SetTriggeringResponse
1895
UA_Client_MonitoredItems_setTriggering(UA_Client *client,
1896
0
                                       const UA_SetTriggeringRequest request) {
1897
0
    UA_SetTriggeringResponse response;
1898
0
    __UA_Client_Service(client, &request, &UA_TYPES[UA_TYPES_SETTRIGGERINGREQUEST],
1899
0
                        &response, &UA_TYPES[UA_TYPES_SETTRIGGERINGRESPONSE]);
1900
0
    return response;
1901
0
}
1902
1903
UA_StatusCode
1904
UA_Client_MonitoredItems_setTriggering_async(UA_Client *client,
1905
                                             const UA_SetTriggeringRequest request,
1906
                                             UA_ClientAsyncSetTriggeringCallback callback,
1907
0
                                             void *userdata, UA_UInt32 *requestId) {
1908
0
    return __UA_Client_AsyncService(client, &request,
1909
0
                                    &UA_TYPES[UA_TYPES_SETTRIGGERINGREQUEST],
1910
0
                                    (UA_ClientAsyncServiceCallback)callback,
1911
0
                                    &UA_TYPES[UA_TYPES_SETTRIGGERINGRESPONSE],
1912
0
                                    userdata, requestId);
1913
0
}