Coverage Report

Created: 2026-09-27 06:56

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