]> git.ipfire.org Git - thirdparty/freeradius-server.git/commitdiff
Change the way workers are signalled to exit
authorArran Cudbard-Bell <a.cudbardb@freeradius.org>
Mon, 16 Dec 2019 10:26:17 +0000 (17:26 +0700)
committerArran Cudbard-Bell <a.cudbardb@freeradius.org>
Mon, 16 Dec 2019 10:26:25 +0000 (17:26 +0700)
They now wait for all their channels to close, then signal themselves to exit via their event loop.

src/lib/io/schedule.c
src/lib/io/worker.c

index f51d92c8207d6fbefff9578379c2439ac02485ce..e439d522bb52b1bf906dd1b808058e14c3d78235 100644 (file)
@@ -125,7 +125,8 @@ struct fr_schedule_s {
 
        unsigned int    num_workers_exited;     //!< number of exited workers
 
-       sem_t           semaphore;              //!< for inter-thread signaling
+       sem_t           worker_sem;             //!< for inter-thread signaling
+       sem_t           network_sem;            //!< for inter-thread signaling
 
        fr_schedule_thread_instantiate_t        worker_thread_instantiate;      //!< thread instantiation callback
        void                                    *worker_instantiate_ctx;        //!< thread instantiation context
@@ -206,7 +207,7 @@ static void *fr_schedule_worker_thread(void *arg)
        /*
         *      Tell the originator that the thread has started.
         */
-       sem_post(&sc->semaphore);
+       sem_post(&sc->worker_sem);
 
        /*
         *      Do all of the work.
@@ -232,7 +233,7 @@ fail:
        /*
         *      Tell the scheduler we're done.
         */
-       sem_post(&sc->semaphore);
+       sem_post(&sc->worker_sem);
 
        return NULL;
 }
@@ -277,7 +278,7 @@ static void *fr_schedule_network_thread(void *arg)
        /*
         *      Tell the originator that the thread has started.
         */
-       sem_post(&sc->semaphore);
+       sem_post(&sc->network_sem);
 
        DEBUG3("Spawned asycn network 0");
 
@@ -296,7 +297,7 @@ fail:
        /*
         *      Tell the scheduler we're done.
         */
-       sem_post(&sc->semaphore);
+       sem_post(&sc->network_sem);
 
        return NULL;
 }
@@ -446,8 +447,15 @@ fr_schedule_t *fr_schedule_create(TALLOC_CTX *ctx, fr_event_list_t *el,
         */
        fr_dlist_init(&sc->workers, fr_schedule_worker_t, entry);
 
-       memset(&sc->semaphore, 0, sizeof(sc->semaphore));
-       if (sem_init(&sc->semaphore, 0, SEMAPHORE_LOCKED) != 0) {
+       memset(&sc->network_sem, 0, sizeof(sc->network_sem));
+       if (sem_init(&sc->network_sem, 0, SEMAPHORE_LOCKED) != 0) {
+               ERROR("Failed creating semaphore: %s", fr_syserror(errno));
+               talloc_free(sc);
+               return NULL;
+       }
+
+       memset(&sc->worker_sem, 0, sizeof(sc->worker_sem));
+       if (sem_init(&sc->worker_sem, 0, SEMAPHORE_LOCKED) != 0) {
                ERROR("Failed creating semaphore: %s", fr_syserror(errno));
                talloc_free(sc);
                return NULL;
@@ -466,7 +474,7 @@ fr_schedule_t *fr_schedule_create(TALLOC_CTX *ctx, fr_event_list_t *el,
                goto fail;
        }
 
-       SEM_WAIT_INTR(&sc->semaphore);
+       SEM_WAIT_INTR(&sc->network_sem);
        if (sc->sn->status != FR_CHILD_RUNNING) {
        fail:
                fr_schedule_destroy(sc);
@@ -510,7 +518,7 @@ fr_schedule_t *fr_schedule_create(TALLOC_CTX *ctx, fr_event_list_t *el,
        for (i = 0; i < (unsigned int)fr_dlist_num_elements(&sc->workers); i++) {
                DEBUG3("Waiting for semaphore from worker %u/%u",
                       i, (unsigned int)fr_dlist_num_elements(&sc->workers));
-               SEM_WAIT_INTR(&sc->semaphore);
+               SEM_WAIT_INTR(&sc->worker_sem);
        }
 
        /*
@@ -599,19 +607,10 @@ int fr_schedule_destroy(fr_schedule_t *sc)
         */
        if (sc->sn->status == FR_CHILD_RUNNING) {
                fr_network_exit(sc->sn->nr);
-               SEM_WAIT_INTR(&sc->semaphore);
+               SEM_WAIT_INTR(&sc->network_sem);
                fr_network_destroy(sc->sn->nr);
        }
 
-       /*
-        *      Signal all of the workers to exit.
-        */
-       for (sw = fr_dlist_head(&sc->workers);
-            sw != NULL;
-            sw = fr_dlist_next(&sc->workers, sw)) {
-               fr_worker_exit(sw->worker);
-       }
-
        /*
         *      Wait for all worker threads to finish.  THEN clean up
         *      modules.  Otherwise, the modules will be removed from
@@ -620,7 +619,7 @@ int fr_schedule_destroy(fr_schedule_t *sc)
        for (i = 0; i < (unsigned int)fr_dlist_num_elements(&sc->workers); i++) {
                DEBUG2("Waiting for semaphore indicating exit %u/%u", i,
                       (unsigned int)fr_dlist_num_elements(&sc->workers));
-               SEM_WAIT_INTR(&sc->semaphore);
+               SEM_WAIT_INTR(&sc->network_sem);
        }
 
        /*
@@ -647,8 +646,8 @@ int fr_schedule_destroy(fr_schedule_t *sc)
 
        TALLOC_FREE(sc->sn->ctx);
 
-       sem_destroy(&sc->semaphore);
-
+       sem_destroy(&sc->network_sem);
+       sem_destroy(&sc->worker_sem);
 done:
        /*
         *      Now that all of the workers are done, we can return to
index fa078800ef2b03629edeb913132072b64d6e44b5..0c6904b6bd8c487ce2467b66d914c7c748cd0e2c 100644 (file)
@@ -288,17 +288,9 @@ static void fr_worker_channel_callback(void *ctx, void const *data, size_t data_
 
                        if (worker->channel[i] != ch) continue;
 
-                       /*
-                        *      @todo check the status, and
-                        *      put the channel into a
-                        *      "closing" list if we can't
-                        *      close it right now.  Then,
-                        *      wake up after a time and try
-                        *      to close it again.
-                        */
-                       (void) fr_channel_responder_ack_close(ch);
-
                        ms = fr_channel_responder_uctx_get(ch);
+
+                       fr_channel_responder_ack_close(ch);
                        rad_assert(ms != NULL);
                        fr_message_set_gc(ms);
                        talloc_free(ms);
@@ -311,6 +303,12 @@ static void fr_worker_channel_callback(void *ctx, void const *data, size_t data_
                }
 
                fr_cond_assert(ok);
+
+               /*
+                *      Our last input channel closed,
+                *      time to die.
+                */
+               if (worker->num_channels == 0) fr_worker_exit(worker);
                break;
        }
 }