/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 = ¶ms->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 | } |