/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 | } |