summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorsteven-na <noreply.github@stvnc.dev>2026-08-15 21:12:11 -0700
committersteven-na <noreply.github@stvnc.dev>2026-08-15 21:12:11 -0700
commit82e164f8796d57aa78c14b20f0a750f4b40c861b (patch)
tree1eb0df74a98ddc3fa872084b02ba309b3b0fb314
parenta4e71241ffbcc8576cd42856db061bb366cec3e5 (diff)
Threadpool timeout
-rw-r--r--Justfile.old99
-rw-r--r--src/threadpool.c86
-rw-r--r--src/threadpool.h11
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 ,