The FreeRADIUS server $Id: f3670dba8951ca10eb4948feb3dc3db9423a334f $
Loading...
Searching...
No Matches
coord.c
Go to the documentation of this file.
1/*
2 * This program is free software; you can redistribute it and/or modify
3 * it under the terms of the GNU General Public License as published by
4 * the Free Software Foundation; either version 2 of the License, or
5 * (at your option) any later version.
6 *
7 * This program is distributed in the hope that it will be useful,
8 * but WITHOUT ANY WARRANTY; without even the implied warranty of
9 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
10 * GNU General Public License for more details.
11 *
12 * You should have received a copy of the GNU General Public License
13 * along with this program; if not, write to the Free Software
14 * Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301, USA
15 */
16
17/**
18 * $Id: 836b13654accd2a0cfb926e1e81523b26bed6b75 $
19 *
20 * @brief Coordination thread management
21 * @file io/coord.c
22 *
23 * @copyright 2026 Network RADIUS SAS (legal@networkradius.com)
24 */
25RCSID("$Id: 836b13654accd2a0cfb926e1e81523b26bed6b75 $")
26
27#include <freeradius-devel/io/listen.h>
28#include <freeradius-devel/io/schedule.h>
29#include <freeradius-devel/io/thread.h>
30#include <freeradius-devel/io/coord_priv.h>
31#include <freeradius-devel/unlang/base.h>
32#include <freeradius-devel/util/syserror.h>
33
34#include <stdalign.h>
35
36#define FR_CONTROL_ID_COORD_WORKER_ATTACH (1) //!< Message sent from worker to attach to a coordinator
37#define FR_CONTROL_ID_COORD_WORKER_DETACH (2) //!< Message sent from worker to detach from a coordinator
38#define FR_CONTROL_ID_COORD_WORKER_ACK (3) //!< Message sent to worker to acknowledge attach / detach
39#define FR_CONTROL_ID_COORD_DATA (4) //!< Worker <-> coordinator message to pass data to a callback
40
41#define MIN_WORKER_ID -1 //!< The minimum value we expect as worker id. -1 is the main thread.
42
46
47/** A coordinator which receives messages from workers
48 */
49struct fr_coord_s {
50 fr_coord_reg_t *coord_reg; //!< Coordinator registration details.
51 fr_event_list_t *el; //!< Coordinator event list.
52 fr_rb_node_t node; //!< Entry in the tree of coordinators.
53 fr_coord_cb_reg_t *callbacks; //!< Array of callbacks for worker -> coordinator messages.
54 uint32_t num_callbacks; //!< Number of callbacks defined.
55 fr_coord_cb_inst_t **cb_inst; //!< Array of callback instance specific data.
56
57 uint32_t max_workers; //!< Maximum number of workers we expect.
58 uint32_t num_workers; //!< How many workers are attached.
59
60 fr_control_t *coord_recv_control; //!< Control plane for worker -> coordinator messages.
61 fr_atomic_queue_t *coord_recv_aq; //!< Atomic queue for worker -> coordinator
62 fr_ring_buffer_t **coord_send_rb; //!< Ring buffers for coordinator -> worker control messages.
63 fr_control_t **coord_send_control; //!< Control planes for coordinator -> worker messages.
64 fr_message_set_t **coord_send_ms; //!< Message sets for coordinator -> worker data.
65 fr_atomic_queue_t **coord_send_aq; //!< Atomic queues for coordinator -> worker data.
66
67 bool exiting; //!< Is this coordinator shutting down.
68 bool single_thread; //!< Are we in single thread mode.
69};
70
71/** The worker end of worker <-> coordinator communication.
72 */
74 fr_coord_t *coord; //!< Coordinator this worker is related to
75 fr_ring_buffer_t *worker_send_rb; //!< Ring buffer for worker -> coordinator control plane
76 fr_message_set_t *worker_send_ms; //!< Message set for worker -> coordinator messages
77 fr_control_t *worker_recv_control; //!< Coordinator -> worker control plane
78 fr_atomic_queue_t *worker_recv_aq; //!< Atomic queue for coordinator -> worker messages
79 fr_coord_worker_cb_reg_t *callbacks; //!< Callbacks for coordinator -> worker messages
80 uint32_t num_callbacks; //!< Number of callbacks registered.
81};
82
83/** A coordinator registration
84 */
86 char const *name; //!< Name for debugging.
87 fr_dlist_t entry; //!< Entry in list of registrations.
88 fr_coord_cb_reg_t *coord_cb; //!< Callbacks for worker -> coordinator messages.
89 fr_coord_worker_cb_reg_t *worker_cb; //!< Callbacks for coordinator -> worker messages.
90 size_t worker_send_size; //!< Initial size for worker -> coordinator ring buffer.
91 size_t coord_send_size; //!< Initial size for coordinator -> worker ring buffer.
92 module_instance_t const *mi; //!< Module instance which registered this coordinator.
93};
94
95/** Scheduler specific information for coordinator threads
96 */
97typedef struct {
98 fr_thread_t thread; //!< common thread information - must be first!
99
100 uint32_t max_workers; //!< Maximum number of workers which will connect to this coordinator.
101 fr_coord_reg_t *coord_reg; //!< Coordinator registration details.
102 fr_coord_t *coord; //!< The coordinator data structure.
103 fr_sem_t *sem; //!< For inter-thread signaling.
105
106/** Control plane message used for workers attaching / detaching to coordinators
107 */
108typedef struct {
109 int32_t worker; //!< Worker ID
110 fr_control_t *worker_recv_control; //!< Control plane to send messages to this worker
111 fr_atomic_queue_t *worker_recv_aq; //!< Atomic queue to send data to this worker
113
114typedef struct {
115 int32_t worker; //!< Worker ID
116 bool exiting; //!< Is the server exiting
118
119/** Compare coordinators by registration
120 */
121static fr_cmp_ret_t coord_cmp(void const *one, void const *two)
122{
123 fr_coord_t const *a = one, *b = two;
124
125 return CMP(a->coord_reg, b->coord_reg);
126}
127
128/** Register a coordinator
129 *
130 * To be called from mod_instantiate of a module which uses a coordinator
131 *
132 * @param reg_ctx Registration data
133 * @return
134 * - coordination registration on success
135 * - NULL on failure
136 */
138{
139 fr_coord_reg_t *coord_reg;
140
141 fr_assert(reg_ctx->coord_cb);
142 fr_assert(reg_ctx->worker_cb);
143 fr_assert(reg_ctx->mi);
144
145 /* Allocate the list of registered coordinators if not already done */
146 if (!coord_regs) {
147 MEM(coord_regs = talloc_zero(NULL, fr_dlist_head_t));
149 }
150
151 MEM(coord_reg = talloc(coord_regs, fr_coord_reg_t));
152 *coord_reg = (fr_coord_reg_t) {
153 .name = reg_ctx->name,
154 .coord_cb = reg_ctx->coord_cb,
155 .worker_cb = reg_ctx->worker_cb,
156 .worker_send_size = reg_ctx->worker_send_size ? reg_ctx->worker_send_size : 4096,
157 .coord_send_size = reg_ctx->coord_send_size ? reg_ctx->coord_send_size : 4096,
158 .mi = reg_ctx->mi,
159 };
160
162
163 return coord_reg;
164}
165
166/** De-register a coordinator
167 *
168 * To be called from mod_detach of a module which uses a coordinator
169 *
170 * @param coord_reg to de-register
171 */
173{
174 fr_dlist_remove(coord_regs, coord_reg);
175
176 talloc_free(coord_reg);
177
178 if (fr_dlist_num_elements(coord_regs) == 0) TALLOC_FREE(coord_regs);
179}
180
181/** Wait for all the coordinator threads to exit
182 *
183 * To be called during the scheduler shutdown in multi-threaded mode.
184 */
186{
187 int ret;
188
189 if (!coord_threads) return;
190
192 if ((ret = pthread_join(sc->thread.pthread_id, NULL)) != 0) {
193 ERROR("Failed joining coordinator %s: %s", sc->coord_reg->name, fr_syserror(ret));
194 } else {
195 DEBUG2("Coordinator %s joined (cleaned up)", sc->coord_reg->name);
196 }
197
200 }
201}
202
203/** Callback for a coordinator receiving data from a worker
204 */
205static void coord_data_recv(void const *data, size_t data_size, fr_time_t now, void *uctx)
206{
207 fr_coord_t *coord = talloc_get_type_abort(uctx, fr_coord_t);
209 fr_coord_data_t *cd;
210 fr_dbuff_t dbuff;
211
212 fr_assert(data_size == sizeof(cm));
213 memcpy(&cm, data, data_size);
214 fr_assert((cm.worker >= MIN_WORKER_ID) && (cm.worker < (int32_t)coord->max_workers));
215
216 if (unlikely(!fr_atomic_queue_pop(coord->coord_recv_aq, (void **)&cd))) return;
217
218 DEBUG3("Coordinator %s got data from worker %d for callback %d",
219 coord->coord_reg->name, cm.worker, cd->coord_cb_id);
220
221 if (cd->coord_cb_id >= coord->num_callbacks) {
222 ERROR("Received data for callback %d which is not defined", cd->coord_cb_id);
223 fr_message_done(&cd->m);
224 return;
225 }
226
227 fr_dbuff_init(&dbuff, (uint8_t const *)cd->m.data, cd->m.data_size);
228 coord->callbacks[cd->coord_cb_id].callback(coord, cm.worker, &dbuff, now,
229 MODULE_CTX(coord->coord_reg->mi,
230 module_thread(coord->coord_reg->mi)->data, NULL, NULL),
231 coord->cb_inst[cd->coord_cb_id] ?
232 coord->cb_inst[cd->coord_cb_id]->inst_data : NULL,
233 coord->callbacks[cd->coord_cb_id].uctx);
234 fr_message_done(&cd->m);
235}
236
237/** Callback for a worker receiving data from a coordinator
238 */
239static void coord_worker_data_recv(void const *data, size_t data_size, fr_time_t now, void *uctx)
240{
241 fr_coord_worker_t *cw = talloc_get_type_abort(uctx, fr_coord_worker_t);
243 fr_coord_data_t *cd;
244 fr_dbuff_t dbuff;
245
246 fr_assert(data_size == sizeof(cm));
247 memcpy(&cm, data, data_size);
248
249 if (unlikely(!fr_atomic_queue_pop(cw->worker_recv_aq, (void **)&cd))) return;
250
251 DEBUG3("Coordinator %s sent message for callback %d", cw->coord->coord_reg->name, cd->coord_cb_id);
252
253 if (cd->coord_cb_id >= cw->num_callbacks) {
254 ERROR("Received message for callback %d which is not defined", cd->coord_cb_id);
255 fr_message_done(&cd->m);
256 return;
257 }
258
259 fr_dbuff_init(&dbuff, (uint8_t const *)cd->m.data, cd->m.data_size);
260 cw->callbacks[cd->coord_cb_id].callback(cw, &dbuff, now,
262 module_thread(cw->coord->coord_reg->mi)->data, NULL, NULL),
263 cw->callbacks[cd->coord_cb_id].uctx);
264 fr_message_done(&cd->m);
265}
266
267/** Callback run by a coordinator when a worker attaches
268 */
269static void coord_worker_attach(void const *data, NDEBUG_UNUSED size_t data_size, UNUSED fr_time_t now, void *uctx)
270{
271 fr_coord_t *coord = talloc_get_type_abort(uctx, fr_coord_t);
273 fr_coord_msg_t ack;
274 uint32_t thread_id;
275
276 fr_assert(data_size == sizeof(fr_coord_worker_attach_msg_t));
277 fr_assert((msg->worker >= MIN_WORKER_ID) && (msg->worker < (int32_t)coord->max_workers));
278
279 DEBUG2("Worker %d attached to %s", msg->worker, coord->coord_reg->name);
280 coord->num_workers++;
281 thread_id = msg->worker - MIN_WORKER_ID;
282 coord->coord_send_control[thread_id] = msg->worker_recv_control;
283 coord->coord_send_aq[thread_id] = msg->worker_recv_aq;
284
285 ack.worker = msg->worker;
286 fr_control_message_send(coord->coord_send_control[thread_id], coord->coord_send_rb[thread_id],
287 FR_CONTROL_ID_COORD_WORKER_ACK, &ack, sizeof(ack));
288}
289
290/** Callback run by a coordinator when a worker detaches
291 */
292static void coord_worker_detach(void const *data, NDEBUG_UNUSED size_t data_size, UNUSED fr_time_t now, void *uctx)
293{
294 fr_coord_t *coord = talloc_get_type_abort(uctx, fr_coord_t);
296 fr_coord_msg_t ack;
297 uint32_t thread_id, i;
298
299 fr_assert(data_size == sizeof(fr_coord_worker_detach_msg_t));
300 fr_assert((msg->worker >= MIN_WORKER_ID) && (msg->worker < (int32_t)coord->max_workers));
301 thread_id = msg->worker - MIN_WORKER_ID;
302
303 DEBUG2("Worker %d detached from %s", msg->worker, coord->coord_reg->name);
304 coord->num_workers--;
305 if (msg->exiting) coord->exiting = true;
306
307 /*
308 * If all workers have detached, and we're exiting, run any exit callbacks.
309 */
310 if (coord->exiting && (coord->num_workers == 0)) {
311 fr_coord_cb_inst_t *cb_inst;
312 for (i = 0; i < coord->num_callbacks; i++) {
313 cb_inst = coord->cb_inst[i];
314 if (!cb_inst || !cb_inst->exit_cb) continue;
315 cb_inst->exit_cb(coord, coord->el, cb_inst->inst_data);
316 }
317 }
318
319 ack.worker = msg->worker;
320 fr_control_message_send(coord->coord_send_control[thread_id], coord->coord_send_rb[thread_id],
322
323 coord->coord_send_control[thread_id] = NULL;
324 coord->coord_send_aq[thread_id] = NULL;
325}
326
327/** Create a coordinator from its registration
328 *
329 * @param ctx to allocate the coordinator in
330 * @param el Event list to run this coordinator
331 * @param coord_reg Registration to configure this coordinator
332 * @param single_thread Is the server in single thread mode
333 * @param max_workers The maximum number of workers which will attach
334 * @return
335 * - the coordinator on success
336 * - NULL on failure
337 */
338static fr_coord_t *fr_coord_create(TALLOC_CTX *ctx, fr_event_list_t *el, fr_coord_reg_t *coord_reg,
339 bool single_thread, uint32_t max_workers)
340{
341 fr_coord_t *coord;
342 uint32_t i, num_threads = max_workers - MIN_WORKER_ID;
343 fr_coord_cb_reg_t *cb = coord_reg->coord_cb;
345
346 MEM(coord = talloc(ctx, fr_coord_t));
347 *coord = (fr_coord_t) {
348 .el = el,
349 .coord_reg = coord_reg,
350 .single_thread = single_thread,
351 .max_workers = max_workers
352 };
353
354 /* Allocate atomic queue / control for receiving messages from workers */
356 if (!aq) {
357 fr_strerror_const("Failed creating worker -> coordinator atomic queue");
358 fail:
359 talloc_free(coord);
360 return NULL;
361 }
362 coord->coord_recv_control = fr_control_create(coord, el, aq, 5);
363 if (!coord->coord_recv_control) {
364 fr_strerror_const("Failed creating worker -> coordinator control plane");
365 goto fail;
366 }
367
368 /* Allocate atomic queue for workers sending data to coordinators */
370 if (!coord->coord_recv_aq) {
371 fr_strerror_const("Failed creating worker -> coordinator data atomic queue");
372 goto fail;
373 }
374
376 coord_worker_attach, coord) < 0) goto fail;
378 coord_worker_detach, coord) < 0) goto fail;
380 coord_data_recv, coord) < 0) goto fail;
381
382 /* Count the number of callbacks defined, for sanity checking messages */
383 while (cb->callback) {
384 coord->num_callbacks++;
385 cb++;
386 }
387 coord->callbacks = coord_reg->coord_cb;
388
389 if (fr_control_open(coord->coord_recv_control) < 0) {
390 fr_strerror_const("Failed opening control plane");
391 goto fail;
392 }
393
394 /*
395 * Coordinator side arrays for holding pointers to worker
396 * specific communication structures. The array sizes are the
397 * number of threads expected to attach which is the number of
398 * workers plus any additional threads, currently just the main
399 * thread (worker id -1)
400 */
401 MEM(coord->coord_send_rb = talloc_array(coord, fr_ring_buffer_t *, num_threads));
402 MEM(coord->coord_send_ms = talloc_array(coord, fr_message_set_t *, num_threads));
403 for (i = 0; i < num_threads; i++) {
405 if (!coord->coord_send_rb[i]) goto fail;
406
408 coord_reg->coord_send_size, true);
409 if (!coord->coord_send_ms[i]) goto fail;
410 }
411 MEM(coord->coord_send_control = talloc_zero_array(coord, fr_control_t *, num_threads));
412 MEM(coord->coord_send_aq = talloc_zero_array(coord, fr_atomic_queue_t *, num_threads));
413
414 MEM(coord->cb_inst = talloc_zero_array(coord, fr_coord_cb_inst_t *, coord->num_callbacks));
415
416 for (i = 0; i < coord->num_callbacks; i++) {
417 if (!coord->callbacks[i].inst_create) continue;
418 coord->cb_inst[i] = coord->callbacks[i].inst_create(coord, coord, coord->el, coord->single_thread,
419 coord->callbacks[i].uctx);
420 if (!coord->cb_inst[i]) goto fail;
421 }
422
423 return coord;
424}
425
426static void fr_coord_destroy(fr_coord_t *coord){
427 uint32_t i;
428
429 for (i = 0; i < coord->num_callbacks; i++) {
430 if (!coord->callbacks[i].inst_destroy) continue;
431 coord->callbacks[i].inst_destroy(coord, coord->cb_inst[i], coord->single_thread,
432 coord->callbacks[i].uctx);
433 }
434}
435
436/** Run the event loop for a coordinator thread when in multi-threaded mode
437 */
438static void fr_coordinate(fr_coord_t *coord)
439{
440 uint32_t i;
441 fr_coord_cb_inst_t *cb_inst;
442
443 /*
444 * Run until we're told to exit AND the number of
445 * workers has dropped to zero.
446 *
447 * Whenever a worker detaches, coord->num_workers
448 * is decremented, so when coord->num_workers == 0,
449 * all workers have detached and are no longer using
450 * the channel.
451 */
452 while (true) {
453 bool wait_for_events = true;
454 int num_events;
455 fr_time_t now = fr_time();
456
457 /*
458 * Check if any coordinator instances report that they have
459 * events to process.
460 */
461 for (i = 0; i < coord->num_callbacks; i++) {
462 cb_inst = coord->cb_inst[i];
463 if (!cb_inst || !cb_inst->event_pre_cb) continue;
464 wait_for_events = (cb_inst->event_pre_cb(now, fr_time_delta_wrap(0), cb_inst->inst_data) == 0);
465 if (!wait_for_events) break;
466 }
467
468 if (wait_for_events) {
469 if (unlikely(coord->exiting) && (coord->num_workers == 0)) break;
470 }
471
472 /*
473 * Check the event list. If there's an error
474 * (e.g. exit), we stop looping and clean up.
475 */
476 DEBUG4("Gathering events");
477 num_events = fr_event_corral(coord->el, fr_time(), true);
478 DEBUG4("%u event(s) pending%s",
479 num_events == -1 ? 0 : num_events, num_events == -1 ? " - event loop exiting" : "");
480 if (num_events < 0) break;
481
482 /*
483 * Service outstanding events.
484 */
485 if (num_events > 0) {
486 DEBUG4("Servicing event(s)");
487 fr_event_service(coord->el);
488 }
489
490 /*
491 * Run any registered instance specific event callbacks
492 */
493 for (i = 0; i < coord->num_callbacks; i++) {
494 cb_inst = coord->cb_inst[i];
495 if (cb_inst && cb_inst->event_post_cb) cb_inst->event_post_cb(coord->el, now, cb_inst->inst_data);
496 }
497 }
498
499 fr_coord_destroy(coord);
500
501 return;
502}
503
504/** Entry point for a coordinator thread
505 */
506static void *fr_coordinate_thread(void *arg)
507{
508 fr_schedule_coord_t *sc = talloc_get_type_abort(arg, fr_schedule_coord_t);
509 fr_coord_reg_t *coord_reg = sc->coord_reg;
511 char coordinate_name[64];
512
513 snprintf(coordinate_name, sizeof(coordinate_name), "Coordinate %s", coord_reg->name);
514
515 if (fr_thread_setup(&sc->thread, coordinate_name) < 0) goto fail;
516
517 sc->coord = fr_coord_create(sc->thread.ctx, sc->thread.el, coord_reg, false, sc->max_workers);
518 if (!sc->coord) {
519 PERROR("%s - Failed creating coordinator thread", coordinate_name);
520 goto fail;
521 }
522
523 /*
524 * Create all the thread specific data for the coordinator thread
525 */
526 if (fr_thread_instantiate(sc->thread.ctx, sc->thread.el) < 0) goto fail;
527
528 /*
529 * Tell the originator that the thread has started.
530 */
531 fr_thread_start(&sc->thread, sc->sem);
532
533 fr_coordinate(sc->coord);
534
535 status = FR_THREAD_EXITED;
536
537fail:
539
540 fr_thread_exit(&sc->thread, status, sc->sem);
541
542 return NULL;
543}
544
545/** Start all registered coordinator threads in multi-threaded mode
546 *
547 * @param num_workers The number of workers which will be attaching
548 * @param sem Semaphore to use signalling the threads are ready
549 * @return
550 * - 0 on success
551 * - -1 on failure
552 */
554{
555 int num = 0;
556
557 if (!coord_regs) return 0;
558
559 MEM(coord_threads = talloc(NULL, fr_dlist_head_t));
562
565
566 MEM(sc = talloc_zero(coord_threads, fr_schedule_coord_t));
567
568 sc->thread.id = num++;
569 sc->coord_reg = coord_reg;
570 sc->max_workers = num_workers;
571 sc->sem = sem;
572
573 if (fr_thread_create(&sc->thread.pthread_id, fr_coordinate_thread, sc) < 0) {
575 PERROR("Failed creating coordinator %s", coord_reg->name);
576 return -1;
577 };
578
580 }
581
582 /*
583 * Wait for all the coordinators to start.
584 */
585 if (fr_thread_wait_list(sem, coord_threads) < 0) {
586 ERROR("Failed creating coordinator threads");
587 return -1;
588 }
589
590 /*
591 * Insert the coordinators in the tree
592 */
594 fr_assert(sc->coord);
595 fr_rb_insert(&coords, sc->coord);
596 }
597
598 return 0;
599}
600
601/** Clean up coordinators in single threaded mode
602 */
604{
605 fr_coord_t *coord;
607
608 if (fr_rb_num_elements(&coords) == 0) return;
609
610 for (coord = fr_rb_iter_init_inorder(&coords, &iter);
611 coord;
612 coord = fr_rb_iter_next_inorder(&coords, &iter)) {
614 fr_coord_destroy(coord);
615 talloc_free(coord);
616 }
617}
618
619/** Start coordinators in single threaded mode
620 */
621int fr_coords_create(TALLOC_CTX *ctx, fr_event_list_t *el)
622{
623 if (!coord_regs) return 0;
624
626
628 char coordinate_name[64];
629 fr_coord_t *coord;
630
631 snprintf(coordinate_name, sizeof(coordinate_name), "Coordinator %s", coord_reg->name);
632
633 INFO("%s - Starting", coordinate_name);
634
635 coord = fr_coord_create(ctx, el, coord_reg, true, 1);
636 if (!coord) {
637 PERROR("%s - Failed creating coordinator thread", coordinate_name);
638 return -1;
639 }
640
641 fr_rb_insert(&coords, coord);
642 }
643
644 return 0;
645}
646
647/** Signal a coordinator that a worker wants to detach
648 *
649 * @param cw Worker which is detaching.
650 * @param exiting Is the server exiting.
651 */
652int fr_coord_detach(fr_coord_worker_t *cw, bool exiting)
653{
655
656 msg = talloc(cw, fr_coord_worker_detach_msg_t);
657 msg->worker = fr_schedule_worker_id();
658 msg->exiting = exiting;
659
662 msg, sizeof(fr_coord_worker_detach_msg_t)) < 0) return -1;
663
665
666 return 0;
667}
668
669/** A worker got an ack from a coordinator in response to attach / detach
670 */
671static void coordinate_worker_ack(NDEBUG_UNUSED void const *data, NDEBUG_UNUSED size_t data_size,
672 UNUSED fr_time_t now, UNUSED void *uctx)
673{
674#ifndef NDEBUG
675 fr_coord_msg_t const *cm = data;
676
677 fr_assert(data_size == sizeof(fr_coord_msg_t));
679#endif
680}
681
682/** Attach a worker to a coordinator
683 *
684 * @param ctx To allocate worker structure in
685 * @param el Event list for control messages
686 * @param coord_reg Coordinator registration to attach to.
687 * @return
688 * - Worker structure for coordinator use on success
689 * - NULL on failure
690 */
692{
694 fr_coord_worker_cb_reg_t *cb_reg = coord_reg->worker_cb;
696 fr_coord_t find;
698
699 cw = talloc_zero(ctx, fr_coord_worker_t);
700
701 find = (fr_coord_t) {
702 .coord_reg = coord_reg
703 };
704 fr_rb_find((void **)&cw->coord, &coords, &find);
705 if (!cw->coord) {
706 fail:
707 talloc_free(cw);
708 return NULL;
709 }
710
711 aq = fr_atomic_queue_talloc(cw, 1024);
716 coord_reg->worker_send_size, true);
717
718 while (cb_reg->callback) {
719 cw->num_callbacks++;
720 cb_reg++;
721 }
722 cw->callbacks = coord_reg->worker_cb;
723
725 coordinate_worker_ack, cw) < 0) goto fail;
727 coord_worker_data_recv, cw) < 0) goto fail;
728
729 if (fr_control_open(cw->worker_recv_control) < 0) goto fail;
730
731 msg.worker_recv_control = cw->worker_recv_control;
732 msg.worker_recv_aq = cw->worker_recv_aq;
733 msg.worker = fr_schedule_worker_id();
734
737 &msg, sizeof(fr_coord_worker_attach_msg_t)) < 0) goto fail;
738
740
741 return cw;
742}
743
744/** Send generic data from a coordinator to a worker
745 *
746 * @param coord Coordinator which is sending the data.
747 * @param worker_id Worker to send data to.
748 * @param cb_id Callback ID for the worker to run.
749 * @param dbuff Buffer containing data to send.
750 * @return
751 * - 0 on success
752 * - <0 on failure
753 */
755{
757 fr_coord_data_t *cd = NULL;
758 uint32_t thread_id = worker_id - MIN_WORKER_ID;
759
760 fr_assert((worker_id >= MIN_WORKER_ID) && (worker_id < (int32_t)coord->max_workers));
761
762 cm = (fr_coord_msg_t) {
764 };
765
766 cd = (fr_coord_data_t *)fr_message_and_data_alloc(coord->coord_send_ms[thread_id], fr_dbuff_used(dbuff));
767 if (!cd) return -1;
768
769 memcpy(cd->m.data, fr_dbuff_buff(dbuff), fr_dbuff_used(dbuff));
770 cd->coord_cb_id = cb_id;
771 if (!fr_atomic_queue_push(coord->coord_send_aq[thread_id], cd)) {
773 return -1;
774 }
775 return fr_control_message_send(coord->coord_send_control[thread_id], coord->coord_send_rb[thread_id],
777 &cm, sizeof(fr_coord_msg_t));
778}
779
780/** Broadcast data from a coordinator to all workers
781 *
782 * @param coord Coordinator which is sending the data.
783 * @param cb_id Callback ID for the workers to run.
784 * @param dbuff Buffer containing data to send.
785 * @return
786 * - 0 on success
787 * - <0 on failure - indicating the number of sends which failed.
788 */
790{
791 int32_t i;
792 int failed = 0;
793
794 for (i = 0; i < (int32_t)coord->max_workers; i++) {
795 if (!coord->coord_send_control[i - MIN_WORKER_ID]) continue;
796 if (fr_coord_to_worker_send(coord, i, cb_id, dbuff) < 0) failed++;
797 }
798
799 return 0 - failed;
800}
801
802/** Send data from a worker to a coordinator
803 *
804 * @param cw Worker side of coordinator sending the data.
805 * @param cb_id Callback ID for the coordinator to run.
806 * @param dbuff Buffer containing data to send.
807 * @return
808 * - 0 on success
809 * - < 0 on failure
810 */
812{
814 fr_coord_data_t *cd = NULL;
815
816 cm = (fr_coord_msg_t) {
818 };
819
821 if (!cd) return -1;
822
823 memcpy(cd->m.data, fr_dbuff_buff(dbuff), fr_dbuff_used(dbuff));
824 cd->coord_cb_id = cb_id;
827 return -1;
828 }
829
832}
833
834/** Insert instance specific pre-event callbacks
835 */
837{
838 fr_coord_t *coord;
840 fr_coord_cb_inst_t *cb_inst;
841 uint32_t i;
842
843 if (!coord_regs) return 0;
844
845 for (coord = fr_rb_iter_init_inorder(&coords, &iter);
846 coord != NULL;
847 coord = fr_rb_iter_next_inorder(&coords, &iter)) {
848 for (i = 0; i < coord->num_callbacks; i++) {
849 cb_inst = coord->cb_inst[i];
850 if (cb_inst && cb_inst->event_pre_cb &&
851 fr_event_pre_insert(el, cb_inst->event_pre_cb, cb_inst->inst_data) < 0) {
852 return -1;
853 }
854 }
855 }
856 return 0;
857}
858
859/** Insert instance specific post-event callbacks
860 */
862{
863 fr_coord_t *coord;
865 fr_coord_cb_inst_t *cb_inst;
866 uint32_t i;
867
868 if (!coord_regs) return 0;
869
870 for (coord = fr_rb_iter_init_inorder(&coords, &iter);
871 coord != NULL;
872 coord = fr_rb_iter_next_inorder(&coords, &iter)) {
873 for (i = 0; i < coord->num_callbacks; i++) {
874 cb_inst = coord->cb_inst[i];
875 if (cb_inst && cb_inst->event_post_cb &&
876 fr_event_post_insert(el, cb_inst->event_post_cb, cb_inst->inst_data) < 0) {
877 return -1;
878 }
879 }
880 }
881 return 0;
882}
883
884/** Event loop callback to exit the loop when all workers have detached from all coordinators
885 *
886 * Used during single threaded shut down to allow the event loop to run any
887 * tidy up needed by coordinators.
888 */
890{
891 fr_coord_t *coord;
893
894 if (fr_rb_num_elements(&coords) == 0) return;
895
896 for (coord = fr_rb_iter_init_inorder(&coords, &iter);
897 coord;
898 coord = fr_rb_iter_next_inorder(&coords, &iter)) {
899 if (coord->num_workers > 0) return;
900 }
901
902 fr_event_loop_exit(el, SIGINT);
903}
904
909
910/** Return the coordinator name
911 */
912char const *fr_coord_name(fr_coord_t const *coord)
913{
914 return coord->coord_reg->name;
915}
log_entry msg
Definition acutest.h:794
bool fr_atomic_queue_pop(fr_atomic_queue_t *aq, void **p_data)
Pop a pointer from the atomic queue.
bool fr_atomic_queue_push(fr_atomic_queue_t *aq, void *data)
Push a pointer into the atomic queue.
fr_atomic_queue_t * fr_atomic_queue_talloc(TALLOC_CTX *ctx, size_t size)
Create fixed-size atomic queue.
Structure to hold the atomic queue.
#define RCSID(id)
Definition build.h:560
#define NDEBUG_UNUSED
Definition build.h:395
#define CMP(_a, _b)
Same as CMP_PREFER_SMALLER use when you don't really care about ordering, you just want an ordering.
Definition build.h:113
#define unlikely(_x)
Definition build.h:455
#define UNUSED
Definition build.h:384
#define FR_CONTROL_MAX_SIZE
Definition control.h:51
#define FR_CONTROL_MAX_MESSAGES
Definition control.h:50
static size_t num_workers
static fr_atomic_queue_t ** aq
fr_dlist_t entry
Entry in list of registrations.
Definition coord.c:87
fr_atomic_queue_t * worker_recv_aq
Atomic queue for coordinator -> worker messages.
Definition coord.c:78
uint32_t max_workers
Maximum number of workers which will connect to this coordinator.
Definition coord.c:100
fr_control_t ** coord_send_control
Control planes for coordinator -> worker messages.
Definition coord.c:63
#define FR_CONTROL_ID_COORD_WORKER_DETACH
Message sent from worker to detach from a coordinator.
Definition coord.c:37
fr_coord_reg_t * coord_reg
Coordinator registration details.
Definition coord.c:101
fr_coord_worker_cb_reg_t * worker_cb
Callbacks for coordinator -> worker messages.
Definition coord.c:89
fr_coord_reg_t * fr_coord_register(fr_coord_reg_ctx_t *reg_ctx)
Register a coordinator.
Definition coord.c:137
int fr_coords_create(TALLOC_CTX *ctx, fr_event_list_t *el)
Start coordinators in single threaded mode.
Definition coord.c:621
fr_control_t * worker_recv_control
Control plane to send messages to this worker.
Definition coord.c:110
static fr_dlist_head_t * coord_threads
Definition coord.c:44
static fr_coord_t * fr_coord_create(TALLOC_CTX *ctx, fr_event_list_t *el, fr_coord_reg_t *coord_reg, bool single_thread, uint32_t max_workers)
Create a coordinator from its registration.
Definition coord.c:338
fr_coord_worker_t * fr_coord_attach(TALLOC_CTX *ctx, fr_event_list_t *el, fr_coord_reg_t *coord_reg)
Attach a worker to a coordinator.
Definition coord.c:691
fr_coord_reg_t * coord_reg
Coordinator registration details.
Definition coord.c:50
size_t worker_send_size
Initial size for worker -> coordinator ring buffer.
Definition coord.c:90
size_t coord_send_size
Initial size for coordinator -> worker ring buffer.
Definition coord.c:91
int32_t worker
Worker ID.
Definition coord.c:109
#define MIN_WORKER_ID
The minimum value we expect as worker id. -1 is the main thread.
Definition coord.c:41
int fr_coord_to_worker_send(fr_coord_t *coord, int32_t worker_id, uint32_t cb_id, fr_dbuff_t *dbuff)
Send generic data from a coordinator to a worker.
Definition coord.c:754
fr_ring_buffer_t ** coord_send_rb
Ring buffers for coordinator -> worker control messages.
Definition coord.c:62
char const * fr_coord_name(fr_coord_t const *coord)
Return the coordinator name.
Definition coord.c:912
uint32_t max_workers
Maximum number of workers we expect.
Definition coord.c:57
bool exiting
Is the server exiting.
Definition coord.c:116
fr_atomic_queue_t * worker_recv_aq
Atomic queue to send data to this worker.
Definition coord.c:111
fr_coord_t * coord
Coordinator this worker is related to.
Definition coord.c:74
fr_coord_cb_inst_t ** cb_inst
Array of callback instance specific data.
Definition coord.c:55
void fr_coord_deregister(fr_coord_reg_t *coord_reg)
De-register a coordinator.
Definition coord.c:172
void fr_coords_destroy(void)
Clean up coordinators in single threaded mode.
Definition coord.c:603
fr_control_t * coord_recv_control
Control plane for worker -> coordinator messages.
Definition coord.c:60
#define FR_CONTROL_ID_COORD_WORKER_ATTACH
Message sent from worker to attach to a coordinator.
Definition coord.c:36
int fr_coord_close_event_insert(fr_event_list_t *el)
Definition coord.c:905
static void coordinate_worker_ack(NDEBUG_UNUSED void const *data, NDEBUG_UNUSED size_t data_size, UNUSED fr_time_t now, UNUSED void *uctx)
A worker got an ack from a coordinator in response to attach / detach.
Definition coord.c:671
int fr_worker_to_coord_send(fr_coord_worker_t *cw, uint32_t cb_id, fr_dbuff_t *dbuff)
Send data from a worker to a coordinator.
Definition coord.c:811
void fr_coord_thread_join(void)
Wait for all the coordinator threads to exit.
Definition coord.c:185
static fr_dlist_head_t * coord_regs
Definition coord.c:43
module_instance_t const * mi
Module instance which registered this coordinator.
Definition coord.c:92
fr_coord_cb_reg_t * callbacks
Array of callbacks for worker -> coordinator messages.
Definition coord.c:53
int fr_coord_post_event_insert(fr_event_list_t *el)
Insert instance specific post-event callbacks.
Definition coord.c:861
int fr_coord_pre_event_insert(fr_event_list_t *el)
Insert instance specific pre-event callbacks.
Definition coord.c:836
uint32_t num_callbacks
Number of callbacks registered.
Definition coord.c:80
fr_coord_cb_reg_t * coord_cb
Callbacks for worker -> coordinator messages.
Definition coord.c:88
fr_atomic_queue_t * coord_recv_aq
Atomic queue for worker -> coordinator.
Definition coord.c:61
static void * fr_coordinate_thread(void *arg)
Entry point for a coordinator thread.
Definition coord.c:506
fr_atomic_queue_t ** coord_send_aq
Atomic queues for coordinator -> worker data.
Definition coord.c:65
uint32_t num_callbacks
Number of callbacks defined.
Definition coord.c:54
static void coord_worker_attach(void const *data, NDEBUG_UNUSED size_t data_size, UNUSED fr_time_t now, void *uctx)
Callback run by a coordinator when a worker attaches.
Definition coord.c:269
bool exiting
Is this coordinator shutting down.
Definition coord.c:67
static void coord_data_recv(void const *data, size_t data_size, fr_time_t now, void *uctx)
Callback for a coordinator receiving data from a worker.
Definition coord.c:205
fr_rb_node_t node
Entry in the tree of coordinators.
Definition coord.c:52
fr_coord_worker_cb_reg_t * callbacks
Callbacks for coordinator -> worker messages.
Definition coord.c:79
fr_event_list_t * el
Coordinator event list.
Definition coord.c:51
int fr_coord_start(uint32_t num_workers, fr_sem_t *sem)
Start all registered coordinator threads in multi-threaded mode.
Definition coord.c:553
#define FR_CONTROL_ID_COORD_WORKER_ACK
Message sent to worker to acknowledge attach / detach.
Definition coord.c:38
#define FR_CONTROL_ID_COORD_DATA
Worker <-> coordinator message to pass data to a callback.
Definition coord.c:39
fr_message_set_t ** coord_send_ms
Message sets for coordinator -> worker data.
Definition coord.c:64
static void coord_worker_detach(void const *data, NDEBUG_UNUSED size_t data_size, UNUSED fr_time_t now, void *uctx)
Callback run by a coordinator when a worker detaches.
Definition coord.c:292
fr_thread_t thread
common thread information - must be first!
Definition coord.c:98
static void fr_coord_close_post_event(fr_event_list_t *el, UNUSED fr_time_t now, UNUSED void *uctx)
Event loop callback to exit the loop when all workers have detached from all coordinators.
Definition coord.c:889
static void fr_coordinate(fr_coord_t *coord)
Run the event loop for a coordinator thread when in multi-threaded mode.
Definition coord.c:438
fr_control_t * worker_recv_control
Coordinator -> worker control plane.
Definition coord.c:77
bool single_thread
Are we in single thread mode.
Definition coord.c:68
char const * name
Name for debugging.
Definition coord.c:86
static void coord_worker_data_recv(void const *data, size_t data_size, fr_time_t now, void *uctx)
Callback for a worker receiving data from a coordinator.
Definition coord.c:239
int fr_coord_to_worker_broadcast(fr_coord_t *coord, uint32_t cb_id, fr_dbuff_t *dbuff)
Broadcast data from a coordinator to all workers.
Definition coord.c:789
int fr_coord_detach(fr_coord_worker_t *cw, bool exiting)
Signal a coordinator that a worker wants to detach.
Definition coord.c:652
int32_t worker
Worker ID.
Definition coord.c:115
static fr_cmp_ret_t coord_cmp(void const *one, void const *two)
Compare coordinators by registration.
Definition coord.c:121
fr_ring_buffer_t * worker_send_rb
Ring buffer for worker -> coordinator control plane.
Definition coord.c:75
fr_coord_t * coord
The coordinator data structure.
Definition coord.c:102
static fr_rb_tree_t coords
Definition coord.c:45
uint32_t num_workers
How many workers are attached.
Definition coord.c:58
fr_message_set_t * worker_send_ms
Message set for worker -> coordinator messages.
Definition coord.c:76
static void fr_coord_destroy(fr_coord_t *coord)
Definition coord.c:426
fr_sem_t * sem
For inter-thread signaling.
Definition coord.c:103
A coordinator registration.
Definition coord.c:85
A coordinator which receives messages from workers.
Definition coord.c:49
Control plane message used for workers attaching / detaching to coordinators.
Definition coord.c:108
The worker end of worker <-> coordinator communication.
Definition coord.c:73
Scheduler specific information for coordinator threads.
Definition coord.c:97
fr_coord_cb_t callback
Callback to run when message is received.
Definition coord.h:48
fr_coord_cb_reg_t * coord_cb
Callbacks for worker -> coordinator messages.
Definition coord.h:64
struct fr_coord_s fr_coord_t
Definition coord.h:37
void * uctx
To pass to all callbacks.
Definition coord.h:51
char const * name
Name for this coordinator.
Definition coord.h:63
module_instance_t const * mi
Module instance registering this coordinator.
Definition coord.h:70
fr_coord_worker_cb_t callback
Definition coord.h:56
size_t coord_send_size
Initial ring buffer size for coordinator -> worker data.
Definition coord.h:68
struct fr_coord_reg_s fr_coord_reg_t
Definition coord.h:35
fr_coord_worker_cb_reg_t * worker_cb
Callbacks for coordinator -> worker messages.
Definition coord.h:65
size_t worker_send_size
Initial ring buffer size for worker -> coordinator data.
Definition coord.h:66
fr_coord_cb_inst_create_t inst_create
Callback to create coordinator instance.
Definition coord.h:49
fr_coord_cb_inst_destroy_t inst_destroy
Callback to destroyed coordinator instance.
Definition coord.h:50
fr_coord_inst_exit_cb_t exit_cb
Callback run when coordinator is asked to exit.
Definition coord_priv.h:49
void * inst_data
Instance data.
Definition coord_priv.h:46
int32_t worker
Worker ID.
Definition coord_priv.h:33
uint32_t coord_cb_id
Callback ID for this message.
Definition coord_priv.h:40
fr_event_status_cb_t event_pre_cb
Pre-event callback.
Definition coord_priv.h:47
fr_event_post_cb_t event_post_cb
Post-event callback.
Definition coord_priv.h:48
fr_message_t m
Message containing data being sent.
Definition coord_priv.h:39
List / data message used between workers and coordinators.
Definition coord_priv.h:38
Generic control message used between workers and coordinators.
Definition coord_priv.h:32
#define fr_dbuff_used(_dbuff_or_marker)
Return the number of bytes remaining between the start of the dbuff or marker and the current positio...
Definition dbuff.h:775
#define fr_dbuff_init(_out, _start, _len_or_end)
Initialise an dbuff for encoding or decoding.
Definition dbuff.h:362
#define fr_dbuff_buff(_dbuff_or_marker)
Return the underlying buffer in a dbuff or one of marker.
Definition dbuff.h:890
#define MEM(x)
Definition debug.h:38
#define ERROR(fmt,...)
Definition dhcpclient.c:40
#define fr_dlist_init(_head, _type, _field)
Initialise the head structure of a doubly linked list.
Definition dlist.h:242
#define fr_dlist_foreach(_list_head, _type, _iter)
Iterate over the contents of a list.
Definition dlist.h:98
static void * fr_dlist_remove(fr_dlist_head_t *list_head, void *ptr)
Remove an item from the list.
Definition dlist.h:620
static unsigned int fr_dlist_num_elements(fr_dlist_head_t const *head)
Return the number of elements in the dlist.
Definition dlist.h:921
static int fr_dlist_insert_tail(fr_dlist_head_t *list_head, void *ptr)
Insert an item into the tail of a list.
Definition dlist.h:360
#define fr_dlist_talloc_init(_head, _type, _field)
Initialise the head structure of a doubly linked list.
Definition dlist.h:257
Head of a doubly linked list.
Definition dlist.h:51
Entry in a doubly linked list.
Definition dlist.h:41
talloc_free(hp)
fr_control_t * fr_control_create(TALLOC_CTX *ctx, fr_event_list_t *el, fr_atomic_queue_t *aq, size_t num_callbacks)
Create a control-plane signaling path.
Definition control.c:152
int fr_control_open(fr_control_t *c)
Open the control-plane signalling path.
Definition control.c:176
int fr_control_message_send(fr_control_t *c, fr_ring_buffer_t *rb, uint32_t id, void *data, size_t data_size)
Send a control-plane message.
Definition control.c:355
int fr_control_callback_add(fr_control_t **c, uint32_t id, fr_control_callback_t callback, void *uctx)
Register a callback for an ID.
Definition control.c:444
void fr_control_wait(fr_control_t *c)
Wait for a plane control to become readable.
Definition control.c:543
The control structure.
Definition control.c:76
#define PERROR(_fmt,...)
Definition log.h:233
#define DEBUG3(_fmt,...)
Definition log.h:271
#define DEBUG4(_fmt,...)
Definition log.h:272
void fr_event_service(fr_event_list_t *el)
Service any outstanding timer or file descriptor events.
Definition event.c:2205
int fr_event_corral(fr_event_list_t *el, fr_time_t now, bool wait)
Gather outstanding timer and file descriptor events.
Definition event.c:2073
#define fr_time()
Definition event.c:60
int fr_event_pre_insert(fr_event_list_t *el, fr_event_status_cb_t callback, void *uctx)
Add a pre-event callback to the event list.
Definition event.c:1981
void fr_event_loop_exit(fr_event_list_t *el, int code)
Signal an event loop exit with the specified code.
Definition event.c:2383
int fr_event_post_insert(fr_event_list_t *el, fr_event_post_cb_t callback, void *uctx)
Add a post-event callback to the event list.
Definition event.c:2028
Stores all information relating to an event list.
Definition event.c:377
unsigned int uint32_t
unsigned char uint8_t
fr_message_set_t * fr_message_set_create(TALLOC_CTX *ctx, int num_messages, size_t message_size, size_t ring_buffer_size, bool unlimited_size)
Create a message set.
Definition message.c:127
int fr_message_done(fr_message_t *m)
Mark a message as done.
Definition message.c:195
fr_message_t * fr_message_and_data_commit(fr_message_set_t *ms, fr_message_t *m, size_t total_size)
Commit a previously reserved message, allocating exactly total_size bytes of packet data.
Definition message.c:1051
fr_message_t * fr_message_and_data_alloc(fr_message_set_t *ms, size_t size)
Reserve and commit a message atomically.
Definition message.c:1108
A Message set, composed of message headers and ring buffer data.
Definition message.c:94
uint8_t * data
pointer to the data in the ring buffer
Definition message.h:49
size_t data_size
size of the data in the ring buffer
Definition message.h:50
fr_cmp_ret_t
Result of an ordering comparison.
Definition misc.h:50
#define MODULE_CTX(_mi, _thread, _env_data, _rctx)
Wrapper to create a module_ctx_t as a compound literal.
Definition module_ctx.h:128
#define fr_assert(_expr)
Definition rad_assert.h:37
#define DEBUG2(fmt,...)
#define INFO(fmt,...)
Definition radict.c:63
uint32_t fr_rb_num_elements(fr_rb_tree_t *tree)
Return how many nodes there are in a tree.
Definition rb.c:807
int fr_rb_find(void **found, fr_rb_tree_t const *tree, void const *data)
Find an element in the tree, returning the data, not the node.
Definition rb.c:586
void * fr_rb_iter_init_inorder(fr_rb_tree_t *tree, fr_rb_iter_inorder_t *iter)
Initialise an in-order iterator.
Definition rb.c:850
void fr_rb_iter_delete_inorder(fr_rb_tree_t *tree, fr_rb_iter_inorder_t *iter)
Remove the current node from the tree.
Definition rb.c:925
int fr_rb_insert(fr_rb_tree_t *tree, void const *data)
Insert data into a tree.
Definition rb.c:637
void * fr_rb_iter_next_inorder(UNUSED fr_rb_tree_t *tree, fr_rb_iter_inorder_t *iter)
Return the next node.
Definition rb.c:876
struct fr_rb_tree_s fr_rb_tree_t
Definition rb.h:51
#define fr_rb_inline_talloc_init(_tree, _type, _field, _data_cmp, _data_free)
Initialises a red black that verifies elements are of a specific talloc type.
Definition rb.h:155
uint32_t num_elements
How many elements are inside the tree.
Definition rb.h:95
Iterator structure for in-order traversal of an rbtree.
Definition rb.h:319
The main red black tree structure.
Definition rb.h:71
fr_ring_buffer_t * fr_ring_buffer_create(TALLOC_CTX *ctx, size_t size)
Create a ring buffer.
Definition ring_buffer.c:64
static _Thread_local int worker_id
Internal ID of the current worker thread.
Definition schedule.c:105
int fr_schedule_worker_id(void)
Return the worker id for the current thread.
Definition schedule.c:111
sem_t fr_sem_t
Definition semaphore.h:53
void * data
Thread specific instance data.
Definition module.h:376
static module_thread_instance_t * module_thread(module_instance_t const *mi)
Retrieve module/thread specific instance for a module.
Definition module.h:515
Module instance data.
Definition module.h:289
static const uchar sc[16]
Definition smbdes.c:115
PUBLIC int snprintf(char *string, size_t length, char *format, va_alist)
Definition snprintf.c:689
char const * fr_syserror(int num)
Guaranteed to be thread-safe version of strerror.
Definition syserror.c:243
void fr_thread_start(fr_thread_t *thread, fr_sem_t *sem)
Signal the parent that we're done.
Definition thread.c:204
int fr_thread_wait_list(fr_sem_t *sem, fr_dlist_head_t *head)
Wait for multiple threads to signal readiness via a semaphore.
Definition thread.c:79
int fr_thread_create(pthread_t *thread, fr_thread_entry_t func, void *arg)
Create a joinable thread.
Definition thread.c:46
int fr_thread_setup(fr_thread_t *out, char const *name)
Common setup for child threads: block signals, allocate a talloc context, and create an event list.
Definition thread.c:112
int fr_thread_instantiate(TALLOC_CTX *ctx, fr_event_list_t *el)
Instantiate thread-specific data for modules, virtual servers, xlats, unlang, and TLS.
Definition thread.c:165
void fr_thread_detach(void)
Detach thread-specific data for modules, virtual servers, xlats.
Definition thread.c:190
void fr_thread_exit(fr_thread_t *thread, fr_thread_status_t status, fr_sem_t *sem)
Signal the parent that we're done.
Definition thread.c:219
fr_thread_status_t
Track the child thread status.
Definition thread.h:38
@ FR_THREAD_EXITED
exited, and in the exited queue
Definition thread.h:42
@ FR_THREAD_FAIL
failed, and in the exited queue
Definition thread.h:43
#define fr_time_delta_wrap(_time)
Definition time.h:152
"server local" time.
Definition time.h:69
static fr_event_list_t * el
#define fr_strerror_const(_msg)
Definition strerror.h:223
static fr_slen_t data
Definition value.h:1367