/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_ */ |