Coverage Report

Created: 2026-09-28 07:11

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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 */