/src/open62541/arch/posix/eventloop_posix.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 2021 (c) Fraunhofer IOSB (Author: Julius Pfrommer) |
6 | | * Copyright 2021 (c) Fraunhofer IOSB (Author: Jan Hermes) |
7 | | * Copyright 2026 (c) o6 Automation GmbH (Author: Julius Pfrommer) |
8 | | */ |
9 | | |
10 | | #include "eventloop_posix.h" |
11 | | #include "open62541/plugin/eventloop.h" |
12 | | |
13 | | #if defined(UA_ARCHITECTURE_POSIX) && !defined(UA_ARCHITECTURE_LWIP) |
14 | | |
15 | | /*********/ |
16 | | /* Timer */ |
17 | | /*********/ |
18 | | |
19 | | UA_DateTime |
20 | 0 | UA_EventLoopPOSIX_nextTimer(UA_EventLoop *public_el) { |
21 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
22 | 0 | if(UA_atomic_load(&el->delayedHead1) > (UA_DelayedCallback *)0x01 || |
23 | 0 | UA_atomic_load(&el->delayedHead2) > (UA_DelayedCallback *)0x01) |
24 | 0 | return el->eventLoop.dateTime_nowMonotonic(&el->eventLoop); |
25 | 0 | return UA_Timer_next(&el->timer); |
26 | 0 | } |
27 | | |
28 | | UA_StatusCode |
29 | | UA_EventLoopPOSIX_addTimer(UA_EventLoop *public_el, UA_Callback cb, |
30 | | void *application, void *data, UA_Double interval_ms, |
31 | | UA_DateTime *baseTime, UA_TimerPolicy timerPolicy, |
32 | 0 | UA_UInt64 *callbackId) { |
33 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
34 | 0 | return UA_Timer_add(&el->timer, cb, application, data, interval_ms, |
35 | 0 | public_el->dateTime_nowMonotonic(public_el), |
36 | 0 | baseTime, timerPolicy, callbackId); |
37 | 0 | } |
38 | | |
39 | | UA_StatusCode |
40 | | UA_EventLoopPOSIX_modifyTimer(UA_EventLoop *public_el, |
41 | | UA_UInt64 callbackId, |
42 | | UA_Double interval_ms, |
43 | | UA_DateTime *baseTime, |
44 | 0 | UA_TimerPolicy timerPolicy) { |
45 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
46 | 0 | return UA_Timer_modify(&el->timer, callbackId, interval_ms, |
47 | 0 | public_el->dateTime_nowMonotonic(public_el), |
48 | 0 | baseTime, timerPolicy); |
49 | 0 | } |
50 | | |
51 | | void |
52 | | UA_EventLoopPOSIX_removeTimer(UA_EventLoop *public_el, |
53 | 0 | UA_UInt64 callbackId) { |
54 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
55 | 0 | UA_Timer_remove(&el->timer, callbackId); |
56 | 0 | } |
57 | | |
58 | | void |
59 | | UA_EventLoopPOSIX_addDelayedCallback(UA_EventLoop *public_el, |
60 | 0 | UA_DelayedCallback *dc) { |
61 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
62 | 0 | dc->next = NULL; |
63 | | |
64 | | /* el->delayedTail points either to prev->next or to the head. In an atomic |
65 | | * xchg-operation we make the tail point to dc. This also gives us |
66 | | * prev->next. Then we make prev->next point to dc. |
67 | | * |
68 | | * This is thread-safe. Another thread might retrieve dc from the tail. |
69 | | * Then he can set dc->next while we are still updating prev->next. |
70 | | * It is ensured that only on thread can updated dc->next. */ |
71 | 0 | UA_atomic(UA_atomic(UA_DelayedCallback*)*) prev_next; |
72 | 0 | UA_atomic_xchg(&el->delayedTail, &dc->next, &prev_next); |
73 | 0 | UA_atomic_store(prev_next, dc); |
74 | 0 | } |
75 | | |
76 | | /* Resets the delayed queue and returns the previous head and tail */ |
77 | | static void |
78 | | resetDelayedQueue(UA_EventLoopPOSIX *el, |
79 | | UA_atomic(UA_DelayedCallback*)* oldHead, |
80 | 0 | UA_atomic(UA_atomic(UA_DelayedCallback*)*)* oldTail) { |
81 | 0 | if(UA_atomic_load(&el->delayedHead1) <= (UA_DelayedCallback *)0x01 && |
82 | 0 | UA_atomic_load(&el->delayedHead2) <= (UA_DelayedCallback *)0x01) |
83 | 0 | return; /* The queue is empty */ |
84 | | |
85 | | /* Get the location of the active and the inactive head */ |
86 | 0 | UA_Boolean active1 = (UA_atomic_load(&el->delayedHead1) != (UA_DelayedCallback*)0x01); |
87 | 0 | UA_atomic(UA_DelayedCallback*)* activeHead = (active1) ? &el->delayedHead1 : &el->delayedHead2; |
88 | 0 | UA_atomic(UA_DelayedCallback*)* inactiveHead = (active1) ? &el->delayedHead2 : &el->delayedHead1; |
89 | | |
90 | | /* Set NULL to the inactive head. This indicates it is now active. */ |
91 | 0 | UA_atomic_store(inactiveHead, NULL); |
92 | | |
93 | | /* Set a sentinel value to "inactivate" the active head. Return the old |
94 | | * active head. Parallel threads may continue to add elements below the old |
95 | | * "activeHead" if they already have a pointer. */ |
96 | 0 | UA_atomic_xchg(activeHead, (UA_DelayedCallback*)0x01, oldHead); |
97 | | |
98 | | /* Make the inactiveHead the new "active" by pointing to it from the tail. |
99 | | * Also return the old tail. From the consumer-thread we can then iterate |
100 | | * the linked-list until we find the old tail as the last element. */ |
101 | 0 | UA_atomic_xchg(&el->delayedTail, inactiveHead, oldTail); |
102 | 0 | } |
103 | | |
104 | | void |
105 | | UA_EventLoopPOSIX_removeDelayedCallback(UA_EventLoop *public_el, |
106 | 0 | UA_DelayedCallback *dc) { |
107 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
108 | 0 | UA_LOCK(&el->elMutex); |
109 | | |
110 | | /* Reset and get the old head and tail */ |
111 | 0 | UA_atomic(UA_DelayedCallback *) cur = NULL; |
112 | 0 | UA_atomic(UA_atomic(UA_DelayedCallback*)*) tail = NULL; |
113 | 0 | resetDelayedQueue(el, &cur, &tail); |
114 | | |
115 | | /* tail points to the location where the next element shall be inserted: The |
116 | | * next-pointer of the last element. Since the next-pointer is the first |
117 | | * struct member, we can directly cast to the last element. */ |
118 | 0 | UA_DelayedCallback *last = (UA_DelayedCallback*)(uintptr_t)tail; |
119 | | |
120 | | /* Loop until we reach the tail (or head and tail are both NULL) */ |
121 | 0 | UA_DelayedCallback *next; |
122 | 0 | for(; cur; cur = next) { |
123 | | /* Spin-loop until the next-pointer of cur is updated. |
124 | | * The element pointed to by tail must appear eventually. */ |
125 | 0 | next = UA_atomic_load(&cur->next); |
126 | 0 | while(!next && cur != last) |
127 | 0 | next = UA_atomic_load(&cur->next); |
128 | 0 | if(cur == dc) |
129 | 0 | continue; |
130 | 0 | UA_EventLoopPOSIX_addDelayedCallback(public_el, cur); |
131 | 0 | } |
132 | |
|
133 | 0 | UA_UNLOCK(&el->elMutex); |
134 | 0 | } |
135 | | |
136 | | void |
137 | 0 | UA_EventLoopPOSIX_processDelayed(UA_EventLoopPOSIX *el) { |
138 | 0 | UA_LOG_TRACE(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
139 | 0 | "Process delayed callbacks"); |
140 | |
|
141 | 0 | UA_LOCK_ASSERT(&el->elMutex); |
142 | | |
143 | | /* Reset and get the old head and tail */ |
144 | 0 | UA_atomic(UA_DelayedCallback *) dc = NULL; |
145 | 0 | UA_atomic(UA_atomic(UA_DelayedCallback*)*) tail = NULL; |
146 | 0 | resetDelayedQueue(el, &dc, &tail); |
147 | | |
148 | | /* tail points to the location where the next element shall be inserted: The |
149 | | * next-pointer of the last element. Since the next-pointer is the first |
150 | | * struct member, we can directly cast to the last element. */ |
151 | 0 | UA_DelayedCallback *last = (UA_DelayedCallback*)(uintptr_t)tail; |
152 | | |
153 | | /* Loop until we reach the tail (or head and tail are both NULL) */ |
154 | 0 | UA_DelayedCallback *next; |
155 | 0 | for(; dc; dc = next) { |
156 | 0 | next = UA_atomic_load(&dc->next); |
157 | 0 | while(!next && dc != last) |
158 | 0 | next = UA_atomic_load(&dc->next); |
159 | 0 | if(!dc->callback) |
160 | 0 | continue; |
161 | 0 | dc->callback(dc->application, dc->context); |
162 | 0 | } |
163 | 0 | } |
164 | | |
165 | | /***********************/ |
166 | | /* EventLoop Lifecycle */ |
167 | | /***********************/ |
168 | | |
169 | | static UA_StatusCode |
170 | 0 | UA_EventLoopPOSIX_start(UA_EventLoop *public_el) { |
171 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
172 | 0 | UA_LOCK(&el->elMutex); |
173 | |
|
174 | 0 | if(el->eventLoop.state != UA_EVENTLOOPSTATE_FRESH && |
175 | 0 | el->eventLoop.state != UA_EVENTLOOPSTATE_STOPPED) { |
176 | 0 | UA_UNLOCK(&el->elMutex); |
177 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
178 | 0 | } |
179 | | |
180 | 0 | UA_LOG_DEBUG(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
181 | 0 | "Starting the EventLoop"); |
182 | | |
183 | | /* Setting custom clock source */ |
184 | 0 | const UA_Int32 *cs = (const UA_Int32*) |
185 | 0 | UA_KeyValueMap_getScalar(&el->eventLoop.params, |
186 | 0 | UA_QUALIFIEDNAME(0, "clock-source"), |
187 | 0 | &UA_TYPES[UA_TYPES_INT32]); |
188 | 0 | if(cs) |
189 | 0 | el->clockSource = *cs; |
190 | |
|
191 | 0 | const UA_Int32 *csm = (const UA_Int32*) |
192 | 0 | UA_KeyValueMap_getScalar(&el->eventLoop.params, |
193 | 0 | UA_QUALIFIEDNAME(0, "clock-source-monotonic"), |
194 | 0 | &UA_TYPES[UA_TYPES_INT32]); |
195 | 0 | if(csm) { |
196 | 0 | if(el->clockSourceMonotonic != *csm && el->timer.idTree.root) { |
197 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
198 | 0 | "Eventloop\t| Setting a different monotonic clock, ", |
199 | 0 | "but existing timers have been registered with a " |
200 | 0 | "different clock source"); |
201 | 0 | } |
202 | 0 | el->clockSourceMonotonic = *csm; |
203 | 0 | } |
204 | | |
205 | | |
206 | | /* Create the self-pipe */ |
207 | 0 | int err = UA_EventLoopPOSIX_pipe(el->selfpipe); |
208 | 0 | if(err != 0) { |
209 | 0 | UA_LOG_SOCKET_ERRNO_WRAP( |
210 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_NETWORK, |
211 | 0 | "Eventloop\t| Could not create the self-pipe (%s)", |
212 | 0 | errno_str)); |
213 | 0 | UA_UNLOCK(&el->elMutex); |
214 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
215 | 0 | } |
216 | | |
217 | | /* Create the epoll socket */ |
218 | 0 | #ifdef UA_HAVE_EPOLL |
219 | 0 | el->epollfd = epoll_create1(0); |
220 | 0 | if(el->epollfd == -1) { |
221 | 0 | UA_LOG_SOCKET_ERRNO_WRAP( |
222 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_NETWORK, |
223 | 0 | "Eventloop\t| Could not create the epoll socket (%s)", |
224 | 0 | errno_str)); |
225 | 0 | UA_close(el->selfpipe[0]); |
226 | 0 | UA_close(el->selfpipe[1]); |
227 | 0 | UA_UNLOCK(&el->elMutex); |
228 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
229 | 0 | } |
230 | | |
231 | | /* epoll always listens on the self-pipe. This is the only epoll_event that |
232 | | * has a NULL data pointer. */ |
233 | 0 | struct epoll_event event; |
234 | 0 | memset(&event, 0, sizeof(struct epoll_event)); |
235 | 0 | event.events = EPOLLIN; |
236 | 0 | err = epoll_ctl(el->epollfd, EPOLL_CTL_ADD, el->selfpipe[0], &event); |
237 | 0 | if(err != 0) { |
238 | 0 | UA_LOG_SOCKET_ERRNO_WRAP( |
239 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_NETWORK, |
240 | 0 | "Eventloop\t| Could not register the self-pipe for epoll (%s)", |
241 | 0 | errno_str)); |
242 | 0 | UA_close(el->selfpipe[0]); |
243 | 0 | UA_close(el->selfpipe[1]); |
244 | 0 | close(el->epollfd); |
245 | 0 | UA_UNLOCK(&el->elMutex); |
246 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
247 | 0 | } |
248 | 0 | #endif |
249 | | |
250 | | /* Start the EventSources */ |
251 | 0 | UA_StatusCode res = UA_STATUSCODE_GOOD; |
252 | 0 | UA_EventSource *es = el->eventLoop.eventSources; |
253 | 0 | while(es) { |
254 | 0 | res |= es->start(es); |
255 | 0 | es = es->next; |
256 | 0 | } |
257 | | |
258 | | /* Dirty-write the state that is const "from the outside" */ |
259 | 0 | *(UA_EventLoopState*)(uintptr_t)&el->eventLoop.state = |
260 | 0 | UA_EVENTLOOPSTATE_STARTED; |
261 | |
|
262 | 0 | UA_UNLOCK(&el->elMutex); |
263 | 0 | return res; |
264 | 0 | } |
265 | | |
266 | | static void |
267 | 0 | checkClosed(UA_EventLoopPOSIX *el) { |
268 | 0 | UA_LOCK_ASSERT(&el->elMutex); |
269 | |
|
270 | 0 | UA_EventSource *es = el->eventLoop.eventSources; |
271 | 0 | while(es) { |
272 | 0 | if(es->state != UA_EVENTSOURCESTATE_STOPPED) |
273 | 0 | return; |
274 | 0 | es = es->next; |
275 | 0 | } |
276 | | |
277 | | /* Not closed until all delayed callbacks are processed */ |
278 | 0 | if(UA_atomic_load(&el->delayedHead1) != NULL && |
279 | 0 | UA_atomic_load(&el->delayedHead2) != NULL) |
280 | 0 | return; |
281 | | |
282 | | /* Close the self-pipe when everything else is done */ |
283 | 0 | UA_close(el->selfpipe[0]); |
284 | 0 | UA_close(el->selfpipe[1]); |
285 | | |
286 | | /* Dirty-write the state that is const "from the outside" */ |
287 | 0 | *(UA_EventLoopState*)(uintptr_t)&el->eventLoop.state = |
288 | 0 | UA_EVENTLOOPSTATE_STOPPED; |
289 | | |
290 | | /* Close the epoll/IOCP socket once all EventSources have shut down */ |
291 | 0 | #ifdef UA_HAVE_EPOLL |
292 | 0 | UA_close(el->epollfd); |
293 | 0 | #endif |
294 | |
|
295 | 0 | UA_LOG_DEBUG(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
296 | 0 | "The EventLoop has stopped"); |
297 | 0 | } |
298 | | |
299 | | static void |
300 | 0 | UA_EventLoopPOSIX_stop(UA_EventLoop *public_el) { |
301 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
302 | 0 | UA_LOCK(&el->elMutex); |
303 | |
|
304 | 0 | if(el->eventLoop.state != UA_EVENTLOOPSTATE_STARTED) { |
305 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
306 | 0 | "The EventLoop is not running, cannot be stopped"); |
307 | 0 | UA_UNLOCK(&el->elMutex); |
308 | 0 | return; |
309 | 0 | } |
310 | | |
311 | 0 | UA_LOG_DEBUG(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
312 | 0 | "Stopping the EventLoop"); |
313 | | |
314 | | /* Set to STOPPING to prevent "normal use" */ |
315 | 0 | *(UA_EventLoopState*)(uintptr_t)&el->eventLoop.state = |
316 | 0 | UA_EVENTLOOPSTATE_STOPPING; |
317 | | |
318 | | /* Stop all event sources (asynchronous) */ |
319 | 0 | UA_EventSource *es = el->eventLoop.eventSources; |
320 | 0 | for(; es; es = es->next) { |
321 | 0 | if(es->state == UA_EVENTSOURCESTATE_STARTING || |
322 | 0 | es->state == UA_EVENTSOURCESTATE_STARTED) { |
323 | 0 | es->stop(es); |
324 | 0 | } |
325 | 0 | } |
326 | | |
327 | | /* Set to STOPPED if all EventSources are STOPPED */ |
328 | 0 | checkClosed(el); |
329 | |
|
330 | 0 | UA_UNLOCK(&el->elMutex); |
331 | 0 | } |
332 | | |
333 | | static UA_StatusCode |
334 | 0 | UA_EventLoopPOSIX_run(UA_EventLoop *public_el, UA_UInt32 timeout) { |
335 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
336 | 0 | UA_LOCK(&el->elMutex); |
337 | |
|
338 | 0 | if(el->executing) { |
339 | 0 | UA_LOG_ERROR(el->eventLoop.logger, |
340 | 0 | UA_LOGCATEGORY_EVENTLOOP, |
341 | 0 | "Cannot run EventLoop from the run method itself"); |
342 | 0 | UA_UNLOCK(&el->elMutex); |
343 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
344 | 0 | } |
345 | | |
346 | 0 | el->executing = true; |
347 | |
|
348 | 0 | if(el->eventLoop.state == UA_EVENTLOOPSTATE_FRESH || |
349 | 0 | el->eventLoop.state == UA_EVENTLOOPSTATE_STOPPED) { |
350 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
351 | 0 | "Cannot run a stopped EventLoop"); |
352 | 0 | el->executing = false; |
353 | 0 | UA_UNLOCK(&el->elMutex); |
354 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
355 | 0 | } |
356 | | |
357 | 0 | UA_LOG_TRACE(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
358 | 0 | "Iterate the EventLoop"); |
359 | | |
360 | | /* Process cyclic callbacks */ |
361 | 0 | UA_DateTime dateBefore = |
362 | 0 | el->eventLoop.dateTime_nowMonotonic(&el->eventLoop); |
363 | |
|
364 | 0 | UA_DateTime dateNext = UA_Timer_process(&el->timer, dateBefore); |
365 | | |
366 | | /* Process delayed callbacks here: |
367 | | * - Removes closed sockets already here instead of polling them again. |
368 | | * - The timeout for polling is selected to be ready in time for the next |
369 | | * cyclic callback. So we want to do little work between the timeout |
370 | | * running out and executing the due cyclic callbacks. */ |
371 | 0 | UA_EventLoopPOSIX_processDelayed(el); |
372 | | |
373 | | /* A delayed callback could create another delayed callback (or re-add |
374 | | * itself). In that case we don't want to wait (indefinitely) for an event |
375 | | * to happen. Process queued events but don't sleep. Then process the |
376 | | * delayed callbacks in the next iteration. */ |
377 | 0 | if(UA_atomic_load(&el->delayedHead1) != NULL && |
378 | 0 | UA_atomic_load(&el->delayedHead2) != NULL) |
379 | 0 | timeout = 0; |
380 | | |
381 | | /* Compute the remaining time */ |
382 | 0 | UA_DateTime maxDate = dateBefore + (timeout * UA_DATETIME_MSEC); |
383 | 0 | if(dateNext > maxDate) |
384 | 0 | dateNext = maxDate; |
385 | 0 | UA_DateTime listenTimeout = |
386 | 0 | dateNext - el->eventLoop.dateTime_nowMonotonic(&el->eventLoop); |
387 | 0 | if(listenTimeout < 0) |
388 | 0 | listenTimeout = 0; |
389 | | |
390 | | /* Listen on the active file-descriptors (sockets) from the |
391 | | * ConnectionManagers */ |
392 | 0 | UA_StatusCode rv = UA_EventLoopPOSIX_pollFDs(el, listenTimeout); |
393 | | |
394 | | /* Check if the last EventSource was successfully stopped */ |
395 | 0 | if(el->eventLoop.state == UA_EVENTLOOPSTATE_STOPPING) |
396 | 0 | checkClosed(el); |
397 | |
|
398 | 0 | el->executing = false; |
399 | 0 | UA_UNLOCK(&el->elMutex); |
400 | 0 | return rv; |
401 | 0 | } |
402 | | |
403 | | /*****************************/ |
404 | | /* Registering Event Sources */ |
405 | | /*****************************/ |
406 | | |
407 | | UA_StatusCode |
408 | | UA_EventLoopPOSIX_registerEventSource(UA_EventLoop *public_el, |
409 | 0 | UA_EventSource *es) { |
410 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
411 | 0 | UA_LOCK(&el->elMutex); |
412 | | |
413 | | /* Already registered? */ |
414 | 0 | if(es->state != UA_EVENTSOURCESTATE_FRESH) { |
415 | 0 | UA_LOG_ERROR(el->eventLoop.logger, UA_LOGCATEGORY_NETWORK, |
416 | 0 | "Cannot register the EventSource \"%.*s\": " |
417 | 0 | "already registered", |
418 | 0 | (int)es->name.length, (char*)es->name.data); |
419 | 0 | UA_UNLOCK(&el->elMutex); |
420 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
421 | 0 | } |
422 | | |
423 | | /* Add to linked list */ |
424 | 0 | es->next = el->eventLoop.eventSources; |
425 | 0 | el->eventLoop.eventSources = es; |
426 | |
|
427 | 0 | es->eventLoop = &el->eventLoop; |
428 | 0 | es->state = UA_EVENTSOURCESTATE_STOPPED; |
429 | | |
430 | | /* Start if the entire EventLoop is started */ |
431 | 0 | UA_StatusCode res = UA_STATUSCODE_GOOD; |
432 | 0 | if(el->eventLoop.state == UA_EVENTLOOPSTATE_STARTED) |
433 | 0 | res = es->start(es); |
434 | |
|
435 | 0 | UA_UNLOCK(&el->elMutex); |
436 | 0 | return res; |
437 | 0 | } |
438 | | |
439 | | UA_StatusCode |
440 | | UA_EventLoopPOSIX_deregisterEventSource(UA_EventLoop *public_el, |
441 | 0 | UA_EventSource *es) { |
442 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
443 | 0 | UA_LOCK(&el->elMutex); |
444 | |
|
445 | 0 | if(es->state != UA_EVENTSOURCESTATE_STOPPED) { |
446 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
447 | 0 | "Cannot deregister the EventSource %.*s: " |
448 | 0 | "Has to be stopped first", |
449 | 0 | (int)es->name.length, es->name.data); |
450 | 0 | UA_UNLOCK(&el->elMutex); |
451 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
452 | 0 | } |
453 | | |
454 | | /* Remove from the linked list */ |
455 | 0 | UA_EventSource **s = &el->eventLoop.eventSources; |
456 | 0 | while(*s) { |
457 | 0 | if(*s == es) { |
458 | 0 | *s = es->next; |
459 | 0 | break; |
460 | 0 | } |
461 | 0 | s = &(*s)->next; |
462 | 0 | } |
463 | | |
464 | | /* Set the state to non-registered */ |
465 | 0 | es->state = UA_EVENTSOURCESTATE_FRESH; |
466 | |
|
467 | 0 | UA_UNLOCK(&el->elMutex); |
468 | 0 | return UA_STATUSCODE_GOOD; |
469 | 0 | } |
470 | | |
471 | | /***************/ |
472 | | /* Time Domain */ |
473 | | /***************/ |
474 | | |
475 | | UA_DateTime |
476 | 0 | UA_EventLoopPOSIX_DateTime_now(UA_EventLoop *el) { |
477 | 0 | UA_EventLoopPOSIX *pel = (UA_EventLoopPOSIX*)el; |
478 | 0 | struct timespec ts; |
479 | 0 | int res = clock_gettime((clockid_t)pel->clockSource, &ts); |
480 | 0 | if(UA_UNLIKELY(res != 0)) |
481 | 0 | return 0; |
482 | 0 | return (ts.tv_sec * UA_DATETIME_SEC) + (ts.tv_nsec / 100) + UA_DATETIME_UNIX_EPOCH; |
483 | 0 | } |
484 | | |
485 | | UA_DateTime |
486 | 0 | UA_EventLoopPOSIX_DateTime_nowMonotonic(UA_EventLoop *el) { |
487 | 0 | UA_EventLoopPOSIX *pel = (UA_EventLoopPOSIX*)el; |
488 | 0 | struct timespec ts; |
489 | 0 | int res = clock_gettime((clockid_t)pel->clockSourceMonotonic, &ts); |
490 | 0 | if(UA_UNLIKELY(res != 0)) |
491 | 0 | return 0; |
492 | | /* Also add the unix epoch for the monotonic clock. So we get a "normal" |
493 | | * output when a "normal" source is configured. */ |
494 | 0 | return (ts.tv_sec * UA_DATETIME_SEC) + (ts.tv_nsec / 100) + UA_DATETIME_UNIX_EPOCH; |
495 | 0 | } |
496 | | |
497 | | UA_Int64 |
498 | 0 | UA_EventLoopPOSIX_DateTime_localTimeUtcOffset(UA_EventLoop *el) { |
499 | | /* TODO: Fix for custom clock sources */ |
500 | 0 | return UA_DateTime_localTimeUtcOffset(); |
501 | 0 | } |
502 | | |
503 | | /*************************/ |
504 | | /* Initialize and Delete */ |
505 | | /*************************/ |
506 | | |
507 | | static UA_StatusCode |
508 | 0 | UA_EventLoopPOSIX_free(UA_EventLoop *public_el) { |
509 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
510 | 0 | UA_LOCK(&el->elMutex); |
511 | | |
512 | | /* Check if the EventLoop can be deleted */ |
513 | 0 | if(el->eventLoop.state != UA_EVENTLOOPSTATE_STOPPED && |
514 | 0 | el->eventLoop.state != UA_EVENTLOOPSTATE_FRESH) { |
515 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
516 | 0 | "Cannot delete a running EventLoop"); |
517 | 0 | UA_UNLOCK(&el->elMutex); |
518 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
519 | 0 | } |
520 | | |
521 | | /* Deregister and delete all the EventSources */ |
522 | 0 | while(el->eventLoop.eventSources) { |
523 | 0 | UA_EventSource *es = el->eventLoop.eventSources; |
524 | 0 | UA_EventLoopPOSIX_deregisterEventSource(public_el, es); |
525 | 0 | es->free(es); |
526 | 0 | } |
527 | | |
528 | | /* Remove the repeated timed callbacks */ |
529 | 0 | UA_Timer_clear(&el->timer); |
530 | |
|
531 | | #ifdef UA_ENABLE_LWS |
532 | | /* The LWS context can only be destroyed synchronously outside an LWS |
533 | | * service callback. All EventSources have released it at this point. */ |
534 | | UA_LWS_destroyContext(public_el); |
535 | | #endif |
536 | | |
537 | | /* Process remaining delayed callbacks */ |
538 | 0 | UA_EventLoopPOSIX_processDelayed(el); |
539 | | |
540 | |
|
541 | 0 | UA_KeyValueMap_clear(&el->eventLoop.params); |
542 | | |
543 | | /* Clean up */ |
544 | 0 | UA_UNLOCK(&el->elMutex); |
545 | 0 | UA_LOCK_DESTROY(&el->elMutex); |
546 | 0 | UA_free(el); |
547 | 0 | return UA_STATUSCODE_GOOD; |
548 | 0 | } |
549 | | |
550 | | void |
551 | 0 | UA_EventLoopPOSIX_lock(UA_EventLoop *public_el) { |
552 | 0 | UA_LOCK(&((UA_EventLoopPOSIX*)public_el)->elMutex); |
553 | 0 | } |
554 | | void |
555 | 0 | UA_EventLoopPOSIX_unlock(UA_EventLoop *public_el) { |
556 | 0 | UA_UNLOCK(&((UA_EventLoopPOSIX*)public_el)->elMutex); |
557 | 0 | } |
558 | | |
559 | | /* Forward declarations for the FD-polling backend implementations further |
560 | | * down in this file (select or epoll, chosen at compile time). */ |
561 | | #if defined(UA_HAVE_EPOLL) |
562 | | static UA_StatusCode registerFD_epoll(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd); |
563 | | static UA_StatusCode modifyFD_epoll(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd); |
564 | | static void deregisterFD_epoll(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd); |
565 | | #else |
566 | | static UA_StatusCode registerFD_select(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd); |
567 | | static UA_StatusCode modifyFD_select(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd); |
568 | | static void deregisterFD_select(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd); |
569 | | #endif |
570 | | |
571 | | UA_EventLoop * |
572 | 0 | UA_EventLoop_new_POSIX(const UA_Logger *logger) { |
573 | |
|
574 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*) |
575 | 0 | UA_calloc(1, sizeof(UA_EventLoopPOSIX)); |
576 | 0 | if(!el) |
577 | 0 | return NULL; |
578 | | |
579 | 0 | UA_LOCK_INIT(&el->elMutex); |
580 | 0 | UA_Timer_init(&el->timer); |
581 | | |
582 | | /* Initialize the queue */ |
583 | 0 | el->delayedTail = &el->delayedHead1; |
584 | 0 | el->delayedHead2 = (UA_DelayedCallback*)0x01; /* sentinel value */ |
585 | | |
586 | | /* Set the public EventLoop content */ |
587 | 0 | el->eventLoop.logger = logger; |
588 | | |
589 | | /* Initialize the clock source to the default */ |
590 | 0 | el->clockSource = CLOCK_REALTIME; |
591 | 0 | # ifdef CLOCK_MONOTONIC_RAW |
592 | 0 | el->clockSourceMonotonic = CLOCK_MONOTONIC_RAW; |
593 | | # else |
594 | | el->clockSourceMonotonic = CLOCK_MONOTONIC; |
595 | | # endif |
596 | | |
597 | | /* Set the method pointers for the interface */ |
598 | 0 | el->eventLoop.start = UA_EventLoopPOSIX_start; |
599 | 0 | el->eventLoop.stop = UA_EventLoopPOSIX_stop; |
600 | 0 | el->eventLoop.free = UA_EventLoopPOSIX_free; |
601 | 0 | el->eventLoop.run = UA_EventLoopPOSIX_run; |
602 | 0 | el->eventLoop.cancel = UA_EventLoopPOSIX_cancel; |
603 | |
|
604 | 0 | el->eventLoop.dateTime_now = UA_EventLoopPOSIX_DateTime_now; |
605 | 0 | el->eventLoop.dateTime_nowMonotonic = |
606 | 0 | UA_EventLoopPOSIX_DateTime_nowMonotonic; |
607 | 0 | el->eventLoop.dateTime_localTimeUtcOffset = |
608 | 0 | UA_EventLoopPOSIX_DateTime_localTimeUtcOffset; |
609 | |
|
610 | 0 | el->eventLoop.nextTimer = UA_EventLoopPOSIX_nextTimer; |
611 | 0 | el->eventLoop.addTimer = UA_EventLoopPOSIX_addTimer; |
612 | 0 | el->eventLoop.modifyTimer = UA_EventLoopPOSIX_modifyTimer; |
613 | 0 | el->eventLoop.removeTimer = UA_EventLoopPOSIX_removeTimer; |
614 | 0 | el->eventLoop.addDelayedCallback = UA_EventLoopPOSIX_addDelayedCallback; |
615 | 0 | el->eventLoop.removeDelayedCallback = UA_EventLoopPOSIX_removeDelayedCallback; |
616 | |
|
617 | 0 | el->eventLoop.registerEventSource = UA_EventLoopPOSIX_registerEventSource; |
618 | 0 | el->eventLoop.deregisterEventSource = UA_EventLoopPOSIX_deregisterEventSource; |
619 | |
|
620 | 0 | el->eventLoop.lock = UA_EventLoopPOSIX_lock; |
621 | 0 | el->eventLoop.unlock = UA_EventLoopPOSIX_unlock; |
622 | | |
623 | | /* Select the FD polling backend */ |
624 | 0 | #if defined(UA_HAVE_EPOLL) |
625 | 0 | el->registerFD = registerFD_epoll; |
626 | 0 | el->modifyFD = modifyFD_epoll; |
627 | 0 | el->deregisterFD = deregisterFD_epoll; |
628 | | #else |
629 | | el->registerFD = registerFD_select; |
630 | | el->modifyFD = modifyFD_select; |
631 | | el->deregisterFD = deregisterFD_select; |
632 | | #endif |
633 | |
|
634 | 0 | return &el->eventLoop; |
635 | 0 | } |
636 | | |
637 | | /***************************/ |
638 | | /* Network Buffer Handling */ |
639 | | /***************************/ |
640 | | |
641 | | UA_StatusCode |
642 | | UA_EventLoopPOSIX_allocNetworkBuffer(UA_ConnectionManager *cm, |
643 | | uintptr_t connectionId, |
644 | | UA_ByteString *buf, |
645 | 0 | size_t bufSize) { |
646 | 0 | UA_POSIXConnectionManager *pcm = (UA_POSIXConnectionManager*)cm; |
647 | | /* Reuse the static tx buffer; fall back to allocation for larger messages. */ |
648 | 0 | if(pcm->txBuffer.length < bufSize) |
649 | 0 | return UA_ByteString_allocBuffer(buf, bufSize); |
650 | 0 | *buf = pcm->txBuffer; |
651 | 0 | buf->length = bufSize; |
652 | 0 | return UA_STATUSCODE_GOOD; |
653 | 0 | } |
654 | | |
655 | | void |
656 | | UA_EventLoopPOSIX_freeNetworkBuffer(UA_ConnectionManager *cm, |
657 | | uintptr_t connectionId, |
658 | 0 | UA_ByteString *buf) { |
659 | 0 | UA_POSIXConnectionManager *pcm = (UA_POSIXConnectionManager*)cm; |
660 | 0 | if(pcm->txBuffer.data == buf->data) |
661 | 0 | UA_ByteString_init(buf); |
662 | 0 | else |
663 | 0 | UA_ByteString_clear(buf); |
664 | 0 | } |
665 | | |
666 | | UA_StatusCode |
667 | 0 | UA_EventLoopPOSIX_allocateStaticBuffers(UA_POSIXConnectionManager *pcm) { |
668 | 0 | UA_StatusCode res = |
669 | 0 | UA_EventLoopCommon_allocStaticBuffer(&pcm->cm.eventSource.params, |
670 | 0 | UA_QUALIFIEDNAME(0, "recv-bufsize"), |
671 | 0 | 1u << 16, /* The default is 64 kb */ |
672 | 0 | &pcm->rxBuffer); |
673 | | |
674 | | /* Default the tx buffer to the rx size so a dedicated static send buffer |
675 | | * always exists. This avoids a malloc/free on every send without reusing |
676 | | * the rx buffer (which may still hold unprocessed received data). */ |
677 | 0 | res |= UA_EventLoopCommon_allocStaticBuffer(&pcm->cm.eventSource.params, |
678 | 0 | UA_QUALIFIEDNAME(0, "send-bufsize"), |
679 | 0 | (UA_UInt32)pcm->rxBuffer.length, |
680 | 0 | &pcm->txBuffer); |
681 | 0 | return res; |
682 | 0 | } |
683 | | |
684 | | /******************/ |
685 | | /* Socket Options */ |
686 | | /******************/ |
687 | | |
688 | | enum ZIP_CMP |
689 | 0 | cmpFD(const UA_FD *a, const UA_FD *b) { |
690 | 0 | if(*a == *b) |
691 | 0 | return ZIP_CMP_EQ; |
692 | 0 | return (*a < *b) ? ZIP_CMP_LESS : ZIP_CMP_MORE; |
693 | 0 | } |
694 | | |
695 | | UA_StatusCode |
696 | 0 | UA_EventLoopPOSIX_setNonBlocking(UA_FD sockfd) { |
697 | 0 | int opts = fcntl(sockfd, F_GETFL); |
698 | 0 | if(opts < 0 || fcntl(sockfd, F_SETFL, opts | O_NONBLOCK) < 0) |
699 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
700 | 0 | return UA_STATUSCODE_GOOD; |
701 | 0 | } |
702 | | |
703 | | UA_StatusCode |
704 | 0 | UA_EventLoopPOSIX_setNoSigPipe(UA_FD sockfd) { |
705 | | #ifdef SO_NOSIGPIPE |
706 | | int val = 1; |
707 | | int res = UA_setsockopt(sockfd, SOL_SOCKET, SO_NOSIGPIPE, &val, sizeof(val)); |
708 | | if(res < 0) |
709 | | return UA_STATUSCODE_BADINTERNALERROR; |
710 | | #endif |
711 | 0 | return UA_STATUSCODE_GOOD; |
712 | 0 | } |
713 | | |
714 | | UA_StatusCode |
715 | 0 | UA_EventLoopPOSIX_setReusable(UA_FD sockfd) { |
716 | 0 | int enableReuseVal = 1; |
717 | 0 | int res = UA_setsockopt(sockfd, SOL_SOCKET, SO_REUSEADDR, |
718 | 0 | (const char*)&enableReuseVal, sizeof(enableReuseVal)); |
719 | 0 | res |= UA_setsockopt(sockfd, SOL_SOCKET, SO_REUSEPORT, |
720 | 0 | (const char*)&enableReuseVal, sizeof(enableReuseVal)); |
721 | 0 | return (res == 0) ? UA_STATUSCODE_GOOD : UA_STATUSCODE_BADINTERNALERROR; |
722 | 0 | } |
723 | | |
724 | | /************************/ |
725 | | /* Select / epoll Logic */ |
726 | | /************************/ |
727 | | |
728 | | /* Re-arm the self-pipe socket for the next signal by reading from it */ |
729 | | static void |
730 | 0 | flushSelfPipe(UA_SOCKET s) { |
731 | 0 | char buf[128]; |
732 | 0 | int i; |
733 | 0 | do { |
734 | 0 | i = UA_recv(s, buf, 128, 0); |
735 | 0 | } while(i > 0); |
736 | 0 | } |
737 | | |
738 | | #if !defined(UA_HAVE_EPOLL) |
739 | | |
740 | | static UA_StatusCode |
741 | | registerFD_select(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd) { |
742 | | UA_LOCK_ASSERT(&el->elMutex); |
743 | | UA_LOG_DEBUG(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
744 | | "Registering fd: %u", (unsigned)rfd->fd); |
745 | | |
746 | | /* Realloc */ |
747 | | UA_RegisteredFD **fds_tmp = (UA_RegisteredFD**) |
748 | | UA_realloc(el->fds, sizeof(UA_RegisteredFD*) * (el->fdsSize + 1)); |
749 | | if(!fds_tmp) { |
750 | | return UA_STATUSCODE_BADOUTOFMEMORY; |
751 | | } |
752 | | el->fds = fds_tmp; |
753 | | |
754 | | /* Add to the last entry */ |
755 | | el->fds[el->fdsSize] = rfd; |
756 | | el->fdsSize++; |
757 | | return UA_STATUSCODE_GOOD; |
758 | | } |
759 | | |
760 | | static UA_StatusCode |
761 | | modifyFD_select(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd) { |
762 | | /* Do nothing, it is enough if the data was changed in the rfd */ |
763 | | UA_LOCK_ASSERT(&el->elMutex); |
764 | | return UA_STATUSCODE_GOOD; |
765 | | } |
766 | | |
767 | | static void |
768 | | deregisterFD_select(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd) { |
769 | | UA_LOCK_ASSERT(&el->elMutex); |
770 | | UA_LOG_DEBUG(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
771 | | "Unregistering fd: %u", (unsigned)rfd->fd); |
772 | | |
773 | | /* Find the entry */ |
774 | | size_t i = 0; |
775 | | for(; i < el->fdsSize; i++) { |
776 | | if(el->fds[i] == rfd) |
777 | | break; |
778 | | } |
779 | | |
780 | | /* Not found? */ |
781 | | if(i == el->fdsSize) |
782 | | return; |
783 | | |
784 | | if(el->fdsSize > 1) { |
785 | | /* Move the last entry in the ith slot and realloc. */ |
786 | | el->fdsSize--; |
787 | | el->fds[i] = el->fds[el->fdsSize]; |
788 | | UA_RegisteredFD **fds_tmp = (UA_RegisteredFD**) |
789 | | UA_realloc(el->fds, sizeof(UA_RegisteredFD*) * el->fdsSize); |
790 | | /* if realloc fails the fds are still in a correct state with |
791 | | * possibly lost memory, so failing silently here is ok */ |
792 | | if(fds_tmp) |
793 | | el->fds = fds_tmp; |
794 | | } else { |
795 | | /* Remove the last entry */ |
796 | | UA_free(el->fds); |
797 | | el->fds = NULL; |
798 | | el->fdsSize = 0; |
799 | | } |
800 | | } |
801 | | |
802 | | static UA_FD |
803 | | setFDSets(UA_EventLoopPOSIX *el, fd_set *readset, fd_set *writeset, fd_set *errset) { |
804 | | UA_LOCK_ASSERT(&el->elMutex); |
805 | | |
806 | | FD_ZERO(readset); |
807 | | FD_ZERO(writeset); |
808 | | FD_ZERO(errset); |
809 | | |
810 | | /* Always listen on the read-end of the pipe */ |
811 | | UA_FD highestfd = el->selfpipe[0]; |
812 | | FD_SET(el->selfpipe[0], readset); |
813 | | |
814 | | for(size_t i = 0; i < el->fdsSize; i++) { |
815 | | UA_FD currentFD = el->fds[i]->fd; |
816 | | |
817 | | /* Add to the fd_sets */ |
818 | | if(el->fds[i]->listenEvents & UA_FDEVENT_IN) |
819 | | FD_SET(currentFD, readset); |
820 | | if(el->fds[i]->listenEvents & UA_FDEVENT_OUT) |
821 | | FD_SET(currentFD, writeset); |
822 | | |
823 | | /* Always return errors */ |
824 | | FD_SET(currentFD, errset); |
825 | | |
826 | | /* Highest fd? */ |
827 | | if(currentFD > highestfd) |
828 | | highestfd = currentFD; |
829 | | } |
830 | | return highestfd; |
831 | | } |
832 | | |
833 | | UA_StatusCode |
834 | | UA_EventLoopPOSIX_pollFDs(UA_EventLoopPOSIX *el, UA_DateTime listenTimeout) { |
835 | | UA_assert(listenTimeout >= 0); |
836 | | UA_LOCK_ASSERT(&el->elMutex); |
837 | | |
838 | | fd_set readset, writeset, errset; |
839 | | UA_FD highestfd = setFDSets(el, &readset, &writeset, &errset); |
840 | | |
841 | | /* Nothing to do? */ |
842 | | if(highestfd == UA_INVALID_FD) { |
843 | | UA_LOG_TRACE(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
844 | | "No valid FDs for processing"); |
845 | | return UA_STATUSCODE_GOOD; |
846 | | } |
847 | | |
848 | | struct timeval tmptv = { |
849 | | (time_t)(listenTimeout / UA_DATETIME_SEC), |
850 | | (suseconds_t)((listenTimeout % UA_DATETIME_SEC) / UA_DATETIME_USEC) |
851 | | }; |
852 | | |
853 | | UA_UNLOCK(&el->elMutex); |
854 | | int selectStatus = UA_select(highestfd+1, &readset, &writeset, &errset, &tmptv); |
855 | | UA_LOCK(&el->elMutex); |
856 | | if(selectStatus < 0) { |
857 | | /* We will retry, only log the error */ |
858 | | UA_LOG_SOCKET_ERRNO_WRAP( |
859 | | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
860 | | "Error during select: %s", errno_str)); |
861 | | return UA_STATUSCODE_GOOD; |
862 | | } |
863 | | |
864 | | /* The self-pipe has received. Clear the buffer by reading. */ |
865 | | if(UA_UNLIKELY(FD_ISSET(el->selfpipe[0], &readset))) |
866 | | flushSelfPipe(el->selfpipe[0]); |
867 | | |
868 | | /* Loop over all registered FD to see if an event arrived. Yes, this is why |
869 | | * select is slow for many open sockets. */ |
870 | | for(size_t i = 0; i < el->fdsSize; i++) { |
871 | | UA_RegisteredFD *rfd = el->fds[i]; |
872 | | |
873 | | /* The rfd is already registered for removal. Don't process incoming |
874 | | * events any longer. */ |
875 | | if(rfd->dc.callback) |
876 | | continue; |
877 | | |
878 | | /* Event signaled for the fd? */ |
879 | | short event = 0; |
880 | | if(FD_ISSET(rfd->fd, &readset)) { |
881 | | event |= UA_FDEVENT_IN; |
882 | | } |
883 | | if(FD_ISSET(rfd->fd, &writeset)) { |
884 | | event |= UA_FDEVENT_OUT; |
885 | | } |
886 | | if(!event && FD_ISSET(rfd->fd, &errset)) { |
887 | | event = UA_FDEVENT_ERR; |
888 | | } |
889 | | if(!event) { |
890 | | continue; |
891 | | } |
892 | | |
893 | | UA_LOG_DEBUG(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
894 | | "Processing event %u on fd %u", (unsigned)event, |
895 | | (unsigned)rfd->fd); |
896 | | |
897 | | /* Call the EventSource callback */ |
898 | | rfd->eventSourceCB(rfd->es, rfd, event); |
899 | | |
900 | | /* The fd has removed itself */ |
901 | | if(i >= el->fdsSize || rfd != el->fds[i]) |
902 | | i--; |
903 | | } |
904 | | return UA_STATUSCODE_GOOD; |
905 | | } |
906 | | |
907 | | #else /* defined(UA_HAVE_EPOLL) */ |
908 | | |
909 | | static UA_StatusCode |
910 | 0 | registerFD_epoll(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd) { |
911 | 0 | struct epoll_event event; |
912 | 0 | memset(&event, 0, sizeof(struct epoll_event)); |
913 | 0 | event.data.ptr = rfd; |
914 | 0 | event.events = 0; |
915 | 0 | if(rfd->listenEvents & UA_FDEVENT_IN) |
916 | 0 | event.events |= EPOLLIN; |
917 | 0 | if(rfd->listenEvents & UA_FDEVENT_OUT) |
918 | 0 | event.events |= EPOLLOUT; |
919 | |
|
920 | 0 | int err = epoll_ctl(el->epollfd, EPOLL_CTL_ADD, rfd->fd, &event); |
921 | 0 | if(err != 0) { |
922 | 0 | UA_LOG_SOCKET_ERRNO_WRAP( |
923 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_NETWORK, |
924 | 0 | "TCP %u\t| Could not register for epoll (%s)", |
925 | 0 | rfd->fd, errno_str)); |
926 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
927 | 0 | } |
928 | 0 | return UA_STATUSCODE_GOOD; |
929 | 0 | } |
930 | | |
931 | | static UA_StatusCode |
932 | 0 | modifyFD_epoll(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd) { |
933 | 0 | struct epoll_event event; |
934 | 0 | memset(&event, 0, sizeof(struct epoll_event)); |
935 | 0 | event.data.ptr = rfd; |
936 | 0 | event.events = 0; |
937 | 0 | if(rfd->listenEvents & UA_FDEVENT_IN) |
938 | 0 | event.events |= EPOLLIN; |
939 | 0 | if(rfd->listenEvents & UA_FDEVENT_OUT) |
940 | 0 | event.events |= EPOLLOUT; |
941 | |
|
942 | 0 | int err = epoll_ctl(el->epollfd, EPOLL_CTL_MOD, rfd->fd, &event); |
943 | 0 | if(err != 0) { |
944 | 0 | UA_LOG_SOCKET_ERRNO_WRAP( |
945 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_NETWORK, |
946 | 0 | "TCP %u\t| Could not modify for epoll (%s)", |
947 | 0 | rfd->fd, errno_str)); |
948 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
949 | 0 | } |
950 | 0 | return UA_STATUSCODE_GOOD; |
951 | 0 | } |
952 | | |
953 | | static void |
954 | 0 | deregisterFD_epoll(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd) { |
955 | 0 | int res = epoll_ctl(el->epollfd, EPOLL_CTL_DEL, rfd->fd, NULL); |
956 | 0 | if(res != 0) { |
957 | 0 | UA_LOG_SOCKET_ERRNO_WRAP( |
958 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_NETWORK, |
959 | 0 | "TCP %u\t| Could not deregister from epoll (%s)", |
960 | 0 | rfd->fd, errno_str)); |
961 | 0 | } |
962 | 0 | } |
963 | | |
964 | | UA_StatusCode |
965 | 0 | UA_EventLoopPOSIX_pollFDs(UA_EventLoopPOSIX *el, UA_DateTime listenTimeout) { |
966 | 0 | UA_assert(listenTimeout >= 0); |
967 | | |
968 | | /* If there is a positive timeout, wait at least one millisecond, the |
969 | | * minimum for blocking epoll_wait. This prevents a busy-loop, as the |
970 | | * open62541 library allows even smaller timeouts, which can result in a |
971 | | * zero timeout due to rounding to an integer here. */ |
972 | 0 | int timeout = (int)(listenTimeout / UA_DATETIME_MSEC); |
973 | 0 | if(timeout == 0 && listenTimeout > 0) |
974 | 0 | timeout = 1; |
975 | | |
976 | | /* Poll the registered sockets */ |
977 | 0 | struct epoll_event epoll_events[64]; |
978 | 0 | UA_UNLOCK(&el->elMutex); |
979 | 0 | int events = epoll_wait(el->epollfd, epoll_events, 64, timeout); |
980 | 0 | UA_LOCK(&el->elMutex); |
981 | | |
982 | | /* TODO: Replace with pwait2 for higher-precision timeouts once this is |
983 | | * available in the standard library. |
984 | | * |
985 | | * struct timespec precisionTimeout = { |
986 | | * (long)(listenTimeout / UA_DATETIME_SEC), |
987 | | * (long)((listenTimeout % UA_DATETIME_SEC) * 100) |
988 | | * }; |
989 | | * int events = epoll_pwait2(epollfd, epoll_events, 64, |
990 | | * precisionTimeout, NULL); */ |
991 | | |
992 | | /* Handle error conditions */ |
993 | 0 | if(events == -1) { |
994 | 0 | if(errno == EINTR) { |
995 | | /* We will retry, only log the error */ |
996 | 0 | UA_LOG_DEBUG(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
997 | 0 | "Timeout during poll"); |
998 | 0 | return UA_STATUSCODE_GOOD; |
999 | 0 | } |
1000 | 0 | UA_LOG_SOCKET_ERRNO_WRAP( |
1001 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_NETWORK, |
1002 | 0 | "TCP\t| Error %s, closing the server socket", |
1003 | 0 | errno_str)); |
1004 | 0 | return UA_STATUSCODE_BADINTERNALERROR; |
1005 | 0 | } |
1006 | | |
1007 | | /* Process all received events */ |
1008 | 0 | for(int i = 0; i < events; i++) { |
1009 | 0 | UA_RegisteredFD *rfd = (UA_RegisteredFD*)epoll_events[i].data.ptr; |
1010 | | |
1011 | | /* The self-pipe has received */ |
1012 | 0 | if(!rfd) { |
1013 | 0 | flushSelfPipe(el->selfpipe[0]); |
1014 | 0 | continue; |
1015 | 0 | } |
1016 | | |
1017 | | /* The rfd is already registered for removal. Don't process incoming |
1018 | | * events any longer. */ |
1019 | 0 | if(rfd->dc.callback) |
1020 | 0 | continue; |
1021 | | |
1022 | | /* Forward both directions so pending input cannot starve writes. */ |
1023 | 0 | short revent = 0; |
1024 | 0 | if(epoll_events[i].events & EPOLLIN) |
1025 | 0 | revent |= UA_FDEVENT_IN; |
1026 | 0 | if(epoll_events[i].events & EPOLLOUT) |
1027 | 0 | revent |= UA_FDEVENT_OUT; |
1028 | 0 | if(!revent) |
1029 | 0 | revent = UA_FDEVENT_ERR; |
1030 | | |
1031 | | /* Call the EventSource callback */ |
1032 | 0 | rfd->eventSourceCB(rfd->es, rfd, revent); |
1033 | 0 | } |
1034 | 0 | return UA_STATUSCODE_GOOD; |
1035 | 0 | } |
1036 | | |
1037 | | #endif /* defined(UA_HAVE_EPOLL) */ |
1038 | | |
1039 | | /* Thin wrappers dispatching through the backend selected in |
1040 | | * UA_EventLoop_new_POSIX / UA_EventLoop_new_GLib. This is what |
1041 | | * ConnectionManagers (TCP, UDP, Ethernet, ...) actually call -- they do not |
1042 | | * need to know which backend is behind a given EventLoop instance. */ |
1043 | | |
1044 | | UA_StatusCode |
1045 | 0 | UA_EventLoopPOSIX_registerFD(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd) { |
1046 | 0 | return el->registerFD(el, rfd); |
1047 | 0 | } |
1048 | | |
1049 | | UA_StatusCode |
1050 | 0 | UA_EventLoopPOSIX_modifyFD(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd) { |
1051 | 0 | return el->modifyFD(el, rfd); |
1052 | 0 | } |
1053 | | |
1054 | | void |
1055 | 0 | UA_EventLoopPOSIX_deregisterFD(UA_EventLoopPOSIX *el, UA_RegisteredFD *rfd) { |
1056 | 0 | el->deregisterFD(el, rfd); |
1057 | 0 | } |
1058 | | |
1059 | 0 | int UA_EventLoopPOSIX_pipe(UA_FD fds[2]) { |
1060 | 0 | int err = socketpair(AF_UNIX, SOCK_STREAM, 0, fds); |
1061 | 0 | if(err != 0) |
1062 | 0 | return err; |
1063 | 0 | UA_EventLoopPOSIX_setNonBlocking(fds[0]); |
1064 | 0 | UA_EventLoopPOSIX_setNonBlocking(fds[1]); |
1065 | 0 | UA_EventLoopPOSIX_setNoSigPipe(fds[0]); |
1066 | 0 | UA_EventLoopPOSIX_setNoSigPipe(fds[1]); |
1067 | 0 | return 0; |
1068 | 0 | } |
1069 | | |
1070 | | void |
1071 | 0 | UA_EventLoopPOSIX_cancel(UA_EventLoop *public_el) { |
1072 | 0 | UA_EventLoopPOSIX *el = (UA_EventLoopPOSIX*)public_el; |
1073 | | /* Nothing to do if the EventLoop is not executing */ |
1074 | 0 | if(!el->executing) |
1075 | 0 | return; |
1076 | | |
1077 | | /* Trigger the self-pipe */ |
1078 | 0 | int err = (int)UA_send(el->selfpipe[1], ".", 1, 0); |
1079 | 0 | if(err <= 0) { |
1080 | | UA_LOG_SOCKET_ERRNO_WRAP( |
1081 | 0 | UA_LOG_WARNING(el->eventLoop.logger, UA_LOGCATEGORY_EVENTLOOP, |
1082 | 0 | "Eventloop\t| Error signaling self-pipe (%s)", errno_str)); |
1083 | 0 | } |
1084 | 0 | } |
1085 | | |
1086 | | #endif |