Coverage Report

Created: 2026-09-27 07:10

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/open62541_15/src/server/ua_subscription.h
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, 2021-2022 (c) Fraunhofer IOSB (Author: Julius Pfrommer)
6
 *    Copyright 2015 (c) Chris Iatrou
7
 *    Copyright 2015-2016 (c) Sten Grüner
8
 *    Copyright 2015 (c) Oleksiy Vasylyev
9
 *    Copyright 2017 (c) Florian Palm
10
 *    Copyright 2017 (c) Stefan Profanter, fortiss GmbH
11
 *    Copyright 2017 (c) Mattias Bornhager
12
 *    Copyright 2019 (c) HMS Industrial Networks AB (Author: Jonas Green)
13
 *    Copyright 2020 (c) Christian von Arnim, ISW University of Stuttgart (for VDW and umati)
14
 *    Copyright 2021 (c) Fraunhofer IOSB (Author: Andreas Ebner)
15
 *    Copyright 2021 (c) Fraunhofer IOSB (Author: Jan Hermes)
16
 *    Copyright 2026 (c) o6 Automation GmbH (Author: Andreas Ebner)
17
 */
18
19
#ifndef UA_SUBSCRIPTION_H_
20
#define UA_SUBSCRIPTION_H_
21
22
#include <open62541/types.h>
23
#include <open62541/types_generated.h>
24
#include <open62541/plugin/nodestore.h>
25
26
#include "ua_session.h"
27
#include "../util/ua_util_internal.h"
28
#include "ziptree.h"
29
30
_UA_BEGIN_DECLS
31
32
#ifdef UA_ENABLE_SUBSCRIPTIONS
33
34
/* MonitoredItems create Notifications. Subscriptions collect Notifications from
35
 * (several) MonitoredItems and publish them to the client.
36
 *
37
 * Notifications are put into two queues at the same time. One for the
38
 * MonitoredItem that generated the notification. Here we can remove it if the
39
 * space reserved for the MonitoredItem runs full. The second queue is the
40
 * "global" queue for all Notifications generated in a Subscription. For
41
 * publication, the notifications are taken out of the Subscription's queue in
42
 * the order of their creation. */
43
44
/*****************/
45
/* Notifications */
46
/*****************/
47
48
/* A notification was not (yet) added to the queue of a Subscription */
49
0
#define UA_SUBSCRIPTION_QUEUE_SENTINEL ((UA_Notification*)0x01)
50
51
typedef struct UA_Notification {
52
    /* The subEntry can be a sentinel value to indicate that the Notification is
53
     * not enqueue in the Subscription. This is the case when the Subscription
54
     * is non-reporting.
55
     *
56
     * The monEntry is always set - the Notification must always be enqueued to
57
     * the MonitoredItem right after its creation with UA_Notification_new. */
58
    TAILQ_ENTRY(UA_Notification) subEntry; /* Notification list of the Subscription */
59
    TAILQ_ENTRY(UA_Notification) monEntry; /* Notification list of the MonitoredItem */
60
    UA_MonitoredItem *mon; /* Always set */
61
62
    /* The event field is used if mon->attributeId is the EventNotifier */
63
    union {
64
        UA_MonitoredItemNotification dataChange;
65
#ifdef UA_ENABLE_SUBSCRIPTIONS_EVENTS
66
        UA_EventFieldList event;
67
#endif
68
    } data;
69
70
#ifdef UA_ENABLE_SUBSCRIPTIONS_EVENTS
71
    UA_Boolean isOverflowEvent; /* Counted manually */
72
#endif
73
} UA_Notification;
74
75
/* Initializes and sets the sentinel pointers. Only create a notification if it
76
 * is also going to be immediately enqueued to a MonitoredItem (see below). */
77
UA_Notification * UA_Notification_new(void);
78
79
/* Notifications are always added to the queue of a MonitoredItem. That queue
80
 * can overflow. If Notifications are reported, they are also added to the queue
81
 * of the Subscription. There they are picked up by the publishing callback.
82
 *
83
 * There are two ways Notifications can be put into the global queue of the
84
 * Subscription: They are added because the MonitoringMode of the MonitoredItem
85
 * is "reporting". Or the MonitoringMode is "sampling" and a link is trigered
86
 * that puts the last Notification into the global queue. */
87
void UA_Notification_enqueueAndTrigger(UA_Server *server,
88
                                       UA_Notification *n);
89
90
/* A NotificationMessage contains an array of notifications.
91
 * Sent NotificationMessages are stored for the republish service. */
92
typedef struct UA_NotificationMessageEntry {
93
    TAILQ_ENTRY(UA_NotificationMessageEntry) listEntry;
94
    UA_NotificationMessage message;
95
} UA_NotificationMessageEntry;
96
97
/* Queue Definitions */
98
typedef TAILQ_HEAD(NotificationQueue, UA_Notification) NotificationQueue;
99
typedef TAILQ_HEAD(NotificationMessageQueue, UA_NotificationMessageEntry)
100
    NotificationMessageQueue;
101
102
/*****************/
103
/* MonitoredItem */
104
/*****************/
105
106
/* Maximum number of outstanding async reads per MonitoredItem.
107
 * This protects against unbounded memory usage. */
108
#define UA_MONITOREDITEM_ASYNC_MAX 8
109
110
/* The type of sampling for MonitoredItems depends on the sampling interval.
111
 *
112
 * >0: Cyclic callback
113
 * =0: Attached to the node. Sampling is triggered after every "write".
114
 * <0: Attached to the subscription. Triggered just before every "publish". */
115
typedef enum {
116
    UA_MONITOREDITEMSAMPLINGTYPE_NONE = 0,
117
    UA_MONITOREDITEMSAMPLINGTYPE_CYCLIC, /* Cyclic callback */
118
    UA_MONITOREDITEMSAMPLINGTYPE_EVENT,  /* Attached to the node. Can be a "write
119
                                          * event" for DataChange MonitoredItems
120
                                          * with a zero sampling interval .*/
121
    UA_MONITOREDITEMSAMPLINGTYPE_PUBLISH /* Attached to the subscription */
122
} UA_MonitoredItemSamplingType;
123
124
struct UA_MonitoredItem {
125
    UA_DelayedCallback delayedFreePointers;
126
    ZIP_ENTRY(UA_MonitoredItem) idTreeEntry; /* Index by Id */
127
    UA_Subscription *subscription;           /* Always non-NULL */
128
    UA_UInt32 monitoredItemId;
129
130
    /* Status and Settings */
131
    UA_ReadValueId itemToMonitor;
132
    UA_MonitoringMode monitoringMode;
133
    UA_TimestampsToReturn timestampsToReturn;
134
    UA_Boolean registered;       /* Registered in the server / Subscription */
135
    UA_DateTime triggeredUntil;  /* If the MonitoringMode is SAMPLING,
136
                                  * triggering the MonitoredItem puts the latest
137
                                  * Notification into the publishing queue (of
138
                                  * the Subscription). In addition, the first
139
                                  * new sample is also published (and not just
140
                                  * sampled) if it occurs within the duration of
141
                                  * one publishing cycle after the triggering. */
142
143
    /* If the filter is a UA_DataChangeFilter: The DataChangeFilter always
144
     * contains an absolute deadband definition. Part 8, §6.2 gives the
145
     * following formula to test for percentage deadbands:
146
     *
147
     * DataChange if (absolute value of (last cached value - current value)
148
     *                > (deadbandValue/100.0) * ((high–low) of EURange)))
149
     *
150
     * So we can convert from a percentage to an absolute deadband and keep
151
     * the hot code path simple.
152
     *
153
     * TODO: Store the percentage deadband to recompute when the UARange is
154
     * changed at runtime of the MonitoredItem */
155
    UA_MonitoringParameters parameters;
156
157
    /* Sampling */
158
    UA_MonitoredItemSamplingType samplingType;
159
    union {
160
        UA_UInt64 callbackId;
161
        LIST_ENTRY(UA_MonitoredItem) nodeListEntry; /* Event-Based: linked into
162
                                                     * the Node's MonitoredItem
163
                                                     * list */
164
        LIST_ENTRY(UA_MonitoredItem) subscriptionSampling; /* Linked to publish
165
                                                            * interval */
166
    } sampling;
167
    UA_DataValue lastValue;
168
    UA_UInt32 outstandingAsyncReads; /* at most UA_MONITOREDITEM_ASYNC_MAX */
169
170
    /* Triggering Links */
171
    size_t triggeringLinksSize;
172
    UA_UInt32 *triggeringLinks;
173
174
    /* Notification Queue */
175
    NotificationQueue queue;
176
    size_t queueSize; /* This is the current size. See also the configured
177
                       * (maximum) queueSize in the parameters. */
178
    size_t eventOverflows; /* Separate counter for the queue. Can at most double
179
                            * the queue size */
180
};
181
182
static UA_INLINE UA_Boolean
183
0
UA_MonitoredItem_isDeleting(const UA_MonitoredItem *mon) {
184
0
    return mon->delayedFreePointers.callback != NULL;
185
0
}
Unexecuted instantiation: ua_session.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_nodes.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_server.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_server_ns0.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_server_ns0_diagnostics.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_server_ns0_gds.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_server_config.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_server_binary.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_server_utils.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_server_async.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_subscription.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_subscription_datachange.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_subscription_event.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_subscription_alarms_conditions.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_services.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_services_view.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_services_method.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_services_session.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_services_attribute.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_services_discovery.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_services_subscription.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_services_monitoreditem.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_services_securechannel.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_services_nodemanagement.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_pubsub_connection.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_pubsub_dataset.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_pubsub_writer.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_pubsub_writergroup.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_pubsub_reader.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_pubsub_readergroup.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_pubsub_manager.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_pubsub_ns0.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_pubsub_ns0_sks.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_discovery_mdns.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: ua_discovery.c:UA_MonitoredItem_isDeleting
Unexecuted instantiation: fuzz_server_services.cc:UA_MonitoredItem_isDeleting(UA_MonitoredItem const*)
Unexecuted instantiation: fuzz_tcp_message.cc:UA_MonitoredItem_isDeleting(UA_MonitoredItem const*)
Unexecuted instantiation: fuzz_msg_message.cc:UA_MonitoredItem_isDeleting(UA_MonitoredItem const*)
Unexecuted instantiation: fuzz_opn_message.cc:UA_MonitoredItem_isDeleting(UA_MonitoredItem const*)
Unexecuted instantiation: fuzz_process_request.cc:UA_MonitoredItem_isDeleting(UA_MonitoredItem const*)
Unexecuted instantiation: fuzz_binary_message.cc:UA_MonitoredItem_isDeleting(UA_MonitoredItem const*)
186
187
void UA_MonitoredItem_init(UA_MonitoredItem *mon);
188
void UA_MonitoredItem_delete(UA_Server *server, UA_MonitoredItem *mon,
189
                             UA_Boolean notify);
190
void UA_MonitoredItem_removeOverflowInfoBits(UA_MonitoredItem *mon);
191
void UA_MonitoredItem_register(UA_Server *server, UA_MonitoredItem *mon);
192
193
void
194
notifyMonitoredItem(UA_Server *server, UA_MonitoredItem *mon,
195
                    UA_ApplicationNotificationType type);
196
197
typedef ZIP_HEAD(UA_MonitoredItemIdTree, UA_MonitoredItem) UA_MonitoredItemIdTree;
198
199
static enum ZIP_CMP
200
UA_MonitoredItemIdTree_cmp(const UA_UInt32 *a,
201
0
                           const UA_UInt32 *b) {
202
0
    if(*a < *b)
203
0
        return ZIP_CMP_LESS;
204
0
    if(*a > *b)
205
0
        return ZIP_CMP_MORE;
206
0
    return ZIP_CMP_EQ;
207
0
}
Unexecuted instantiation: ua_session.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_nodes.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_server.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_server_ns0.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_server_ns0_diagnostics.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_server_ns0_gds.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_server_config.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_server_binary.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_server_utils.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_server_async.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_subscription.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_subscription_datachange.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_subscription_event.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_subscription_alarms_conditions.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_services.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_services_view.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_services_method.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_services_session.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_services_attribute.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_services_discovery.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_services_subscription.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_services_monitoreditem.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_services_securechannel.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_services_nodemanagement.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_pubsub_connection.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_pubsub_dataset.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_pubsub_writer.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_pubsub_writergroup.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_pubsub_reader.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_pubsub_readergroup.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_pubsub_manager.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_pubsub_ns0.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_pubsub_ns0_sks.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_discovery_mdns.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: ua_discovery.c:UA_MonitoredItemIdTree_cmp
Unexecuted instantiation: fuzz_server_services.cc:UA_MonitoredItemIdTree_cmp(unsigned int const*, unsigned int const*)
Unexecuted instantiation: fuzz_tcp_message.cc:UA_MonitoredItemIdTree_cmp(unsigned int const*, unsigned int const*)
Unexecuted instantiation: fuzz_msg_message.cc:UA_MonitoredItemIdTree_cmp(unsigned int const*, unsigned int const*)
Unexecuted instantiation: fuzz_opn_message.cc:UA_MonitoredItemIdTree_cmp(unsigned int const*, unsigned int const*)
Unexecuted instantiation: fuzz_process_request.cc:UA_MonitoredItemIdTree_cmp(unsigned int const*, unsigned int const*)
Unexecuted instantiation: fuzz_binary_message.cc:UA_MonitoredItemIdTree_cmp(unsigned int const*, unsigned int const*)
208
209
19.6k
ZIP_FUNCTIONS(UA_MonitoredItemIdTree, UA_MonitoredItem, idTreeEntry,
Unexecuted instantiation: ua_server_ns0.c:UA_MonitoredItemIdTree_ZIP_ITER
Unexecuted instantiation: ua_server_ns0_diagnostics.c:UA_MonitoredItemIdTree_ZIP_ITER
ua_subscription.c:UA_MonitoredItemIdTree_ZIP_ITER
Line
Count
Source
209
ZIP_FUNCTIONS(UA_MonitoredItemIdTree, UA_MonitoredItem, idTreeEntry,
Unexecuted instantiation: ua_subscription.c:UA_MonitoredItemIdTree_ZIP_INSERT
Unexecuted instantiation: ua_subscription.c:UA_MonitoredItemIdTree_ZIP_REMOVE
Unexecuted instantiation: ua_services_subscription.c:UA_MonitoredItemIdTree_ZIP_ITER
210
              UA_UInt32, monitoredItemId, UA_MonitoredItemIdTree_cmp)
211
212
/* UA_NodeHead.monitoredItems is a bare UA_MonitoredItem* in the public
213
 * nodestore.h, which is layout-compatible with this LIST_HEAD. The server runs
214
 * the LIST_* macros over the per-node MonitoredItem list by casting the field's
215
 * address to UA_MonitoredItemList*. */
216
typedef LIST_HEAD(UA_MonitoredItemList, UA_MonitoredItem) UA_MonitoredItemList;
217
218
/* Register sampling. Either by adding a repeated callback or by adding the
219
 * MonitoredItem to a linked list in the node. */
220
UA_StatusCode
221
UA_MonitoredItem_registerSampling(UA_Server *server, UA_MonitoredItem *mon);
222
223
void
224
UA_MonitoredItem_unregisterSampling(UA_Server *server, UA_MonitoredItem *mon);
225
226
UA_StatusCode
227
UA_MonitoredItem_setMonitoringMode(UA_Server *server, UA_MonitoredItem *mon,
228
                                   UA_MonitoringMode monitoringMode);
229
230
void
231
UA_MonitoredItem_sample(UA_Server *server, UA_MonitoredItem *mon);
232
233
/* Do not use the value after calling this. It will be moved to mon or freed. */
234
void
235
UA_MonitoredItem_processSampledValue(UA_Server *server, UA_MonitoredItem *mon,
236
                                     UA_DataValue *value);
237
238
UA_StatusCode
239
UA_MonitoredItem_removeLink(UA_Subscription *sub, UA_MonitoredItem *mon,
240
                            UA_UInt32 linkId);
241
242
UA_StatusCode
243
UA_MonitoredItem_addLink(UA_Subscription *sub, UA_MonitoredItem *mon,
244
                         UA_UInt32 linkId);
245
246
UA_StatusCode
247
UA_MonitoredItem_createDataChangeNotification(UA_Server *server, UA_MonitoredItem *mon,
248
                                              const UA_DataValue *value);
249
250
/* Remove entries until mon->maxQueueSize is reached. Sets infobits for lost
251
 * data if required. */
252
void UA_MonitoredItem_ensureQueueSpace(UA_Server *server, UA_MonitoredItem *mon);
253
254
/****************/
255
/* Subscription */
256
/****************/
257
258
/* We use only a subset of the states defined in the standard */
259
typedef enum {
260
    UA_SUBSCRIPTIONSTATE_STOPPED = 0,
261
    UA_SUBSCRIPTIONSTATE_REMOVING,
262
    UA_SUBSCRIPTIONSTATE_ENABLED_NOPUBLISH, /* only keepalive */
263
    UA_SUBSCRIPTIONSTATE_ENABLED
264
} UA_SubscriptionState;
265
266
/* Subscriptions are managed in a server-wide linked list. If they are attached
267
 * to a Session, then they are additionaly in the per-Session linked-list. A
268
 * subscription is always generated for a Session. But the CloseSession Service
269
 * may keep Subscriptions intact beyond the Session lifetime. They can then be
270
 * re-bound to a new Session with the TransferSubscription Service. */
271
struct UA_Subscription {
272
    UA_DelayedCallback delayedFreePointers;
273
    LIST_ENTRY(UA_Subscription) serverListEntry;
274
    /* Ordered according to the priority byte and round-robin scheduling for
275
     * late subscriptions. See ua_session.h. Only set if session != NULL. */
276
    TAILQ_ENTRY(UA_Subscription) sessionListEntry;
277
    UA_Session *session; /* May be NULL if no session is attached. */
278
    UA_UInt32 subscriptionId;
279
280
    /* Settings */
281
    UA_UInt32 lifeTimeCount;
282
    UA_UInt32 maxKeepAliveCount;
283
    UA_Double publishingInterval; /* in ms */
284
    UA_UInt32 notificationsPerPublish;
285
    UA_Byte priority;
286
287
    /* Runtime information */
288
    UA_SubscriptionState state;
289
    UA_Boolean late;
290
    UA_Boolean wasTransferred; /* Set to true when this subscription was
291
                                * transferred to another session. Used to
292
                                * prevent incorrect diagnostic counter updates. */
293
    UA_StatusCode statusChange; /* If set, a notification is generated and the
294
                                 * Subscription is deleted within
295
                                 * UA_Subscription_publish. */
296
    UA_UInt32 nextSequenceNumber;
297
    UA_UInt32 currentKeepAliveCount;
298
    UA_UInt32 currentLifetimeCount;
299
300
    /* Publish Callback. Registered if id > 0. */
301
    UA_UInt64 publishCallbackId;
302
303
    /* Delayed callback to schedule publication of more notifications */
304
    UA_Boolean delayedCallbackRegistered;
305
    UA_DelayedCallback delayedMoreNotifications;
306
307
    /* MonitoredItems */
308
    UA_UInt32 lastMonitoredItemId; /* increase the identifiers */
309
    UA_MonitoredItemIdTree monitoredItemsById;
310
    UA_UInt32 monitoredItemsSize;
311
312
    /* MonitoredItems that are sampled in every publish callback (with the
313
     * publish interval of the subscription) */
314
    LIST_HEAD(, UA_MonitoredItem) samplingMonitoredItems;
315
316
    /* Global list of notifications from the MonitoredItems */
317
    TAILQ_HEAD(, UA_Notification) notificationQueue;
318
    UA_UInt32 notificationQueueSize; /* Total queue size */
319
    UA_UInt32 dataChangeNotifications;
320
    UA_UInt32 eventNotifications;
321
322
    /* Retransmission Queue */
323
    NotificationMessageQueue retransmissionQueue;
324
    size_t retransmissionQueueSize;
325
326
    /* Statistics for the server diagnostics. The fields are defined according
327
     * to the SubscriptionDiagnosticsDataType (Part 5, §12.15). */
328
#ifdef UA_ENABLE_DIAGNOSTICS
329
    UA_NodeId ns0Id; /* Representation in the Session object */
330
331
    UA_UInt32 modifyCount;
332
    UA_UInt32 enableCount;
333
    UA_UInt32 disableCount;
334
    UA_UInt32 republishRequestCount;
335
    UA_UInt32 republishMessageCount;
336
    UA_UInt32 transferRequestCount;
337
    UA_UInt32 transferredToAltClientCount;
338
    UA_UInt32 transferredToSameClientCount;
339
    UA_UInt32 publishRequestCount;
340
    UA_UInt32 dataChangeNotificationsCount;
341
    UA_UInt32 eventNotificationsCount;
342
    UA_UInt32 notificationsCount;
343
    UA_UInt32 latePublishRequestCount;
344
    UA_UInt32 discardedMessageCount;
345
    UA_UInt32 monitoringQueueOverflowCount;
346
    UA_UInt32 eventQueueOverflowCount;
347
#endif
348
};
349
350
UA_Subscription * UA_Subscription_new(void);
351
352
void
353
UA_Subscription_delete(UA_Server *server, UA_Subscription *sub);
354
355
UA_StatusCode
356
Subscription_setState(UA_Server *server, UA_Subscription *sub,
357
                      UA_SubscriptionState state);
358
359
void
360
Subscription_resetLifetime(UA_Subscription *sub);
361
362
UA_Subscription *
363
getSubscriptionById(UA_Server *server, UA_UInt32 subscriptionId);
364
365
UA_MonitoredItem *
366
UA_Subscription_getMonitoredItem(UA_Subscription *sub,
367
                                 UA_UInt32 monitoredItemId);
368
369
void
370
UA_Subscription_publish(UA_Server *server, UA_Subscription *sub);
371
372
void
373
UA_Subscription_localPublish(UA_Server *server, UA_Subscription *sub);
374
375
void
376
UA_Subscription_resendData(UA_Server *server, UA_Subscription *sub);
377
378
UA_StatusCode
379
UA_Subscription_removeRetransmissionMessage(UA_Subscription *sub,
380
                                            UA_UInt32 sequenceNumber);
381
382
void
383
UA_Session_ensurePublishQueueSpace(UA_Server *server, UA_Session *session);
384
385
/* Forward declaration for A&C used in ua_server_internal.h" */
386
struct UA_ConditionSource;
387
typedef struct UA_ConditionSource UA_ConditionSource;
388
389
/* Event Handling */
390
#ifdef UA_ENABLE_SUBSCRIPTIONS_EVENTS
391
392
0
#define UA_EVENTFILTER_MAXELEMENTS 64 /* Max operator elements */
393
0
#define UA_EVENTFILTER_MAXOPERANDS 64 /* Max operands per operator */
394
0
#define UA_EVENTFILTER_MAXSELECT   64 /* Max select clauses */
395
396
/* Static validation when the filter is registered */
397
UA_StatusCode
398
UA_SimpleAttributeOperandValidation(UA_Server *server,
399
                                    const UA_SimpleAttributeOperand *sao);
400
401
/* Static validation when the filter is registered */
402
UA_ContentFilterElementResult
403
UA_ContentFilterElementValidation(UA_Server *server, size_t operatorIndex,
404
                                  size_t operatorsCount,
405
                                  const UA_ContentFilterElement *ef);
406
407
/* If the sessionId, subscriptionId or monitoredItemId are non-NULL, they are
408
 * used to filter the subscriptions and monitoreditems that emit the event. */
409
UA_StatusCode
410
createEvent(UA_Server *server, const UA_EventDescription *ed,
411
            UA_ByteString *outEventId);
412
413
typedef struct {
414
    UA_Server *server;
415
    UA_Session *session; /* may be NULL if no session is attached. */
416
    UA_EventDescription ed; /* shallow copy */
417
    UA_EventFilter filter;  /* shallow copy */
418
419
    /* Don't reread or regenerate values that have already been gotten during
420
     * the same filter evaluation */
421
    UA_KeyValueMap fieldCache;
422
423
    /* Generated / cached EventId */
424
    UA_ByteString eventId;
425
    UA_Byte eventIdBuf[16];
426
427
    /* <-- Variables used only by evaluateWhereClause --> */
428
429
    UA_Variant operatorResults[UA_EVENTFILTER_MAXELEMENTS];
430
431
    /* The operand stack contains temporary variants. Cleaned up after the
432
     * evaluation of each operator. */
433
    size_t top;
434
    UA_Variant operandStack[UA_EVENTFILTER_MAXOPERANDS];
435
} UA_FilterEvalContext;
436
437
/* The _reset method resets the filter between evaluations for different
438
 * MonitoredItems. The EventId is reset *only* if eventId.data != eventIdBuf.
439
 * This does not lead to memleaks and ensures we do not regenerate random
440
 * EventIds for the same event on different MonitoredItems. */
441
void UA_FilterEvalContext_init(UA_FilterEvalContext *ctx);
442
void UA_FilterEvalContext_reset(UA_FilterEvalContext *ctx);
443
444
UA_StatusCode
445
resolveSAO(UA_FilterEvalContext *ctx, const UA_SimpleAttributeOperand *sao,
446
           UA_Variant *out);
447
448
/* Retrieve or generate the unique EventId and cache it. Takes the EventId from
449
 * the EventDescription if explicitly set by the user. If successful,
450
 * ctx->eventId contains the cached EventId */
451
UA_StatusCode cacheEventId(UA_FilterEvalContext *ctx);
452
453
/* Evaluate content filter, exported only for unit testing */
454
UA_StatusCode
455
evaluateWhereClause(UA_FilterEvalContext *ctx);
456
457
/* Applies the select clause and resolves the result fields */
458
UA_StatusCode
459
evaluateSelectClause(UA_FilterEvalContext *ctx, UA_EventFieldList *efl);
460
461
#endif /* UA_ENABLE_SUBSCRIPTIONS_EVENTS */
462
463
/***********/
464
/* Helpers */
465
/***********/
466
467
/* Setting an integer value within bounds */
468
#define UA_BOUNDEDVALUE_SETWBOUNDS(BOUNDS, SRC, DST)   \
469
720
    do { \
470
720
        if(SRC > BOUNDS.max) DST = BOUNDS.max;         \
471
720
        else if(SRC < BOUNDS.min) DST = BOUNDS.min;    \
472
325
        else DST = SRC;                                \
473
720
    } while (0)
474
475
/* Logging
476
 * See a description of the tricks used in ua_session.h */
477
#define UA_LOG_SUBSCRIPTION_INTERNAL(LOGGER, LEVEL, SUB, MSG, ...)      \
478
    do {                                                                \
479
        if((SUB) && (SUB)->session) {                                   \
480
            UA_LOG_##LEVEL##_SESSION(LOGGER, (SUB)->session,            \
481
                                     "Subscription %" PRIu32 " | " MSG "%.0s", \
482
                                     (SUB)->subscriptionId, __VA_ARGS__); \
483
        } else {                                                        \
484
            UA_LOG_##LEVEL(LOGGER, UA_LOGCATEGORY_SERVER,               \
485
                           "Subscription %" PRIu32 " | " MSG "%.0s",    \
486
                           (SUB) ? (SUB)->subscriptionId : 0, __VA_ARGS__); \
487
        }                                                               \
488
    } while(0)
489
490
#if UA_LOGLEVEL <= 100
491
# define UA_LOG_TRACE_SUBSCRIPTION(LOGGER, SUB, ...)                     \
492
    UA_MACRO_EXPAND(UA_LOG_SUBSCRIPTION_INTERNAL(LOGGER, TRACE, SUB, __VA_ARGS__, ""))
493
#else
494
# define UA_LOG_TRACE_SUBSCRIPTION(LOGGER, SUB, ...) do {} while(0)
495
#endif
496
497
#if UA_LOGLEVEL <= 200
498
# define UA_LOG_DEBUG_SUBSCRIPTION(LOGGER, SUB, ...)                     \
499
    UA_MACRO_EXPAND(UA_LOG_SUBSCRIPTION_INTERNAL(LOGGER, DEBUG, SUB, __VA_ARGS__, ""))
500
#else
501
0
# define UA_LOG_DEBUG_SUBSCRIPTION(LOGGER, SUB, ...) do {} while(0)
502
#endif
503
504
#if UA_LOGLEVEL <= 300
505
# define UA_LOG_INFO_SUBSCRIPTION(LOGGER, SUB, ...)                     \
506
    UA_MACRO_EXPAND(UA_LOG_SUBSCRIPTION_INTERNAL(LOGGER, INFO, SUB, __VA_ARGS__, ""))
507
#else
508
6.69k
# define UA_LOG_INFO_SUBSCRIPTION(LOGGER, SUB, ...) do {} while(0)
509
#endif
510
511
#if UA_LOGLEVEL <= 400
512
# define UA_LOG_WARNING_SUBSCRIPTION(LOGGER, SUB, ...)                     \
513
    UA_MACRO_EXPAND(UA_LOG_SUBSCRIPTION_INTERNAL(LOGGER, WARNING, SUB, __VA_ARGS__, ""))
514
#else
515
0
# define UA_LOG_WARNING_SUBSCRIPTION(LOGGER, SUB, ...) do {} while(0)
516
#endif
517
518
#if UA_LOGLEVEL <= 500
519
# define UA_LOG_ERROR_SUBSCRIPTION(LOGGER, SUB, ...)                     \
520
    UA_MACRO_EXPAND(UA_LOG_SUBSCRIPTION_INTERNAL(LOGGER, ERROR, SUB, __VA_ARGS__, ""))
521
#else
522
0
# define UA_LOG_ERROR_SUBSCRIPTION(LOGGER, SUB, ...) do {} while(0)
523
#endif
524
525
#if UA_LOGLEVEL <= 600
526
# define UA_LOG_FATAL_SUBSCRIPTION(LOGGER, SUB, ...)                     \
527
    UA_MACRO_EXPAND(UA_LOG_SUBSCRIPTION_INTERNAL(LOGGER, FATAL, SUB, __VA_ARGS__, ""))
528
#else
529
# define UA_LOG_FATAL_SUBSCRIPTION(LOGGER, SUB, ...) do {} while(0)
530
#endif
531
532
#endif /* UA_ENABLE_SUBSCRIPTIONS */
533
534
_UA_END_DECLS
535
536
#endif /* UA_SUBSCRIPTION_H_ */