diff options
| author | steven-na <noreply.github@stvnc.dev> | 2026-08-15 21:12:11 -0700 |
|---|---|---|
| committer | steven-na <noreply.github@stvnc.dev> | 2026-08-15 21:12:11 -0700 |
| commit | 82e164f8796d57aa78c14b20f0a750f4b40c861b (patch) | |
| tree | 1eb0df74a98ddc3fa872084b02ba309b3b0fb314 | |
| parent | a4e71241ffbcc8576cd42856db061bb366cec3e5 (diff) | |
Threadpool timeout
| -rw-r--r-- | Justfile.old | 99 | ||||
| -rw-r--r-- | src/threadpool.c | 86 | ||||
| -rw-r--r-- | src/threadpool.h | 11 |
3 files changed, 70 insertions, 126 deletions
diff --git a/Justfile.old b/Justfile.old deleted file mode 100644 index d1f2d94..0000000 --- a/Justfile.old +++ /dev/null @@ -1,99 +0,0 @@ -lib := "steez" -src_dir := "src" -test_dir := "tests" -build_dir := "build" -cc := "cc" -cflags := "-Wall -Wextra -Wpedantic -Werror -std=c23 -fPIC" -lflags := "-lm" -prefix := "/usr/local" -extra_cflags := "" - -default: - @just --list - -build: build-static build-shared - -gen-version type="release": - ./scripts/gen-version.sh {{type}} - -build-static: (gen-version "release") - mkdir -p {{build_dir}} - cd {{build_dir}} && {{cc}} {{cflags}} {{extra_cflags}} -c ../{{src_dir}}/*.c - ar rcs {{build_dir}}/lib{{lib}}.a {{build_dir}}/*.o - -build-shared: (gen-version "release") - mkdir -p {{build_dir}} - cd {{build_dir}} && {{cc}} {{cflags}} {{extra_cflags}} -c ../{{src_dir}}/*.c - {{cc}} -shared -o {{build_dir}}/lib{{lib}}.so {{build_dir}}/*.o {{lflags}} - -build-release: (gen-version "release") - mkdir -p {{build_dir}} - cd {{build_dir}} && {{cc}} {{cflags}} {{extra_cflags}} -O3 -DNDEBUG -c ../{{src_dir}}/*.c - ar rcs {{build_dir}}/lib{{lib}}-release.a {{build_dir}}/*.o - {{cc}} -shared -o {{build_dir}}/lib{{lib}}-release.so {{build_dir}}/*.o {{lflags}} - -build-debug: (gen-version "debug") - mkdir -p {{build_dir}} - cd {{build_dir}} && {{cc}} {{cflags}} {{extra_cflags}} -g -O0 -fsanitize=address,undefined -fno-omit-frame-pointer -c ../{{src_dir}}/*.c - ar rcs {{build_dir}}/lib{{lib}}-debug.a {{build_dir}}/*.o - {{cc}} -shared -fsanitize=address,undefined -o {{build_dir}}/lib{{lib}}-debug.so {{build_dir}}/*.o {{lflags}} - -install: build - sudo mkdir -p {{prefix}}/include/{{lib}} {{prefix}}/lib - sudo cp {{src_dir}}/*.h {{prefix}}/include/{{lib}}/ - sudo cp {{build_dir}}/lib{{lib}}.a {{build_dir}}/lib{{lib}}.so {{prefix}}/lib/ - sudo ldconfig - -uninstall: - sudo rm -rf {{prefix}}/include/{{lib}} - sudo rm -f {{prefix}}/lib/lib{{lib}}.a {{prefix}}/lib/lib{{lib}}.so - sudo ldconfig - -install-vendor dir: build - #!/usr/bin/env bash - set -euo pipefail - dir="{{dir}}" - dir="${dir%/}" - vendor_dir="$dir/vendor/{{lib}}" - mkdir -p "$vendor_dir/include/{{lib}}" "$vendor_dir/lib" - cp {{src_dir}}/*.h "$vendor_dir/include/{{lib}}/" - cp {{build_dir}}/lib{{lib}}.a {{build_dir}}/lib{{lib}}.so "$vendor_dir/lib/" - echo "Installed lib{{lib}} into $vendor_dir" - - target_justfile="$dir/Justfile" - if [ -f "$target_justfile" ] && grep -q '^cflags :=' "$target_justfile" && grep -q '^lflags :=' "$target_justfile"; then - echo "Detected ctmplt-layout project at $dir; wiring up build flags" - inc_flag="-Ivendor/{{lib}}/include" - static_lib="vendor/{{lib}}/lib/lib{{lib}}.a" - if ! grep -qF -- "$inc_flag" "$target_justfile"; then - sed -i "s|^cflags := \"\(.*\)\"|cflags := \"\1 $inc_flag\"|" "$target_justfile" - fi - if ! grep -qF -- "$static_lib" "$target_justfile"; then - sed -i "s|^lflags := \"\(.*\)\"|lflags := \"\1 $static_lib\"|" "$target_justfile" - fi - echo "Updated $target_justfile (cflags/lflags)" - else - echo "$target_justfile doesn't match the ctmplt bin/lib layout (no cflags:=/lflags:= Justfile); skipping auto-link wiring" - fi - -uninstall-vendor dir: - rm -rf {{dir}}/vendor/{{lib}} - -test-all: (gen-version "debug") - mkdir -p {{build_dir}} - {{cc}} {{cflags}} {{extra_cflags}} {{test_dir}}/*.c {{src_dir}}/*.c -o {{build_dir}}/test -lcriterion {{lflags}} - ./{{build_dir}}/test - -coverage: (gen-version "debug") - mkdir -p {{build_dir}}/cov-src {{build_dir}}/cov-tests - cd {{build_dir}}/cov-src && {{cc}} {{cflags}} --coverage -c ../../{{src_dir}}/*.c - cd {{build_dir}}/cov-tests && {{cc}} {{cflags}} --coverage -c ../../{{test_dir}}/*.c - {{cc}} {{cflags}} --coverage {{build_dir}}/cov-src/*.o {{build_dir}}/cov-tests/*.o -o {{build_dir}}/test-coverage -lcriterion {{lflags}} - ./{{build_dir}}/test-coverage - gcovr --root . --filter '{{src_dir}}/' {{build_dir}}/cov-src --html-details {{build_dir}}/coverage.html --print-summary - -clean: - rm -rf {{build_dir}} - -compiledb: - bear -- just build diff --git a/src/threadpool.c b/src/threadpool.c index 977f73c..152ca66 100644 --- a/src/threadpool.c +++ b/src/threadpool.c @@ -1,3 +1,9 @@ +#if defined(__linux__) + #ifndef _DEFAULT_SOURCE + #define _DEFAULT_SOURCE + #endif /* ifndef _DEFAULT_SOURCE */ +#endif + #include "common.h" #include "threadpool.h" #include "smrt_arena.h" @@ -5,8 +11,11 @@ #include "log.h" #include <bits/pthreadtypes.h> +#include <bits/time.h> #include <pthread.h> #include <string.h> +#include <time.h> +#include <asm-generic/errno.h> typedef struct { tp_job_proc proc; @@ -27,22 +36,28 @@ void *worker_proc(void *args) { #endif /* ifndef NLOG_TRACE */ // Signal to tp_create that threads are done spinning up - pthread_mutex_lock(tp->count_mtx); { + pthread_mutex_lock(tp->access_mtx); { if (++(tp->num_alive) == tp->num_threads) { pthread_cond_broadcast(tp->done_signal); } - } pthread_mutex_unlock(tp->count_mtx); + } pthread_mutex_unlock(tp->access_mtx); for (;;) { tp_job_t *j = (tp_job_t*)ts_deque_pop(tp->jobs); - pthread_mutex_lock(tp->count_mtx); + pthread_mutex_lock(tp->access_mtx); tp->num_active++; - pthread_mutex_unlock(tp->count_mtx); + pthread_mutex_unlock(tp->access_mtx); j->proc(j->args); - pthread_mutex_lock(tp->count_mtx); + pthread_mutex_lock(tp->access_mtx); + + if (tp->is_abandoned) { + pthread_mutex_unlock(tp->access_mtx); + return NULL; + } + // If this is the last active thread finishing, tell someone about it if (--(tp->num_active) == 0) { pthread_cond_broadcast(tp->done_signal); @@ -53,21 +68,21 @@ void *worker_proc(void *args) { } if (!tp->is_running) { - pthread_mutex_unlock(tp->count_mtx); + pthread_mutex_unlock(tp->access_mtx); break; } - pthread_mutex_unlock(tp->count_mtx); + pthread_mutex_unlock(tp->access_mtx); } #ifndef NLOG_TRACE log_trace("Spinning down thread %lu", w_args->id); #endif /* ifndef NLOG_TRACE */ - pthread_mutex_lock(tp->count_mtx); { + pthread_mutex_lock(tp->access_mtx); { if (--(tp->num_alive) == 0) { pthread_cond_broadcast(tp->done_signal); } - } pthread_mutex_unlock(tp->count_mtx); + } pthread_mutex_unlock(tp->access_mtx); return NULL; } @@ -94,13 +109,13 @@ thread_pool_t *tp_create(smrt_arena_t *arena, u64 max_jobs, u64 num_threads) { .jobs=tsq, .is_running=true, - .count_mtx=count_mtx, + .access_mtx=count_mtx, .done_signal=all_done, .num_alive=0, .num_active=0, }; - pthread_mutex_lock(tp->count_mtx); + pthread_mutex_lock(tp->access_mtx); for (u64 i = 0; i < num_threads; i++) { args[i] = (worker_args){ .id = i, @@ -110,9 +125,9 @@ thread_pool_t *tp_create(smrt_arena_t *arena, u64 max_jobs, u64 num_threads) { } while (tp->num_alive != num_threads) { - pthread_cond_wait(tp->done_signal, tp->count_mtx); + pthread_cond_wait(tp->done_signal, tp->access_mtx); } - pthread_mutex_unlock(tp->count_mtx); + pthread_mutex_unlock(tp->access_mtx); #ifndef NLOG_TRACE @@ -124,33 +139,54 @@ thread_pool_t *tp_create(smrt_arena_t *arena, u64 max_jobs, u64 num_threads) { void no_op(void *nothing) { (void)nothing; } -i32 tp_destroy(thread_pool_t *tp) { +i32 tp_destroy(thread_pool_t *tp, i32 timeout) { #ifndef NLOG_TRACE log_trace("Destroying threadpool"); #endif /* ifndef NLOG_TRACE */ - pthread_mutex_lock(tp->count_mtx); + pthread_mutex_lock(tp->access_mtx); tp->is_running = false; + b32 timeout_enable = timeout != INT32_MAX; + struct timespec deadline; + if (timeout_enable) { + clock_gettime(CLOCK_REALTIME, &deadline); + deadline.tv_sec += 5; + } + // Hand out Kool Aid u64 pills = tp->num_alive; for (u64 i = 0; i < pills; i++) tp_push_job(tp, no_op, NULL); - while (tp->num_alive != 0) { - pthread_cond_wait(tp->done_signal, tp->count_mtx); - } - for (u64 t = 0; t < tp->num_threads; t++) { - pthread_join(tp->threads[t], NULL); + while (tp->num_alive != 0) { + if (timeout_enable) { + if (pthread_cond_timedwait(tp->done_signal, + tp->access_mtx, + &deadline) == ETIMEDOUT) + { goto abandon; } + } /* Don't destroy infrastructure if abandoning */ + else { + pthread_cond_wait(tp->done_signal, tp->access_mtx); + } } - pthread_mutex_unlock(tp->count_mtx); + pthread_mutex_unlock(tp->access_mtx); + for (u64 t = 0; t < tp->num_threads; t++) pthread_join(tp->threads[t], NULL); ts_deque_destroy(tp->jobs); - pthread_mutex_destroy(tp->count_mtx); + pthread_mutex_destroy(tp->access_mtx); pthread_cond_destroy(tp->done_signal); #ifndef NLOG_TRACE log_trace("Destroyed threadpool"); #endif /* ifndef NLOG_TRACE */ return 0; + +abandon: + log_error("Threadpool shutdown timed out with %lu threads still alive", tp->num_alive); + tp->is_abandoned = true; + for (u64 t = 0; t < tp->num_threads; t++) pthread_detach(tp->threads[t]); + pthread_mutex_unlock(tp->access_mtx); + + return ETIMEDOUT; } i32 tp_push_job(thread_pool_t *tp, tp_job_proc job, void *args) { @@ -171,11 +207,11 @@ void tp_wait(thread_pool_t *tp) { log_trace("Threadpool waiting for %lu jobs", tp->num_active + tp->jobs.queue->occupied); #endif /* ifndef NLOG_TRACE */ - pthread_mutex_lock(tp->count_mtx); + pthread_mutex_lock(tp->access_mtx); while (tp->jobs.queue->occupied || tp->num_active) { - pthread_cond_wait(tp->done_signal, tp->count_mtx); + pthread_cond_wait(tp->done_signal, tp->access_mtx); } - pthread_mutex_unlock(tp->count_mtx); + pthread_mutex_unlock(tp->access_mtx); #ifndef NLOG_TRACE log_trace("Threadpool finished waiting"); diff --git a/src/threadpool.h b/src/threadpool.h index b009008..190f971 100644 --- a/src/threadpool.h +++ b/src/threadpool.h @@ -14,13 +14,15 @@ typedef struct { // Threads check this before waiting for a job b32 is_running; // Write mutex for num_active/alive -pthread_mutex_t * count_mtx; +pthread_mutex_t * access_mtx; // Signal to tp_wait that all jobs are complete pthread_cond_t *done_signal; // Threads currently doing a job u64 num_active; // Threads who are alive u64 num_alive; +// Forceful teardown status + b32 is_abandoned; // Job queue ts_deque_t jobs; } thread_pool_t; @@ -30,7 +32,12 @@ thread_pool_t *tp_create(smrt_arena_t * arena , u64 max_jobs , u64 num_threads); -i32 tp_destroy(thread_pool_t *tp); +// Pass INT32_MAX for infinite timeout +// May return ETIMEDOUT, in which case, +// *tp should not be deallocaed until +// the program exits. +i32 tp_destroy(thread_pool_t *tp , + i32 timeout ); i32 tp_push_job(thread_pool_t * tp , tp_job_proc job , |