Coverage Report

Created: 2026-09-28 06:27

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/freeradius-server/src/lib/io/schedule.c
Line
Count
Source
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: 40e679483855d5ca2ab4f75b1ea0382444dc2459 $
19
 *
20
 * @brief Network / worker thread scheduling
21
 * @file io/schedule.c
22
 *
23
 * @copyright 2016 Alan DeKok (aland@freeradius.org)
24
 */
25
RCSID("$Id: 40e679483855d5ca2ab4f75b1ea0382444dc2459 $")
26
27
#define LOG_DST sc->log
28
29
#include <freeradius-devel/autoconf.h>
30
31
#include <freeradius-devel/io/schedule.h>
32
#include <freeradius-devel/io/thread.h>
33
#include <freeradius-devel/util/dlist.h>
34
#include <freeradius-devel/util/rb.h>
35
#include <freeradius-devel/util/syserror.h>
36
#include <freeradius-devel/server/module_rlm.h>
37
#include <freeradius-devel/server/trigger.h>
38
#include <freeradius-devel/util/semaphore.h>
39
40
#include <pthread.h>
41
42
/** Scheduler specific information for worker threads
43
 *
44
 * Wraps a fr_worker_t, tracking additional information that
45
 * the scheduler uses.
46
 */
47
typedef struct {
48
  fr_thread_t thread;     //!< common thread structure - must be first!
49
50
  int   uses;     //!< how many network threads are using it
51
  fr_time_t cpu_time;   //!< how much CPU time this worker has used
52
53
  fr_schedule_t *sc;      //!< the scheduler we are running under
54
55
  fr_worker_t *worker;    //!< the worker data structure
56
} fr_schedule_worker_t;
57
58
/** Scheduler specific information for network threads
59
 *
60
 * Wraps a fr_network_t, tracking additional information that
61
 * the scheduler uses.
62
 */
63
typedef struct {
64
  fr_thread_t thread;     //!< common thread structure - must be first!
65
66
  fr_schedule_t *sc;      //!< the scheduler we are running under
67
68
  fr_network_t  *nr;      //!< the receive data structure
69
70
  fr_timer_t  *ev;    //!< timer for stats_interval
71
} fr_schedule_network_t;
72
73
74
/**
75
 *  The scheduler
76
 */
77
struct fr_schedule_s {
78
  bool    running;    //!< is the scheduler running?
79
80
  CONF_SECTION  *cs;      //!< thread pool configuration section
81
  fr_event_list_t *el;      //!< event list for single-threaded mode.
82
  bool    single_threaded;  //!< true if running in single-threaded mode.
83
84
  fr_log_t  *log;     //!< log destination
85
  fr_log_lvl_t  lvl;      //!< log level
86
87
  fr_schedule_config_t *config;   //!< configuration
88
89
  unsigned int  num_workers_exited; //!< number of exited workers
90
91
  fr_sem_t  *worker_sem;    //!< for inter-thread signaling
92
  fr_sem_t  *network_sem;   //!< for inter-thread signaling
93
  fr_sem_t  *coord_sem;   //!< for inter-thread signaling
94
95
  fr_schedule_thread_instantiate_t  worker_thread_instantiate;  //!< thread instantiation callback
96
  fr_schedule_thread_detach_t   worker_thread_detach;
97
98
  fr_dlist_head_t workers;    //!< list of workers
99
  fr_dlist_head_t networks;   //!< list of networks
100
101
  fr_network_t  *single_network;  //!< for single-threaded mode
102
  fr_worker_t *single_worker;   //!< for single-threaded mode
103
};
104
105
static _Thread_local int worker_id = -1;  //!< Internal ID of the current worker thread.
106
107
/** Return the worker id for the current thread
108
 *
109
 * @return worker ID
110
 */
111
int fr_schedule_worker_id(void)
112
0
{
113
0
  return worker_id;
114
0
}
115
116
/** Explicitly set the worker id for the current thread
117
 *
118
 * **Only to be used in test programs like unit_test_module**
119
 */
120
void fr_schedule_worker_id_set(int id)
121
0
{
122
0
  worker_id = id;
123
0
}
124
125
/** Entry point for worker threads
126
 *
127
 * @param[in] arg the fr_schedule_worker_t
128
 * @return NULL
129
 */
130
static void *fr_schedule_worker_thread(void *arg)
131
0
{
132
0
  fr_schedule_worker_t    *sw = talloc_get_type_abort(arg, fr_schedule_worker_t);
133
0
  fr_schedule_t     *sc = sw->sc;
134
0
  fr_thread_status_t    status = FR_THREAD_FAIL;
135
0
  char        worker_name[32];
136
137
0
  worker_id = sw->thread.id;    /* Store the current worker ID */
138
139
0
  snprintf(worker_name, sizeof(worker_name), "Worker %d", sw->thread.id);
140
141
0
#ifdef HAVE_PTHREAD_SETNAME_NP
142
#  ifdef __APPLE__
143
  pthread_setname_np(worker_name);
144
#  else
145
0
  pthread_setname_np(pthread_self(), worker_name);
146
0
#  endif
147
0
#endif
148
149
0
  if (fr_thread_setup(&sw->thread, worker_name) < 0) goto fail;
150
151
0
  sw->worker = fr_worker_alloc(sw->thread.ctx, sw->thread.el, worker_name, sc->log, sc->lvl, &sc->config->worker);
152
0
  if (!sw->worker) {
153
0
    PERROR("%s - Failed creating worker", worker_name);
154
0
    goto fail;
155
0
  }
156
157
  /*
158
   *  @todo make this a registry
159
   */
160
0
  if (sc->worker_thread_instantiate) {
161
0
    CONF_SECTION  *cs;
162
0
    char    section_name[32];
163
164
0
    snprintf(section_name, sizeof(section_name), "%u", sw->thread.id);
165
166
0
    cs = cf_section_find(sc->cs, "worker", section_name);
167
0
    if (!cs) cs = cf_section_find(sc->cs, "worker", NULL);
168
169
0
    if (sc->worker_thread_instantiate(sw->thread.ctx, sw->thread.el, cs) < 0) {
170
0
      PERROR("%s - Worker thread instantiation failed", worker_name);
171
0
      goto fail;
172
0
    }
173
0
  }
174
175
  /*
176
   *  Add this worker to all network threads.
177
   */
178
0
  fr_dlist_foreach(&sc->networks, fr_schedule_network_t, sn) {
179
0
    if (unlikely(fr_network_worker_add(sn->nr, sw->worker) < 0)) {
180
0
      PERROR("%s - Failed adding worker to network %u", worker_name, sn->thread.id);
181
0
      goto fail;  /* FIXME - Should maybe try to undo partial adds? */
182
0
    }
183
0
  }
184
185
  /*
186
   *  Tell the originator that the thread has started.
187
   */
188
0
  fr_thread_start(&sw->thread, sc->worker_sem);
189
190
  /*
191
   *  Do all of the work.
192
   */
193
0
  fr_worker(sw->worker);
194
195
0
  status = FR_THREAD_EXITED;
196
197
0
fail:
198
0
  if (sw->worker) {
199
0
    fr_worker_exit(sw->worker);
200
0
    sw->worker = NULL;
201
0
  }
202
203
0
  if (sc->worker_thread_detach) sc->worker_thread_detach(NULL); /* Fixme once we figure out what uctx should be */
204
205
0
  fr_thread_exit(&sw->thread, status, sc->worker_sem);
206
207
0
  return NULL;
208
0
}
209
210
211
static void stats_timer(fr_timer_list_t *tl, fr_time_t now, void *uctx)
212
0
{
213
0
  fr_schedule_network_t   *sn = talloc_get_type_abort(uctx, fr_schedule_network_t);
214
215
0
  fr_network_stats_log(sn->nr, sn->sc->log);
216
217
0
  (void) fr_timer_at(sn, tl, &sn->ev, fr_time_add(now, sn->sc->config->stats_interval), false, stats_timer, sn);
218
0
}
219
220
/** Initialize and run the network thread.
221
 *
222
 * @param[in] arg the fr_schedule_network_t
223
 * @return NULL
224
 */
225
static void *fr_schedule_network_thread(void *arg)
226
0
{
227
0
  fr_schedule_network_t   *sn = talloc_get_type_abort(arg, fr_schedule_network_t);
228
0
  fr_schedule_t     *sc = sn->sc;
229
0
  fr_thread_status_t    status = FR_THREAD_FAIL;
230
0
  char        network_name[32];
231
232
0
  snprintf(network_name, sizeof(network_name), "Network %d", sn->thread.id);
233
234
0
#ifdef HAVE_PTHREAD_SETNAME_NP
235
#  ifdef __APPLE__
236
  pthread_setname_np(network_name);
237
#  else
238
0
  pthread_setname_np(pthread_self(), network_name);
239
0
#  endif
240
0
#endif
241
242
0
  if (fr_thread_setup(&sn->thread, network_name) < 0) goto fail;
243
244
0
  sn->nr = fr_network_create(sn->thread.ctx, sn->thread.el, network_name, sc->log, sc->lvl, &sc->config->network);
245
0
  if (!sn->nr) {
246
0
    PERROR("%s - Failed creating network", network_name);
247
0
    goto fail;
248
0
  }
249
250
  /*
251
   *  Tell the originator that the thread has started.
252
   */
253
0
  fr_thread_start(&sn->thread, sc->network_sem);
254
255
  /*
256
   *  Print out statistics for this network IO handler.
257
   */
258
0
  if (fr_time_delta_ispos(sc->config->stats_interval)) {
259
0
    (void) fr_timer_in(sn, sn->thread.el->tl, &sn->ev, sn->sc->config->stats_interval, false, stats_timer, sn);
260
0
  }
261
262
  /*
263
   *  Call the main event processing loop of the network
264
   *  thread Will not return until the worker is about
265
   *      to exit.
266
   */
267
0
  fr_network(sn->nr);
268
269
0
  status = FR_THREAD_EXITED;
270
271
0
fail:
272
0
  fr_thread_exit(&sn->thread, status, sc->network_sem);
273
274
0
  return NULL;
275
0
}
276
277
/** Create a scheduler and spawn the child threads.
278
 *
279
 * @param[in] ctx       talloc context.
280
 * @param[in] single_threaded     no workers are spawned, everything runs in a common event loop.
281
 * @param[in] el        event list, only for single-threaded mode.
282
 * @param[in] logger        destination for all logging messages.
283
 * @param[in] lvl       log level.
284
 * @param[in] worker_thread_instantiate   callback for new worker threads.
285
 * @param[in] worker_thread_detach    callback to destroy resources
286
 *            allocated by worker_thread_instantiate.
287
 * @param[in] config        configuration for the scheduler
288
 * @return
289
 *  - NULL on error
290
 *  - fr_schedule_t new scheduler
291
 */
292
fr_schedule_t *fr_schedule_create(TALLOC_CTX *ctx,
293
          bool single_threaded,
294
          fr_event_list_t *el,
295
          fr_log_t *logger, fr_log_lvl_t lvl,
296
          fr_schedule_thread_instantiate_t worker_thread_instantiate,
297
          fr_schedule_thread_detach_t worker_thread_detach,
298
          fr_schedule_config_t *config)
299
0
{
300
0
  unsigned int i;
301
0
  fr_schedule_worker_t *sw, *next_sw;
302
0
  fr_schedule_network_t *sn, *next_sn;
303
0
  fr_schedule_t *sc;
304
305
0
  sc = talloc_zero(ctx, fr_schedule_t);
306
0
  if (!sc) {
307
0
    fr_strerror_const("Failed allocating memory");
308
0
    return NULL;
309
0
  }
310
311
  /*
312
   *  Parse any scheduler-specific configuration.
313
   */
314
0
  if (!config) {
315
0
    MEM(sc->config = talloc_zero(sc, fr_schedule_config_t));
316
0
    sc->config->max_networks = 1;
317
0
    sc->config->max_workers = 4;
318
0
  } else {
319
0
    sc->config = config;
320
321
0
    if (sc->config->max_networks < 1) sc->config->max_networks = 1;
322
0
    if (sc->config->max_networks > 64) sc->config->max_networks = 64;
323
0
    if (sc->config->max_workers < 1) sc->config->max_workers = 1;
324
0
    if (sc->config->max_workers > 64) sc->config->max_workers = 64;
325
0
  }
326
327
0
  sc->el = el;
328
0
  sc->single_threaded = single_threaded;
329
0
  sc->log = logger;
330
0
  sc->lvl = lvl;
331
0
  sc->cs = sc->config->cs;
332
333
0
  sc->worker_thread_instantiate = worker_thread_instantiate;
334
0
  sc->worker_thread_detach = worker_thread_detach;
335
0
  sc->running = true;
336
337
  /*
338
   *  If we're single-threaded, create network / worker, and insert them into the event loop.
339
   */
340
0
  if (single_threaded) {
341
0
    sc->single_network = fr_network_create(sc, el, "Network", sc->log, sc->lvl, &sc->config->network);
342
0
    if (!sc->single_network) {
343
0
      PERROR("Failed creating network");
344
0
    pre_instantiate_st_fail:
345
0
      talloc_free(sc);
346
0
      return NULL;
347
0
    }
348
349
0
    if (fr_coords_create(sc, el) < 0) {
350
0
      PERROR("Failed creating coordinators");
351
0
      if (unlikely(fr_network_destroy(sc->single_network) < 0)) {
352
0
        PERROR("Failed destroying network");
353
0
      }
354
0
      goto pre_instantiate_st_fail;
355
0
    }
356
357
0
    worker_id = 0;
358
0
    sc->single_worker = fr_worker_alloc(sc, el, "Worker", sc->log, sc->lvl, &sc->config->worker);
359
0
    if (!sc->single_worker) {
360
0
      PERROR("Failed creating worker");
361
0
      if (unlikely(fr_network_destroy(sc->single_network) < 0)) {
362
0
        PERROR("Failed destroying network");
363
0
      }
364
0
      goto pre_instantiate_st_fail;
365
0
    }
366
367
    /*
368
     *  Parent thread-specific data from the single_worker
369
     */
370
0
    if (sc->worker_thread_instantiate) {
371
0
      CONF_SECTION *subcs;
372
373
0
      subcs = cf_section_find(sc->cs, "worker", "0");
374
0
      if (!subcs) subcs = cf_section_find(sc->cs, "worker", NULL);
375
376
0
      if (sc->worker_thread_instantiate(sc->single_worker, el, subcs) < 0) {
377
0
        PERROR("Worker thread instantiation failed");
378
0
      destroy_both:
379
0
        if (unlikely(fr_network_destroy(sc->single_network) < 0)) {
380
0
          PERROR("Failed destroying network");
381
0
        }
382
0
        fr_worker_exit(sc->single_worker);
383
0
        goto pre_instantiate_st_fail;
384
0
      }
385
0
    }
386
387
0
    if (fr_command_register_hook(NULL, "0", sc->single_worker, cmd_worker_table) < 0) {
388
0
      PERROR("Failed adding worker commands");
389
0
    st_fail:
390
0
      if (sc->worker_thread_detach) sc->worker_thread_detach(NULL);
391
0
      goto destroy_both;
392
0
    }
393
394
0
    if (fr_command_register_hook(NULL, "0", sc->single_network, cmd_network_table) < 0) {
395
0
      PERROR("Failed adding network commands");
396
0
      goto st_fail;
397
0
    }
398
399
    /*
400
     *  Register the worker with the network, so
401
     *  things like fr_network_send_request() work.
402
     */
403
0
    fr_network_worker_add_self(sc->single_network, sc->single_worker);
404
0
    DEBUG("Scheduler created in single-threaded mode");
405
406
0
    if (fr_event_pre_insert(el, fr_worker_pre_event, sc->single_worker) < 0) {
407
0
      fr_strerror_const("Failed adding pre-check to event list");
408
0
      goto st_fail;
409
0
    }
410
411
0
    if (fr_coord_pre_event_insert(el) < 0) {
412
0
      fr_strerror_const("Failed adding coordinator pre-check to event list");
413
0
      goto st_fail;
414
0
    }
415
416
    /*
417
     *  Add the event which processes request_t packets.
418
     */
419
0
    if (fr_event_post_insert(el, fr_worker_post_event, sc->single_worker) < 0) {
420
0
      fr_strerror_const("Failed inserting post-processing event");
421
0
      goto st_fail;
422
0
    }
423
424
0
    if (fr_coord_post_event_insert(el) < 0) {
425
0
      fr_strerror_const("Failed adding coordinator post-processing to event list");
426
0
      goto st_fail;
427
0
    }
428
429
0
    return sc;
430
0
  }
431
432
  /*
433
   *  Create the lists which hold the workers and networks.
434
   */
435
0
  fr_dlist_init(&sc->workers, fr_schedule_worker_t, thread.entry);
436
0
  fr_dlist_init(&sc->networks, fr_schedule_network_t, thread.entry);
437
438
0
  sc->network_sem = fr_sem_alloc();
439
0
  if (!sc->network_sem) {
440
0
  sem_fail:
441
0
    ERROR("Failed creating semaphore: %s", fr_syserror(errno));
442
0
    fr_sem_free(sc->network_sem);
443
0
    fr_sem_free(sc->worker_sem);
444
0
    talloc_free(sc);
445
0
    return NULL;
446
0
  }
447
448
0
  sc->worker_sem = fr_sem_alloc();
449
0
  if (!sc->worker_sem) goto sem_fail;
450
451
0
  sc->coord_sem = fr_sem_alloc();
452
0
  if (!sc->coord_sem) goto sem_fail;
453
454
  /*
455
   *  Create the network threads first.
456
   */
457
0
  for (i = 0; i < sc->config->max_networks; i++) {
458
0
    DEBUG3("Creating %u/%u networks", i + 1, sc->config->max_networks);
459
460
    /*
461
     *  Create a worker "glue" structure
462
     */
463
0
    sn = talloc_zero(sc, fr_schedule_network_t);
464
0
    if (!sn) {
465
0
      ERROR("Network %u - Failed allocating memory", i);
466
0
      break;
467
0
    }
468
469
0
    sn->thread.id = i;
470
0
    sn->sc = sc;
471
0
    sn->thread.status = FR_THREAD_INITIALIZING;
472
473
0
    if (fr_thread_create(&sn->thread.pthread_id, fr_schedule_network_thread, sn) < 0) {
474
0
      talloc_free(sn);
475
0
      PERROR("Failed creating network %u", i);
476
0
      break;
477
0
    }
478
479
0
    fr_dlist_insert_head(&sc->networks, sn);
480
0
  }
481
482
  /*
483
   *  Wait for all of the networks to signal us that either
484
   *  they've started, OR there's been a problem and they
485
   *  can't start.
486
   */
487
0
  if (fr_thread_wait_list(sc->network_sem, &sc->networks) < 0) {
488
0
    fr_schedule_destroy(&sc);
489
0
    return NULL;
490
0
  }
491
492
  /*
493
   *  Create the coordination threads
494
   */
495
0
  if (fr_coord_start(sc->config->max_workers, sc->coord_sem) < 0) {
496
0
    fr_schedule_destroy(&sc);
497
0
    return NULL;
498
0
  };
499
500
  /*
501
   *  Create all of the workers.
502
   */
503
0
  for (i = 0; i < sc->config->max_workers; i++) {
504
0
    DEBUG3("Creating %u/%u workers", i + 1, sc->config->max_workers);
505
506
    /*
507
     *  Create a worker "glue" structure
508
     */
509
0
    sw = talloc_zero(sc, fr_schedule_worker_t);
510
0
    if (!sw) {
511
0
      ERROR("Worker %u - Failed allocating memory", i);
512
0
      break;
513
0
    }
514
515
0
    sw->thread.id = i;
516
0
    sw->sc = sc;
517
0
    sw->thread.status = FR_THREAD_INITIALIZING;
518
519
0
    if (fr_thread_create(&sw->thread.pthread_id, fr_schedule_worker_thread, sw) < 0) {
520
0
      talloc_free(sw);
521
0
      PERROR("Failed creating worker %u", i);
522
0
      break;
523
0
    }
524
525
0
    fr_dlist_insert_head(&sc->workers, sw);
526
0
  }
527
528
  /*
529
   *  Wait for all of the workers to signal us that either
530
   *  they've started, OR there's been a problem and they
531
   *  can't start.
532
   */
533
0
  if (fr_thread_wait_list(sc->worker_sem, &sc->workers) < 0) {
534
0
    fr_schedule_destroy(&sc);
535
0
    return NULL;
536
0
  }
537
538
0
  for (sw = fr_dlist_head(&sc->workers), i = 0;
539
0
       sw != NULL;
540
0
       sw = next_sw, i++) {
541
0
    char buffer[32];
542
543
0
    next_sw = fr_dlist_next(&sc->workers, sw);
544
545
0
    snprintf(buffer, sizeof(buffer), "%d", i);
546
0
    if (fr_command_register_hook(NULL, buffer, sw->worker, cmd_worker_table) < 0) {
547
0
      PERROR("Failed adding worker commands");
548
0
    mt_fail:
549
0
      fr_schedule_destroy(&sc);
550
0
      return NULL;
551
0
    }
552
0
  }
553
554
0
  for (sn = fr_dlist_head(&sc->networks), i = 0;
555
0
       sn != NULL;
556
0
       sn = next_sn, i++) {
557
0
    char buffer[32];
558
559
0
    next_sn = fr_dlist_next(&sc->networks, sn);
560
561
0
    snprintf(buffer, sizeof(buffer), "%d", i);
562
0
    if (fr_command_register_hook(NULL, buffer, sn->nr, cmd_network_table) < 0) {
563
0
      PERROR("Failed adding network commands");
564
0
      goto mt_fail;
565
0
    }
566
0
  }
567
568
0
  if (sc) INFO("Scheduler created successfully with %u networks and %u workers",
569
0
         sc->config->max_networks, (unsigned int)fr_dlist_num_elements(&sc->workers));
570
571
  /*
572
   *  Instantiate thread-local data for the main thread too.
573
   *  In single-threaded mode this is done above.  In
574
   *  multi-worker mode the main thread also needs module
575
   *  thread data so that triggers can use module xlats.
576
   */
577
0
  if (sc->worker_thread_instantiate &&
578
0
      unlikely((sc->worker_thread_instantiate(sc, el, NULL) < 0))) {
579
0
    PERROR("Main thread instantiation failed");
580
0
    goto mt_fail;
581
0
  }
582
583
0
  return sc;
584
0
}
585
586
/** Destroy a scheduler, and tell its child threads to exit.
587
 *
588
 * @note This may be called with no worker or network threads in the case of a
589
 *   instantiation error.  This function _should_ deal with that condition
590
 *   gracefully.
591
 *
592
 * @param[in] sc_to_free the scheduler
593
 * @return
594
 *  - <0 on error
595
 *  - 0 on success
596
 */
597
int fr_schedule_destroy(fr_schedule_t **sc_to_free)
598
0
{
599
0
  fr_schedule_t   *sc = *sc_to_free;
600
0
  unsigned int    i;
601
0
  fr_schedule_worker_t  *sw;
602
0
  fr_schedule_network_t *sn;
603
0
  int     ret;
604
605
0
  if (!sc) return 0;
606
607
0
  sc->running = false;
608
609
610
611
  /*
612
   *  Single threaded mode: kill the only network / worker we have.
613
   */
614
0
  if (sc->single_threaded) {
615
    /*
616
     *  Destroy the network side first.  It tells the
617
     *  workers to close.
618
     */
619
0
    if (unlikely(fr_network_destroy(sc->single_network) < 0)) {
620
0
      ERROR("Failed destroying network");
621
0
    }
622
623
    /*
624
     *  Add events to handle close down gracefully.
625
     */
626
0
    if (unlikely((fr_network_close_event_insert(sc->single_network) < 0) ||
627
0
           (fr_worker_close_event_insert(sc->single_worker) < 0))) {
628
0
      ERROR("Failed setting up close events");
629
0
    }
630
631
    /*
632
     *  Run the event loop so the worker gets the signal from
633
     *  the network and shuts down gracefully.
634
     */
635
0
    fr_event_loop(sc->el);
636
637
    /*
638
     *  Detach worker from coordinators.  This needs to be done
639
     *  before the worker is freed.
640
     */
641
0
    if (modules_rlm_coord_detach() > 0) {
642
      /*
643
       *  Run the event loop again to handle coordinator detach
644
       *  messages, with different callbacks to determine when
645
       *  to exit the loop.
646
       */
647
0
      fr_network_close_event_delete(sc->single_network);
648
0
      if (unlikely(fr_coord_close_event_insert(sc->el) < 0)) {
649
0
        ERROR("Failed setting up coordinator close events");
650
0
      }
651
0
      fr_event_loop(sc->el);
652
0
    }
653
654
0
    fr_worker_exit(sc->single_worker);
655
0
    fr_coords_destroy();
656
657
0
    goto done;
658
0
  } else {
659
    /*
660
     *  Detach thread-local data for the main thread.
661
     *  Worker threads handle their own detach, but
662
     *  the main thread was instantiated explicitly
663
     *  by fr_schedule_create.
664
     */
665
0
    if (sc->worker_thread_detach) sc->worker_thread_detach(NULL);
666
0
  }
667
668
  /*
669
   *  Signal each network thread to exit.
670
   */
671
0
  fr_dlist_foreach(&sc->networks, fr_schedule_network_t, sne) {
672
0
    if (fr_network_exit(sne->nr) < 0) {
673
0
      PERROR("Failed signaling network %i to exit", sne->thread.id);
674
0
    }
675
0
  }
676
677
  /*
678
   *  If the network threads are running, tell them to exit,
679
   *  and wait for them to do so.  Each network thread tells
680
   *  all of its worker threads that it's exiting.  It then
681
   *  closes the channels.  When the workers see that there
682
   *  are no input channels, they exit, too.
683
   */
684
0
  for (i = 0; i < (unsigned int)fr_dlist_num_elements(&sc->networks); i++) {
685
0
    DEBUG2("Scheduler - Waiting for semaphore indicating network exit %u/%u", i + 1,
686
0
           (unsigned int)fr_dlist_num_elements(&sc->networks));
687
0
    SEM_WAIT_INTR(sc->network_sem);
688
0
  }
689
0
  DEBUG2("Scheduler - All networks indicated exit complete");
690
691
0
  while ((sn = fr_dlist_pop_head(&sc->networks)) != NULL) {
692
    /*
693
     *  Ensure that the thread has exited before
694
     *  cleaning up the context.
695
     *
696
     *  This also ensures that the child threads have
697
     *  exited before the main thread cleans up the
698
     *  module instances.
699
     */
700
0
    if ((ret = pthread_join(sn->thread.pthread_id, NULL)) != 0) {
701
0
      ERROR("Failed joining network %i: %s", sn->thread.id, fr_syserror(ret));
702
0
    } else {
703
0
      DEBUG2("Network %i joined (cleaned up)", sn->thread.id);
704
0
    }
705
0
  }
706
707
  /*
708
   *  Wait for all worker threads to finish.  THEN clean up
709
   *  modules.  Otherwise, the modules will be removed from
710
   *  underneath the workers!
711
   */
712
0
  for (i = 0; i < (unsigned int)fr_dlist_num_elements(&sc->workers); i++) {
713
0
    DEBUG2("Scheduler - Waiting for semaphore indicating worker exit %u/%u", i + 1,
714
0
           (unsigned int)fr_dlist_num_elements(&sc->workers));
715
0
    SEM_WAIT_INTR(sc->worker_sem);
716
0
  }
717
0
  DEBUG2("Scheduler - All workers indicated exit complete");
718
719
  /*
720
   *  Clean up the exited workers.
721
   */
722
0
  while ((sw = fr_dlist_pop_head(&sc->workers)) != NULL) {
723
    /*
724
     *  Ensure that the thread has exited before
725
     *  cleaning up the context.
726
     *
727
     *  This also ensures that the child threads have
728
     *  exited before the main thread cleans up the
729
     *  module instances.
730
     */
731
0
    if ((ret = pthread_join(sw->thread.pthread_id, NULL)) != 0) {
732
0
      ERROR("Failed joining worker %i: %s", sw->thread.id, fr_syserror(ret));
733
0
    } else {
734
0
      DEBUG2("Worker %i joined (cleaned up)", sw->thread.id);
735
0
    }
736
0
  }
737
738
0
  fr_coord_thread_join();
739
740
0
  fr_sem_free(sc->coord_sem);
741
0
  fr_sem_free(sc->network_sem);
742
0
  fr_sem_free(sc->worker_sem);
743
744
0
done:
745
  /*
746
   *  Now that all of the workers are done, we can return to
747
   *  the caller, and have it dlclose() the modules.
748
   */
749
0
  talloc_free(sc);
750
0
  *sc_to_free = NULL;
751
752
0
  return 0;
753
0
}
754
755
/** Add a fr_listen_t to a scheduler.
756
 *
757
 * @param[in] sc the scheduler
758
 * @param[in] li the ctx and callbacks for the transport.
759
 * @return
760
 *  - NULL on error
761
 *  - the fr_network_t that the socket was added to.
762
 */
763
fr_network_t *fr_schedule_listen_add(fr_schedule_t *sc, fr_listen_t *li)
764
0
{
765
0
  fr_network_t *nr;
766
767
0
  (void) talloc_get_type_abort(sc, fr_schedule_t);
768
769
0
  if (sc->single_threaded) {
770
0
    nr = sc->single_network;
771
0
  } else {
772
0
    fr_schedule_network_t *sn;
773
774
    /*
775
     *  @todo - round robin it among the listeners?
776
     *  or maybe add it to the same parent thread?
777
     */
778
0
    sn = fr_dlist_head(&sc->networks);
779
0
    nr = sn->nr;
780
0
  }
781
782
0
  if (fr_network_listen_add(nr, li) < 0) return NULL;
783
784
0
  return nr;
785
0
}
786
787
/** Add a directory NOTE_EXTEND to a scheduler.
788
 *
789
 * @param[in] sc the scheduler
790
 * @param[in] li the ctx and callbacks for the transport.
791
 * @return
792
 *  - NULL on error
793
 *  - the fr_network_t that the socket was added to.
794
 */
795
fr_network_t *fr_schedule_directory_add(fr_schedule_t *sc, fr_listen_t *li)
796
0
{
797
0
  fr_network_t *nr;
798
799
0
  (void) talloc_get_type_abort(sc, fr_schedule_t);
800
801
0
  if (sc->single_threaded) {
802
0
    nr = sc->single_network;
803
0
  } else {
804
0
    fr_schedule_network_t *sn;
805
806
    /*
807
     *  @todo - round robin it among the listeners?
808
     *  or maybe add it to the same parent thread?
809
     */
810
0
    sn = fr_dlist_head(&sc->networks);
811
0
    nr = sn->nr;
812
0
  }
813
814
0
  if (fr_network_directory_add(nr, li) < 0) return NULL;
815
816
0
  return nr;
817
0
}