/src/open62541_15/src/server/ua_services_subscription.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 2014-2018, 2022 (c) Fraunhofer IOSB (Author: Julius Pfrommer) |
6 | | * Copyright 2016-2017 (c) Florian Palm |
7 | | * Copyright 2015 (c) Chris Iatrou |
8 | | * Copyright 2015-2016 (c) Sten GrĂ¼ner |
9 | | * Copyright 2015-2016 (c) Oleksiy Vasylyev |
10 | | * Copyright 2017 (c) Stefan Profanter, fortiss GmbH |
11 | | * Copyright 2018 (c) Ari Breitkreuz, fortiss GmbH |
12 | | * Copyright 2017 (c) Mattias Bornhager |
13 | | * Copyright 2017 (c) Henrik Norrman |
14 | | * Copyright 2017-2018 (c) Thomas Stalder, Blue Time Concept SA |
15 | | * Copyright 2018 (c) Fabian Arndt, Root-Core |
16 | | * Copyright 2017-2019 (c) HMS Industrial Networks AB (Author: Jonas Green) |
17 | | * Copyright 2026 (c) o6 Automation GmbH (Author: Andreas Ebner) |
18 | | */ |
19 | | |
20 | | #include "ua_server_internal.h" |
21 | | #include "ua_services.h" |
22 | | #include "ua_subscription.h" |
23 | | |
24 | | #ifdef UA_ENABLE_SUBSCRIPTIONS /* conditional compilation */ |
25 | | |
26 | | static void |
27 | | setSubscriptionSettings(UA_Server *server, UA_Subscription *subscription, |
28 | | UA_Double requestedPublishingInterval, |
29 | | UA_UInt32 requestedLifetimeCount, |
30 | | UA_UInt32 requestedMaxKeepAliveCount, |
31 | | UA_UInt32 maxNotificationsPerPublish, |
32 | 243 | UA_Byte priority) { |
33 | 243 | UA_LOCK_ASSERT(&server->serviceMutex); |
34 | | |
35 | | /* re-parameterize the subscription */ |
36 | 243 | UA_BOUNDEDVALUE_SETWBOUNDS(server->config.publishingIntervalLimits, |
37 | 243 | requestedPublishingInterval, |
38 | 243 | subscription->publishingInterval); |
39 | | /* check for nan*/ |
40 | 243 | if(requestedPublishingInterval != requestedPublishingInterval) |
41 | 22 | subscription->publishingInterval = server->config.publishingIntervalLimits.min; |
42 | 243 | UA_BOUNDEDVALUE_SETWBOUNDS(server->config.keepAliveCountLimits, |
43 | 243 | requestedMaxKeepAliveCount, subscription->maxKeepAliveCount); |
44 | 243 | UA_BOUNDEDVALUE_SETWBOUNDS(server->config.lifeTimeCountLimits, |
45 | 243 | requestedLifetimeCount, subscription->lifeTimeCount); |
46 | 243 | if(subscription->lifeTimeCount < 3 * subscription->maxKeepAliveCount) |
47 | 30 | subscription->lifeTimeCount = 3 * subscription->maxKeepAliveCount; |
48 | 243 | subscription->notificationsPerPublish = maxNotificationsPerPublish; |
49 | 243 | if(maxNotificationsPerPublish == 0 || |
50 | 230 | maxNotificationsPerPublish > server->config.maxNotificationsPerPublish) |
51 | 236 | subscription->notificationsPerPublish = server->config.maxNotificationsPerPublish; |
52 | 243 | subscription->priority = priority; |
53 | 243 | } |
54 | | |
55 | | static void |
56 | | notifySubscription(UA_Server *server, UA_Subscription *sub, |
57 | 120 | UA_ApplicationNotificationType type) { |
58 | 120 | if(!server->config.subscriptionNotificationCallback && |
59 | 120 | !server->config.globalNotificationCallback) |
60 | 120 | return; |
61 | | |
62 | 0 | static UA_THREAD_LOCAL UA_KeyValuePair createSubData[8] = { |
63 | 0 | {{0, UA_STRING_STATIC("session-id")}, {0}}, |
64 | 0 | {{0, UA_STRING_STATIC("subscription-id")}, {0}}, |
65 | 0 | {{0, UA_STRING_STATIC("publishing-interval")}, {0}}, |
66 | 0 | {{0, UA_STRING_STATIC("lifetime-count")}, {0}}, |
67 | 0 | {{0, UA_STRING_STATIC("max-keepalive-count")}, {0}}, |
68 | 0 | {{0, UA_STRING_STATIC("max-notifications-per-publish")}, {0}}, |
69 | 0 | {{0, UA_STRING_STATIC("priority")}, {0}}, |
70 | 0 | {{0, UA_STRING_STATIC("publishing-enabled")}, {0}} |
71 | 0 | }; |
72 | 0 | UA_KeyValueMap createSubMap = {8, createSubData}; |
73 | |
|
74 | 0 | static UA_THREAD_LOCAL UA_NodeId sessionId; |
75 | 0 | sessionId = (sub->session) ? sub->session->sessionId : UA_NODEID_NULL; |
76 | 0 | static UA_THREAD_LOCAL UA_Boolean enabled; |
77 | 0 | enabled = (sub->state == UA_SUBSCRIPTIONSTATE_ENABLED); |
78 | |
|
79 | 0 | UA_Variant_setScalar(&createSubData[0].value, &sessionId, |
80 | 0 | &UA_TYPES[UA_TYPES_NODEID]); |
81 | 0 | UA_Variant_setScalar(&createSubData[1].value, &sub->subscriptionId, |
82 | 0 | &UA_TYPES[UA_TYPES_UINT32]); |
83 | 0 | UA_Variant_setScalar(&createSubData[2].value, &sub->publishingInterval, |
84 | 0 | &UA_TYPES[UA_TYPES_DOUBLE]); |
85 | 0 | UA_Variant_setScalar(&createSubData[3].value, &sub->lifeTimeCount, |
86 | 0 | &UA_TYPES[UA_TYPES_UINT32]); |
87 | 0 | UA_Variant_setScalar(&createSubData[4].value, &sub->maxKeepAliveCount, |
88 | 0 | &UA_TYPES[UA_TYPES_UINT32]); |
89 | 0 | UA_Variant_setScalar(&createSubData[5].value, &sub->notificationsPerPublish, |
90 | 0 | &UA_TYPES[UA_TYPES_UINT32]); |
91 | 0 | UA_Variant_setScalar(&createSubData[6].value, &sub->priority, |
92 | 0 | &UA_TYPES[UA_TYPES_BYTE]); |
93 | 0 | UA_Variant_setScalar(&createSubData[7].value, &enabled, |
94 | 0 | &UA_TYPES[UA_TYPES_BOOLEAN]); |
95 | |
|
96 | 0 | if(server->config.subscriptionNotificationCallback) |
97 | 0 | server->config.subscriptionNotificationCallback(server, type, createSubMap); |
98 | 0 | if(server->config.globalNotificationCallback) |
99 | 0 | server->config.globalNotificationCallback(server, type, createSubMap); |
100 | 0 | } |
101 | | |
102 | | UA_Boolean |
103 | | Service_CreateSubscription(UA_Server *server, UA_Session *session, |
104 | | const UA_CreateSubscriptionRequest *request, |
105 | 120 | UA_CreateSubscriptionResponse *response) { |
106 | 120 | UA_LOCK_ASSERT(&server->serviceMutex); |
107 | | |
108 | | /* Check limits for the number of subscriptions */ |
109 | 120 | if(((server->config.maxSubscriptions != 0) && |
110 | 0 | (server->subscriptionsSize >= server->config.maxSubscriptions)) || |
111 | 120 | ((server->config.maxSubscriptionsPerSession != 0) && |
112 | 0 | (session->subscriptionsSize >= server->config.maxSubscriptionsPerSession))) { |
113 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADTOOMANYSUBSCRIPTIONS; |
114 | 0 | return true; |
115 | 0 | } |
116 | | |
117 | | /* Create the subscription */ |
118 | 120 | UA_Subscription *sub = UA_Subscription_new(); |
119 | 120 | if(!sub) { |
120 | 0 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
121 | 0 | "Processing CreateSubscriptionRequest failed"); |
122 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADOUTOFMEMORY; |
123 | 0 | return true; |
124 | 0 | } |
125 | | |
126 | | /* Set the subscription parameters */ |
127 | 120 | setSubscriptionSettings(server, sub, request->requestedPublishingInterval, |
128 | 120 | request->requestedLifetimeCount, |
129 | 120 | request->requestedMaxKeepAliveCount, |
130 | 120 | request->maxNotificationsPerPublish, request->priority); |
131 | 120 | sub->subscriptionId = ++server->lastSubscriptionId; /* Assign the SubscriptionId */ |
132 | | |
133 | | /* Register the subscription in the server */ |
134 | 120 | LIST_INSERT_HEAD(&server->subscriptions, sub, serverListEntry); |
135 | 120 | server->subscriptionsSize++; |
136 | | |
137 | | /* Update the server statistics */ |
138 | 120 | server->serverDiagnosticsSummary.currentSubscriptionCount++; |
139 | 120 | server->serverDiagnosticsSummary.cumulatedSubscriptionCount++; |
140 | | |
141 | | /* Attach the Subscription to the session */ |
142 | 120 | UA_Session_attachSubscription(session, sub); |
143 | | |
144 | | /* Create representation in the Session object */ |
145 | 120 | #ifdef UA_ENABLE_DIAGNOSTICS |
146 | 120 | createSubscriptionObject(server, session, sub); |
147 | 120 | #endif |
148 | | |
149 | | /* Set the subscription state. This also registers the callback. |
150 | | * Note that also a disabled subscription publishes keepalives. */ |
151 | 120 | UA_SubscriptionState sState = (request->publishingEnabled) ? |
152 | 81 | UA_SUBSCRIPTIONSTATE_ENABLED : UA_SUBSCRIPTIONSTATE_ENABLED_NOPUBLISH; |
153 | 120 | UA_StatusCode res = Subscription_setState(server, sub, sState); |
154 | 120 | if(res != UA_STATUSCODE_GOOD) { |
155 | 0 | UA_LOG_DEBUG_SESSION(server->config.logging, sub->session, |
156 | 0 | "Subscription %" PRIu32 " | Could not register " |
157 | 0 | "publish callback with error code %s", |
158 | 0 | sub->subscriptionId, UA_StatusCode_name(res)); |
159 | 0 | response->responseHeader.serviceResult = res; |
160 | 0 | UA_Subscription_delete(server, sub); |
161 | 0 | return true; |
162 | 0 | } |
163 | | |
164 | 120 | UA_LOG_INFO_SUBSCRIPTION(server->config.logging, sub, |
165 | 120 | "Subscription created (Publishing interval %.2fms, " |
166 | 120 | "max %lu notifications per publish)", |
167 | 120 | sub->publishingInterval, |
168 | 120 | (long unsigned)sub->notificationsPerPublish); |
169 | | |
170 | | /* Notify the application */ |
171 | 120 | notifySubscription(server, sub, |
172 | 120 | UA_APPLICATIONNOTIFICATIONTYPE_SUBSCRIPTION_CREATED); |
173 | | |
174 | | /* Prepare the response */ |
175 | 120 | response->subscriptionId = sub->subscriptionId; |
176 | 120 | response->revisedPublishingInterval = sub->publishingInterval; |
177 | 120 | response->revisedLifetimeCount = sub->lifeTimeCount; |
178 | 120 | response->revisedMaxKeepAliveCount = sub->maxKeepAliveCount; |
179 | | |
180 | 120 | return true; |
181 | 120 | } |
182 | | |
183 | | struct UpdateSamplingContext { |
184 | | UA_Server *server; |
185 | | UA_Double oldPublishingInterval; |
186 | | }; |
187 | | |
188 | | static void * |
189 | 0 | updateSamplingIntervalVisitor(void *context, UA_MonitoredItem *mon) { |
190 | 0 | struct UpdateSamplingContext *ctx = |
191 | 0 | (struct UpdateSamplingContext*)context; |
192 | |
|
193 | 0 | if(mon->parameters.samplingInterval == mon->subscription->publishingInterval || |
194 | 0 | mon->parameters.samplingInterval == ctx->oldPublishingInterval) { |
195 | 0 | UA_MonitoredItem_unregisterSampling(ctx->server, mon); |
196 | 0 | UA_MonitoredItem_registerSampling(ctx->server, mon); |
197 | 0 | } |
198 | |
|
199 | 0 | return NULL; |
200 | 0 | } |
201 | | |
202 | | UA_Boolean |
203 | | Service_ModifySubscription(UA_Server *server, UA_Session *session, |
204 | | const UA_ModifySubscriptionRequest *request, |
205 | 8 | UA_ModifySubscriptionResponse *response) { |
206 | 8 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
207 | 8 | "Processing ModifySubscriptionRequest"); |
208 | 8 | UA_LOCK_ASSERT(&server->serviceMutex); |
209 | | |
210 | 8 | UA_Subscription *sub = UA_Session_getSubscriptionById(session, request->subscriptionId); |
211 | 8 | if(!sub) { |
212 | 8 | response->responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID; |
213 | 8 | return true; |
214 | 8 | } |
215 | | |
216 | | /* Store the old publishing interval */ |
217 | 0 | UA_Double oldPublishingInterval = sub->publishingInterval; |
218 | 0 | UA_Byte oldPriority = sub->priority; |
219 | | |
220 | | /* Change the Subscription settings */ |
221 | 0 | setSubscriptionSettings(server, sub, request->requestedPublishingInterval, |
222 | 0 | request->requestedLifetimeCount, |
223 | 0 | request->requestedMaxKeepAliveCount, |
224 | 0 | request->maxNotificationsPerPublish, request->priority); |
225 | | |
226 | | /* Reset the subscription lifetime */ |
227 | 0 | Subscription_resetLifetime(sub); |
228 | | |
229 | | /* The publish interval has changed */ |
230 | 0 | if(sub->publishingInterval != oldPublishingInterval) { |
231 | | /* Change the repeated callback to the new interval. This cannot fail as |
232 | | * memory is reused. */ |
233 | 0 | if(sub->publishCallbackId > 0) |
234 | 0 | changeRepeatedCallbackInterval(server, sub->publishCallbackId, |
235 | 0 | sub->publishingInterval); |
236 | | |
237 | | /* For each MonitoredItem check if it was/shall be attached to the |
238 | | * publish interval. This ensures that we have less cyclic callbacks |
239 | | * registered and that the notifications are fresh. */ |
240 | 0 | struct UpdateSamplingContext ctx; |
241 | 0 | ctx.server = server; |
242 | 0 | ctx.oldPublishingInterval = oldPublishingInterval; |
243 | |
|
244 | 0 | ZIP_ITER(UA_MonitoredItemIdTree, &sub->monitoredItemsById, |
245 | 0 | updateSamplingIntervalVisitor, &ctx); |
246 | 0 | } |
247 | | |
248 | | /* If the priority has changed, re-enter the subscription to the |
249 | | * priority-ordered queue in the session. */ |
250 | 0 | if(oldPriority != sub->priority) { |
251 | 0 | UA_Session_detachSubscription(server, session, sub, false); |
252 | 0 | UA_Session_attachSubscription(session, sub); |
253 | 0 | } |
254 | | |
255 | | /* Notify the application */ |
256 | 0 | notifySubscription(server, sub, |
257 | 0 | UA_APPLICATIONNOTIFICATIONTYPE_SUBSCRIPTION_MODIFIED); |
258 | | |
259 | | /* Set the response */ |
260 | 0 | response->revisedPublishingInterval = sub->publishingInterval; |
261 | 0 | response->revisedLifetimeCount = sub->lifeTimeCount; |
262 | 0 | response->revisedMaxKeepAliveCount = sub->maxKeepAliveCount; |
263 | | |
264 | | /* Update the diagnostics statistics */ |
265 | 0 | #ifdef UA_ENABLE_DIAGNOSTICS |
266 | 0 | sub->modifyCount++; |
267 | 0 | #endif |
268 | |
|
269 | 0 | return true; |
270 | 8 | } |
271 | | |
272 | | static void |
273 | | Operation_SetPublishingMode(UA_Server *server, UA_Session *session, |
274 | | const UA_Boolean *publishingEnabled, |
275 | | const UA_UInt32 *subscriptionId, |
276 | 78.8k | UA_StatusCode *result) { |
277 | 78.8k | UA_LOCK_ASSERT(&server->serviceMutex); |
278 | 78.8k | UA_Subscription *sub = UA_Session_getSubscriptionById(session, *subscriptionId); |
279 | 78.8k | if(!sub) { |
280 | 78.8k | *result = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID; |
281 | 78.8k | return; |
282 | 78.8k | } |
283 | | |
284 | | /* Enable/disable */ |
285 | 0 | UA_SubscriptionState sState = (*publishingEnabled) ? |
286 | 0 | UA_SUBSCRIPTIONSTATE_ENABLED : UA_SUBSCRIPTIONSTATE_ENABLED_NOPUBLISH; |
287 | 0 | *result = Subscription_setState(server, sub, sState); |
288 | | |
289 | | /* Reset the lifetime counter */ |
290 | 0 | Subscription_resetLifetime(sub); |
291 | | |
292 | | /* Notify the application */ |
293 | 0 | notifySubscription(server, sub, |
294 | 0 | UA_APPLICATIONNOTIFICATIONTYPE_SUBSCRIPTION_PUBLISHINGMODE); |
295 | 0 | } |
296 | | |
297 | | UA_Boolean |
298 | | Service_SetPublishingMode(UA_Server *server, UA_Session *session, |
299 | | const UA_SetPublishingModeRequest *request, |
300 | 28 | UA_SetPublishingModeResponse *response) { |
301 | 28 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
302 | 28 | "Processing SetPublishingModeRequest"); |
303 | 28 | UA_LOCK_ASSERT(&server->serviceMutex); |
304 | | |
305 | 28 | UA_Boolean publishingEnabled = request->publishingEnabled; /* request is const */ |
306 | 28 | response->responseHeader.serviceResult = |
307 | 28 | allocProcessServiceOperations(server, session, |
308 | 28 | (UA_ServiceOperation)Operation_SetPublishingMode, |
309 | 28 | &publishingEnabled, &request->subscriptionIdsSize, |
310 | 28 | &UA_TYPES[UA_TYPES_UINT32], &response->resultsSize, |
311 | 28 | &UA_TYPES[UA_TYPES_STATUSCODE]); |
312 | | |
313 | 28 | return true; |
314 | 28 | } |
315 | | |
316 | | UA_Boolean |
317 | | Service_Publish(UA_Server *server, UA_Session *session, |
318 | | const UA_PublishRequest *request, |
319 | 1 | UA_PublishResponse *response) { |
320 | 1 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
321 | 1 | "Processing PublishRequest with RequestId %u", |
322 | 1 | server->asyncManager.currentRequestId); |
323 | 1 | UA_LOCK_ASSERT(&server->serviceMutex); |
324 | | |
325 | | /* Return an error if the session has no subscription */ |
326 | 1 | if(TAILQ_EMPTY(&session->subscriptions)) { |
327 | 1 | response->responseHeader.serviceResult = UA_STATUSCODE_BADNOSUBSCRIPTION; |
328 | 1 | return true; |
329 | 1 | } |
330 | | |
331 | | /* Handle too many subscriptions to free resources before trying to allocate |
332 | | * resources for the new publish request. If the limit has been reached the |
333 | | * oldest publish request are returned with an error message. */ |
334 | 0 | UA_Session_ensurePublishQueueSpace(server, session); |
335 | | |
336 | | /* Allocate the response to store it in the retransmission queue */ |
337 | 0 | UA_PublishResponseEntry *entry = (UA_PublishResponseEntry *) |
338 | 0 | UA_malloc(sizeof(UA_PublishResponseEntry)); |
339 | 0 | if(!entry) { |
340 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADOUTOFMEMORY; |
341 | 0 | return true; |
342 | 0 | } |
343 | | |
344 | | /* Prepare the response */ |
345 | 0 | UA_PublishResponse *entry_response = &entry->response; |
346 | 0 | UA_PublishResponse_init(entry_response); |
347 | | |
348 | | /* Allocate the results array to acknowledge the acknowledge */ |
349 | 0 | if(request->subscriptionAcknowledgementsSize > 0) { |
350 | 0 | entry_response->results = (UA_StatusCode *) |
351 | 0 | UA_Array_new(request->subscriptionAcknowledgementsSize, |
352 | 0 | &UA_TYPES[UA_TYPES_STATUSCODE]); |
353 | 0 | if(!entry_response->results) { |
354 | 0 | UA_free(entry); |
355 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADOUTOFMEMORY; |
356 | 0 | return true; |
357 | 0 | } |
358 | 0 | entry_response->resultsSize = request->subscriptionAcknowledgementsSize; |
359 | 0 | } |
360 | | |
361 | | /* <--- Async response from here on ---> */ |
362 | | |
363 | 0 | entry->requestId = server->asyncManager.currentRequestId; |
364 | 0 | entry_response->responseHeader.requestHandle = request->requestHeader.requestHandle; |
365 | | |
366 | | /* Delete Acknowledged Subscription Messages */ |
367 | 0 | for(size_t i = 0; i < request->subscriptionAcknowledgementsSize; ++i) { |
368 | 0 | UA_SubscriptionAcknowledgement *ack = &request->subscriptionAcknowledgements[i]; |
369 | 0 | UA_Subscription *sub = UA_Session_getSubscriptionById(session, ack->subscriptionId); |
370 | 0 | if(!sub) { |
371 | 0 | entry_response->results[i] = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID; |
372 | 0 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
373 | 0 | "Cannot process acknowledgements subscription %u" PRIu32, |
374 | 0 | ack->subscriptionId); |
375 | 0 | continue; |
376 | 0 | } |
377 | | /* Remove the acked transmission from the retransmission queue */ |
378 | 0 | entry_response->results[i] = |
379 | 0 | UA_Subscription_removeRetransmissionMessage(sub, ack->sequenceNumber); |
380 | 0 | } |
381 | | |
382 | | /* Set the maxTime if a timeout hint is defined */ |
383 | 0 | entry->maxTime = UA_INT64_MAX; |
384 | 0 | if(request->requestHeader.timeoutHint > 0) { |
385 | 0 | UA_EventLoop *el = server->config.eventLoop; |
386 | 0 | entry->maxTime = el->dateTime_nowMonotonic(el) + |
387 | 0 | (request->requestHeader.timeoutHint * UA_DATETIME_MSEC); |
388 | 0 | } |
389 | | |
390 | | /* Queue the publish response. It will be dequeued in a repeated publish |
391 | | * callback. This can also be triggered right now for a late |
392 | | * subscription. */ |
393 | 0 | UA_Session_queuePublishReq(session, entry, false); |
394 | 0 | UA_LOG_DEBUG_SESSION(server->config.logging, session, "Queued a publication message"); |
395 | | |
396 | | /* If there are late subscriptions, the new publish request is used to |
397 | | * answer them immediately. Late subscriptions with higher priority are |
398 | | * considered earlier. However, a single subscription that generates many |
399 | | * notifications must not "starve" other late subscriptions. Hence we move |
400 | | * it to the end of the queue for the subscriptions of that priority. */ |
401 | 0 | UA_Subscription *late, *late_tmp; |
402 | 0 | TAILQ_FOREACH_SAFE(late, &session->subscriptions, sessionListEntry, late_tmp) { |
403 | | /* Skip non-late subscriptions */ |
404 | 0 | if(!late->late) |
405 | 0 | continue; |
406 | | |
407 | | /* Call publish on the late subscription */ |
408 | 0 | UA_LOG_DEBUG_SUBSCRIPTION(server->config.logging, late, |
409 | 0 | "Send PublishResponse on a late subscription"); |
410 | 0 | UA_Subscription_publish(server, late); |
411 | | |
412 | | /* Skip re-insert if the subscription was deleted or deactivated during |
413 | | * _publish */ |
414 | 0 | if(late->state >= UA_SUBSCRIPTIONSTATE_ENABLED_NOPUBLISH) { |
415 | | /* Find the first element with smaller priority and insert before |
416 | | * that. If there is none, insert at the end of the queue. */ |
417 | 0 | UA_Subscription *after = TAILQ_NEXT(late, sessionListEntry); |
418 | 0 | while(after && after->priority >= late->priority) |
419 | 0 | after = TAILQ_NEXT(after, sessionListEntry); |
420 | 0 | TAILQ_REMOVE(&session->subscriptions, late, sessionListEntry); |
421 | 0 | if(after) |
422 | 0 | TAILQ_INSERT_BEFORE(after, late, sessionListEntry); |
423 | 0 | else |
424 | 0 | TAILQ_INSERT_TAIL(&session->subscriptions, late, sessionListEntry); |
425 | 0 | } |
426 | | |
427 | | /* Responses left in the queue? */ |
428 | 0 | if(session->responseQueueSize == 0) |
429 | 0 | break; |
430 | 0 | } |
431 | |
|
432 | 0 | return false; |
433 | 0 | } |
434 | | |
435 | | static void |
436 | | Operation_DeleteSubscription(UA_Server *server, UA_Session *session, void *_, |
437 | 231 | const UA_UInt32 *subscriptionId, UA_StatusCode *result) { |
438 | | /* Find the Subscription */ |
439 | 231 | UA_Subscription *sub = UA_Session_getSubscriptionById(session, *subscriptionId); |
440 | 231 | if(!sub) { |
441 | 231 | *result = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID; |
442 | 231 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
443 | 231 | "Deleting Subscription with Id %" PRIu32 |
444 | 231 | " failed with error code %s", |
445 | 231 | *subscriptionId, UA_StatusCode_name(*result)); |
446 | 231 | return; |
447 | 231 | } |
448 | | |
449 | | /* Notify the application */ |
450 | 0 | notifySubscription(server, sub, |
451 | 0 | UA_APPLICATIONNOTIFICATIONTYPE_SUBSCRIPTION_DELETED); |
452 | | |
453 | | /* Delete the Subscription */ |
454 | 0 | UA_Subscription_delete(server, sub); |
455 | 0 | *result = UA_STATUSCODE_GOOD; |
456 | |
|
457 | 0 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
458 | 0 | "Subscription %" PRIu32 " | Subscription deleted", |
459 | 0 | *subscriptionId); |
460 | 0 | } |
461 | | |
462 | | UA_Boolean |
463 | | Service_DeleteSubscriptions(UA_Server *server, UA_Session *session, |
464 | | const UA_DeleteSubscriptionsRequest *request, |
465 | 24 | UA_DeleteSubscriptionsResponse *response) { |
466 | 24 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
467 | 24 | "Processing DeleteSubscriptionsRequest"); |
468 | 24 | UA_LOCK_ASSERT(&server->serviceMutex); |
469 | | |
470 | 24 | response->responseHeader.serviceResult = |
471 | 24 | allocProcessServiceOperations(server, session, |
472 | 24 | (UA_ServiceOperation)Operation_DeleteSubscription, |
473 | 24 | NULL, &request->subscriptionIdsSize, |
474 | 24 | &UA_TYPES[UA_TYPES_UINT32], &response->resultsSize, |
475 | 24 | &UA_TYPES[UA_TYPES_STATUSCODE]); |
476 | 24 | return true; |
477 | 24 | } |
478 | | |
479 | | UA_Boolean |
480 | | Service_Republish(UA_Server *server, UA_Session *session, |
481 | | const UA_RepublishRequest *request, |
482 | 0 | UA_RepublishResponse *response) { |
483 | 0 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
484 | 0 | "Processing RepublishRequest"); |
485 | 0 | UA_LOCK_ASSERT(&server->serviceMutex); |
486 | | |
487 | | /* Get the subscription */ |
488 | 0 | UA_Subscription *sub = UA_Session_getSubscriptionById(session, request->subscriptionId); |
489 | 0 | if(!sub) { |
490 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID; |
491 | 0 | return true; |
492 | 0 | } |
493 | | |
494 | | /* Reset the lifetime counter */ |
495 | 0 | Subscription_resetLifetime(sub); |
496 | | |
497 | | /* Update the subscription statistics */ |
498 | 0 | #ifdef UA_ENABLE_DIAGNOSTICS |
499 | 0 | sub->republishRequestCount++; |
500 | 0 | #endif |
501 | | |
502 | | /* Find the notification in the retransmission queue */ |
503 | 0 | UA_NotificationMessageEntry *entry; |
504 | 0 | TAILQ_FOREACH(entry, &sub->retransmissionQueue, listEntry) { |
505 | 0 | if(entry->message.sequenceNumber == request->retransmitSequenceNumber) |
506 | 0 | break; |
507 | 0 | } |
508 | 0 | if(!entry) { |
509 | 0 | response->responseHeader.serviceResult = UA_STATUSCODE_BADMESSAGENOTAVAILABLE; |
510 | 0 | return true; |
511 | 0 | } |
512 | | |
513 | 0 | response->responseHeader.serviceResult = |
514 | 0 | UA_NotificationMessage_copy(&entry->message, &response->notificationMessage); |
515 | | |
516 | | /* Update the subscription statistics for the case where we return a message */ |
517 | 0 | #ifdef UA_ENABLE_DIAGNOSTICS |
518 | 0 | sub->republishMessageCount++; |
519 | 0 | #endif |
520 | |
|
521 | 0 | return true; |
522 | 0 | } |
523 | | |
524 | | static UA_StatusCode |
525 | 0 | setTransferredSequenceNumbers(const UA_Subscription *sub, UA_TransferResult *result) { |
526 | | /* Allocate memory */ |
527 | 0 | result->availableSequenceNumbers = (UA_UInt32*) |
528 | 0 | UA_Array_new(sub->retransmissionQueueSize, &UA_TYPES[UA_TYPES_UINT32]); |
529 | 0 | if(!result->availableSequenceNumbers) |
530 | 0 | return UA_STATUSCODE_BADOUTOFMEMORY; |
531 | 0 | result->availableSequenceNumbersSize = sub->retransmissionQueueSize; |
532 | | |
533 | | /* Copy over the sequence numbers */ |
534 | 0 | UA_NotificationMessageEntry *entry; |
535 | 0 | size_t i = 0; |
536 | 0 | TAILQ_FOREACH(entry, &sub->retransmissionQueue, listEntry) { |
537 | 0 | result->availableSequenceNumbers[i] = entry->message.sequenceNumber; |
538 | 0 | i++; |
539 | 0 | } |
540 | |
|
541 | 0 | UA_assert(i == result->availableSequenceNumbersSize); |
542 | | |
543 | 0 | return UA_STATUSCODE_GOOD; |
544 | 0 | } |
545 | | |
546 | | static void * |
547 | 0 | setMonitoredItemSubscriptionVisitor(void *context, UA_MonitoredItem *mon) { |
548 | 0 | mon->subscription = (UA_Subscription*)context; |
549 | 0 | return NULL; |
550 | 0 | } |
551 | | |
552 | | static void |
553 | | Operation_TransferSubscription(UA_Server *server, UA_Session *session, |
554 | | const UA_Boolean *sendInitialValues, |
555 | | const UA_UInt32 *subscriptionId, |
556 | 90 | UA_TransferResult *result) { |
557 | 90 | UA_LOCK_ASSERT(&server->serviceMutex); |
558 | | |
559 | | /* Get the subscription. This requires a server-wide lookup instead of the |
560 | | * usual session-wide lookup. */ |
561 | 90 | UA_Subscription *sub = getSubscriptionById(server, *subscriptionId); |
562 | 90 | if(!sub) { |
563 | 90 | result->statusCode = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID; |
564 | 90 | return; |
565 | 90 | } |
566 | | |
567 | | /* Update the diagnostics statistics */ |
568 | 0 | #ifdef UA_ENABLE_DIAGNOSTICS |
569 | 0 | sub->transferRequestCount++; |
570 | 0 | #endif |
571 | | |
572 | | /* Is this the same session? Return the sequence numbers and do nothing else. */ |
573 | 0 | UA_Session *oldSession = sub->session; |
574 | 0 | if(oldSession == session) { |
575 | 0 | result->statusCode = setTransferredSequenceNumbers(sub, result); |
576 | 0 | #ifdef UA_ENABLE_DIAGNOSTICS |
577 | 0 | sub->transferredToSameClientCount++; |
578 | 0 | #endif |
579 | 0 | return; |
580 | 0 | } |
581 | | |
582 | | /* Check with AccessControl if the transfer is allowed */ |
583 | 0 | if(server->config.accessControl.allowTransferSubscription) { |
584 | 0 | if(!server->config.accessControl. |
585 | 0 | allowTransferSubscription(server, &server->config.accessControl, |
586 | 0 | oldSession ? &oldSession->sessionId : NULL, |
587 | 0 | oldSession ? oldSession->context : NULL, |
588 | 0 | &session->sessionId, session->context)) { |
589 | 0 | result->statusCode = UA_STATUSCODE_BADUSERACCESSDENIED; |
590 | 0 | return; |
591 | 0 | } |
592 | 0 | } else { |
593 | 0 | result->statusCode = UA_STATUSCODE_BADUSERACCESSDENIED; |
594 | 0 | return; |
595 | 0 | } |
596 | | |
597 | | /* Check limits for the number of subscriptions for this Session */ |
598 | 0 | if((server->config.maxSubscriptionsPerSession != 0) && |
599 | 0 | (session->subscriptionsSize >= server->config.maxSubscriptionsPerSession)) { |
600 | 0 | result->statusCode = UA_STATUSCODE_BADTOOMANYSUBSCRIPTIONS; |
601 | 0 | return; |
602 | 0 | } |
603 | | |
604 | | /* Allocate memory for the new subscription */ |
605 | 0 | UA_Subscription *newSub = (UA_Subscription*)UA_malloc(sizeof(UA_Subscription)); |
606 | 0 | if(!newSub) { |
607 | 0 | result->statusCode = UA_STATUSCODE_BADOUTOFMEMORY; |
608 | 0 | return; |
609 | 0 | } |
610 | | |
611 | | /* Set the available sequence numbers */ |
612 | 0 | result->statusCode = setTransferredSequenceNumbers(sub, result); |
613 | 0 | if(result->statusCode != UA_STATUSCODE_GOOD) { |
614 | 0 | UA_free(newSub); |
615 | 0 | return; |
616 | 0 | } |
617 | | |
618 | | /* Create an identical copy of the Subscription struct. The original |
619 | | * subscription remains in place until a StatusChange notification has been |
620 | | * sent. The elements for lists and queues are moved over manually to ensure |
621 | | * that all backpointers are set correctly. */ |
622 | 0 | memcpy(newSub, sub, sizeof(UA_Subscription)); |
623 | | |
624 | | /* Set to the same state as the original subscription */ |
625 | 0 | newSub->publishCallbackId = 0; |
626 | 0 | result->statusCode = Subscription_setState(server, newSub, sub->state); |
627 | 0 | if(result->statusCode != UA_STATUSCODE_GOOD) { |
628 | 0 | UA_Array_delete(result->availableSequenceNumbers, |
629 | 0 | sub->retransmissionQueueSize, &UA_TYPES[UA_TYPES_UINT32]); |
630 | 0 | result->availableSequenceNumbers = NULL; |
631 | 0 | result->availableSequenceNumbersSize = 0; |
632 | 0 | UA_free(newSub); |
633 | 0 | return; |
634 | 0 | } |
635 | | |
636 | | /* <-- The point of no return --> */ |
637 | | |
638 | | /* Mark the old subscription as transferred to prevent incorrect |
639 | | * diagnostic counter updates when it is deleted */ |
640 | 0 | sub->wasTransferred = true; |
641 | | |
642 | | /* Move over the MonitoredItems and adjust the backpointers */ |
643 | 0 | newSub->monitoredItemsById = sub->monitoredItemsById; |
644 | 0 | ZIP_INIT(&sub->monitoredItemsById); |
645 | 0 | sub->monitoredItemsSize = 0; |
646 | 0 | ZIP_ITER(UA_MonitoredItemIdTree, &newSub->monitoredItemsById, |
647 | 0 | setMonitoredItemSubscriptionVisitor, newSub); |
648 | | |
649 | | /* Move over the samplingMonitoredItems and adjust the backpointers */ |
650 | 0 | LIST_INIT(&newSub->samplingMonitoredItems); |
651 | 0 | UA_MonitoredItem *smon, *smon_tmp; |
652 | 0 | LIST_FOREACH_SAFE(smon, &sub->samplingMonitoredItems, sampling.subscriptionSampling, smon_tmp) { |
653 | 0 | LIST_REMOVE(smon, sampling.subscriptionSampling); |
654 | 0 | LIST_INSERT_HEAD(&newSub->samplingMonitoredItems, smon, |
655 | 0 | sampling.subscriptionSampling); |
656 | 0 | } |
657 | | |
658 | | /* Move over the notification queue */ |
659 | 0 | TAILQ_INIT(&newSub->notificationQueue); |
660 | 0 | UA_Notification *nn, *nn_tmp; |
661 | 0 | TAILQ_FOREACH_SAFE(nn, &sub->notificationQueue, subEntry, nn_tmp) { |
662 | 0 | TAILQ_REMOVE(&sub->notificationQueue, nn, subEntry); |
663 | 0 | TAILQ_INSERT_TAIL(&newSub->notificationQueue, nn, subEntry); |
664 | 0 | } |
665 | 0 | sub->notificationQueueSize = 0; |
666 | 0 | sub->dataChangeNotifications = 0; |
667 | 0 | sub->eventNotifications = 0; |
668 | |
|
669 | 0 | TAILQ_INIT(&newSub->retransmissionQueue); |
670 | 0 | UA_NotificationMessageEntry *nme, *nme_tmp; |
671 | 0 | TAILQ_FOREACH_SAFE(nme, &sub->retransmissionQueue, listEntry, nme_tmp) { |
672 | 0 | TAILQ_REMOVE(&sub->retransmissionQueue, nme, listEntry); |
673 | 0 | TAILQ_INSERT_TAIL(&newSub->retransmissionQueue, nme, listEntry); |
674 | 0 | if(oldSession) |
675 | 0 | oldSession->totalRetransmissionQueueSize -= 1; |
676 | 0 | sub->retransmissionQueueSize -= 1; |
677 | 0 | } |
678 | 0 | UA_assert(sub->retransmissionQueueSize == 0); |
679 | 0 | sub->retransmissionQueueSize = 0; |
680 | | |
681 | | /* Add to the server */ |
682 | 0 | UA_assert(newSub->subscriptionId == sub->subscriptionId); |
683 | 0 | LIST_INSERT_HEAD(&server->subscriptions, newSub, serverListEntry); |
684 | 0 | server->subscriptionsSize++; |
685 | | |
686 | | /* Attach to the session */ |
687 | 0 | UA_Session_attachSubscription(session, newSub); |
688 | | |
689 | | /* Notify the application */ |
690 | 0 | notifySubscription(server, newSub, |
691 | 0 | UA_APPLICATIONNOTIFICATIONTYPE_SUBSCRIPTION_TRANSFERRED); |
692 | |
|
693 | 0 | UA_LOG_INFO_SUBSCRIPTION(server->config.logging, newSub, |
694 | 0 | "Transferred to this Session"); |
695 | | |
696 | | /* Set StatusChange in the original subscription and force publish. This |
697 | | * also removes the Subscription, even if there was no PublishResponse |
698 | | * queued to send a StatusChangeNotification. */ |
699 | 0 | sub->statusChange = UA_STATUSCODE_GOODSUBSCRIPTIONTRANSFERRED; |
700 | 0 | UA_Subscription_publish(server, sub); |
701 | | |
702 | | /* Re-create notifications with the current values for the new subscription */ |
703 | 0 | if(*sendInitialValues) |
704 | 0 | UA_Subscription_resendData(server, newSub); |
705 | | |
706 | | /* Do not update the statistics for the number of Subscriptions here. The |
707 | | * fact that we duplicate the subscription and move over the content is just |
708 | | * an implementtion detail. |
709 | | * server->serverDiagnosticsSummary.currentSubscriptionCount++; |
710 | | * server->serverDiagnosticsSummary.cumulatedSubscriptionCount++; |
711 | | * |
712 | | * Update the diagnostics statistics: */ |
713 | 0 | #ifdef UA_ENABLE_DIAGNOSTICS |
714 | 0 | if(oldSession && |
715 | 0 | UA_equal(&oldSession->clientDescription, &session->clientDescription, |
716 | 0 | &UA_TYPES[UA_TYPES_APPLICATIONDESCRIPTION])) |
717 | 0 | sub->transferredToSameClientCount++; |
718 | 0 | else |
719 | 0 | sub->transferredToAltClientCount++; |
720 | 0 | #endif |
721 | 0 | } |
722 | | |
723 | | UA_Boolean |
724 | | Service_TransferSubscriptions(UA_Server *server, UA_Session *session, |
725 | | const UA_TransferSubscriptionsRequest *request, |
726 | 126 | UA_TransferSubscriptionsResponse *response) { |
727 | 126 | UA_LOG_DEBUG_SESSION(server->config.logging, session, |
728 | 126 | "Processing TransferSubscriptionsRequest"); |
729 | 126 | UA_LOCK_ASSERT(&server->serviceMutex); |
730 | | |
731 | 126 | response->responseHeader.serviceResult = |
732 | 126 | allocProcessServiceOperations(server, session, |
733 | 126 | (UA_ServiceOperation)Operation_TransferSubscription, |
734 | 126 | &request->sendInitialValues, |
735 | 126 | &request->subscriptionIdsSize, |
736 | 126 | &UA_TYPES[UA_TYPES_UINT32], |
737 | 126 | &response->resultsSize, |
738 | 126 | &UA_TYPES[UA_TYPES_TRANSFERRESULT]); |
739 | | return true; |
740 | 126 | } |
741 | | |
742 | | #endif /* UA_ENABLE_SUBSCRIPTIONS */ |