From 5ac070e0eb285468c46891381ea84c7eca2dc883 Mon Sep 17 00:00:00 2001 From: GauthierMalfilatre Date: Thu, 30 Jul 2026 17:18:02 +0200 Subject: [PATCH 1/5] [WIP] Threadpool and task queue in progress --- README.md | 4 +- include/kronkflow/macros/optimization.h | 16 ++-- include/kronkflow/macros/types.h | 12 +-- include/kronkflow/scheduler.h | 10 +-- include/kronkflow/task.h | 10 +-- src/scheduler/clear/scheduler_destroy.c | 2 +- src/scheduler/init/scheduler_create.c | 2 +- src/scheduler/scheduler.h | 2 +- src/scheduler/scheduler_addTask.c | 6 +- src/scheduler/scheduler_currentTick.c | 4 +- src/scheduler/scheduler_update.c | 2 +- src/task/task.h | 8 +- src/task/task_create.c | 6 +- src/task/task_opt.c | 6 +- src/threadpool/queue/queue.h | 39 +++++++++ src/threadpool/queue/queue_create.c | 29 +++++++ src/threadpool/queue/queue_destroy.c | 33 ++++++++ src/threadpool/queue/queue_empty.c | 13 +++ src/threadpool/queue/queue_front.c | 15 ++++ src/threadpool/queue/queue_pop.c | 26 ++++++ src/threadpool/queue/queue_push.c | 31 +++++++ src/threadpool/queue/queue_size.c | 15 ++++ src/threadpool/threadpool.h | 48 +++++++++++ src/threadpool/threadpool_create.c | 108 ++++++++++++++++++++++++ 24 files changed, 402 insertions(+), 45 deletions(-) create mode 100644 src/threadpool/queue/queue.h create mode 100644 src/threadpool/queue/queue_create.c create mode 100644 src/threadpool/queue/queue_destroy.c create mode 100644 src/threadpool/queue/queue_empty.c create mode 100644 src/threadpool/queue/queue_front.c create mode 100644 src/threadpool/queue/queue_pop.c create mode 100644 src/threadpool/queue/queue_push.c create mode 100644 src/threadpool/queue/queue_size.c create mode 100644 src/threadpool/threadpool.h create mode 100644 src/threadpool/threadpool_create.c diff --git a/README.md b/README.md index b6373f7..7a50ead 100644 --- a/README.md +++ b/README.md @@ -70,8 +70,8 @@ int main(void) - `kfScheduler_tick(kfScheduler *sch, void *context)`: Advance the scheduler by one tick and execute ready tasks. ### Task Management -- `kfTask_opt(prHandler handler, void *data, prClearer clearer)`: Build a task opt structure. -- `kfScheduler_addTask(kfScheduler *sch, kfTask task, prTick delay, prTick interval)`: Register a task in the scheduler. +- `kfTask_opt(kfHandler handler, void *data, kfClearer clearer)`: Build a task opt structure. +- `kfScheduler_addTask(kfScheduler *sch, kfTask task, kfTick delay, kfTick interval)`: Register a task in the scheduler. ## License diff --git a/include/kronkflow/macros/optimization.h b/include/kronkflow/macros/optimization.h index b166795..ecbe98c 100644 --- a/include/kronkflow/macros/optimization.h +++ b/include/kronkflow/macros/optimization.h @@ -8,27 +8,27 @@ #define PROPHECY_MACROS_OPTIMIZATION_H #if defined(__GNUC__) && (__GNUC__ >= 4) || defined(__has_attribute) && __has_attribute(visibility) - #define PR_API __attribute__((visibility("default"))) + #define KF_API __attribute__((visibility("default"))) #else - #define PR_API + #define KF_API #endif #if defined(__GNUC__) && (__GNUC__ >= 4) || defined(__has_attribute) && __has_attribute(unused) - #define PR_UNUSED __attribute__((unused)) + #define KF_UNUSED __attribute__((unused)) #else - #define PR_UNUSED + #define KF_UNUSED #endif #if defined(__GNUC__) && (__GNUC__ >= 4) || defined(__has_attribute) && __has_attribute(hot) - #define PR_HOT __attribute__((hot)) + #define KF_HOT __attribute__((hot)) #else - #define PR_HOT + #define KF_HOT #endif #if defined(__GNUC__) && (__GNUC__ >= 4) || defined(__has_attribute) && __has_attribute(cold) - #define PR_COLD __attribute__((cold)) + #define KF_COLD __attribute__((cold)) #else - #define PR_COLD + #define KF_COLD #endif #endif /* PROPHECY_MACROS_OPTIMIZATION_H */ diff --git a/include/kronkflow/macros/types.h b/include/kronkflow/macros/types.h index 8c4edd7..644f45a 100644 --- a/include/kronkflow/macros/types.h +++ b/include/kronkflow/macros/types.h @@ -10,13 +10,13 @@ #include #include - typedef int prBool; - #define knTrue 1 - #define knFalse 0 + typedef int kfBool; + #define kfTrue 1 + #define kfFalse 0 - typedef prBool (*prHandler)(void *, void *); - typedef void (*prClearer)(void *); - typedef uint64_t prTick; + typedef kfBool (*kfHandler)(void *, void *); + typedef void (*kfClearer)(void *); + typedef uint64_t kfTick; typedef size_t kfTaskID; #endif /* PROPHECY_MACROS_TYPES_H */ diff --git a/include/kronkflow/scheduler.h b/include/kronkflow/scheduler.h index 8867eb5..8771611 100644 --- a/include/kronkflow/scheduler.h +++ b/include/kronkflow/scheduler.h @@ -35,7 +35,7 @@ typedef struct prophecy_scheduler_s kfScheduler; * @return Returns newly allocated kfScheduler, or NULL on error */ /////////////////////////////////////////////////////////////////////////////// -PR_API kfScheduler *kfScheduler_create(size_t size); +KF_API kfScheduler *kfScheduler_create(size_t size); /////////////////////////////////////////////////////////////////////////////// @@ -46,7 +46,7 @@ PR_API kfScheduler *kfScheduler_create(size_t size); * @param sch The scheduler to destroy */ /////////////////////////////////////////////////////////////////////////////// -PR_API void kfScheduler_destroy(kfScheduler *sch); +KF_API void kfScheduler_destroy(kfScheduler *sch); /////////////////////////////////////////////////////////////////////////////// @@ -61,7 +61,7 @@ PR_API void kfScheduler_destroy(kfScheduler *sch); * @return The id of the task added (>1), 0 if failed */ /////////////////////////////////////////////////////////////////////////////// -PR_API kfTaskID kfScheduler_addTask(kfScheduler *sch, kfTaskOpt opt, prTick delay, prTick interval); +KF_API kfTaskID kfScheduler_addTask(kfScheduler *sch, kfTaskOpt opt, kfTick delay, kfTick interval); /////////////////////////////////////////////////////////////////////////////// @@ -74,7 +74,7 @@ PR_API kfTaskID kfScheduler_addTask(kfScheduler *sch, kfTaskOpt opt, prTick dela * @return Returns the number of tasks executed this tick */ /////////////////////////////////////////////////////////////////////////////// -PR_API size_t kfScheduler_tick(kfScheduler *sch, void *context); +KF_API size_t kfScheduler_tick(kfScheduler *sch, void *context); /////////////////////////////////////////////////////////////////////////////// @@ -86,7 +86,7 @@ PR_API size_t kfScheduler_tick(kfScheduler *sch, void *context); * @return The current tick */ /////////////////////////////////////////////////////////////////////////////// -PR_API prTick kfScheduler_currentTick(const kfScheduler *sch); +KF_API kfTick kfScheduler_currentTick(const kfScheduler *sch); /////////////////////////////////////////////////////////////////////////////// #endif /* PROPHECY_SCHEDULER_H */ diff --git a/include/kronkflow/task.h b/include/kronkflow/task.h index 303e137..2016361 100644 --- a/include/kronkflow/task.h +++ b/include/kronkflow/task.h @@ -27,9 +27,9 @@ typedef struct prophecy_task_s kfTask; /////////////////////////////////////////////////////////////////////////////// typedef struct prophecy_task_opt_s { - prHandler handler; //!< The task handler (function ptr) + kfHandler handler; //!< The task handler (function ptr) void* data; //!< The data to give to the handler - prClearer clearer; //!< The clearer of the data if allocated + kfClearer clearer; //!< The clearer of the data if allocated } kfTaskOpt; /////////////////////////////////////////////////////////////////////////////// @@ -39,12 +39,12 @@ typedef struct prophecy_task_opt_s { /** * @brief Create task opts * - * @param handler The function ptr handler prBool (*)(void *, void *) + * @param handler The function ptr handler kfBool (*)(void *, void *) * @param data The data to give to the handler * @param clearer The clearer to clear data if allocated */ /////////////////////////////////////////////////////////////////////////////// -PR_API kfTaskOpt kfTask_opt(prHandler handler, void *data, prClearer clearer); +KF_API kfTaskOpt kfTask_opt(kfHandler handler, void *data, kfClearer clearer); /////////////////////////////////////////////////////////////////////////////// @@ -58,7 +58,7 @@ PR_API kfTaskOpt kfTask_opt(prHandler handler, void *data, prClearer clearer); * @return Returns the new task */ /////////////////////////////////////////////////////////////////////////////// -PR_API kfTask kfTask_create(kfTaskOpt *data, prTick delay, prTick interval); +KF_API kfTask kfTask_create(kfTaskOpt *data, kfTick delay, kfTick interval); /////////////////////////////////////////////////////////////////////////////// #endif /* PROPHECY_TASKS_H */ diff --git a/src/scheduler/clear/scheduler_destroy.c b/src/scheduler/clear/scheduler_destroy.c index 78a56bf..e147ad5 100644 --- a/src/scheduler/clear/scheduler_destroy.c +++ b/src/scheduler/clear/scheduler_destroy.c @@ -8,7 +8,7 @@ #include "kronkflow/macros/optimization.h" #include -PR_API +KF_API void kfScheduler_destroy( kfScheduler *sch ) diff --git a/src/scheduler/init/scheduler_create.c b/src/scheduler/init/scheduler_create.c index 844d7c0..07e5fdf 100644 --- a/src/scheduler/init/scheduler_create.c +++ b/src/scheduler/init/scheduler_create.c @@ -9,7 +9,7 @@ #include #include -PR_API +KF_API kfScheduler *kfScheduler_create( size_t size ) diff --git a/src/scheduler/scheduler.h b/src/scheduler/scheduler.h index f1213e9..2214590 100644 --- a/src/scheduler/scheduler.h +++ b/src/scheduler/scheduler.h @@ -21,7 +21,7 @@ typedef struct prophecy_scheduler_s { kfTask *tasks; //!< The raw array of tasks size_t size; //!< The size of tasks raw array size_t count; //!< The number of tasks pushed - prTick tick; //!< The current tick (please tick scheduler at each loop) + kfTick tick; //!< The current tick (please tick scheduler at each loop) } kfScheduler; /////////////////////////////////////////////////////////////////////////////// diff --git a/src/scheduler/scheduler_addTask.c b/src/scheduler/scheduler_addTask.c index 2818075..de548fc 100644 --- a/src/scheduler/scheduler_addTask.c +++ b/src/scheduler/scheduler_addTask.c @@ -28,12 +28,12 @@ static int __kfScheduler_ensureCapacity( return 0; } -PR_API +KF_API size_t kfScheduler_addTask( kfScheduler *sch, kfTaskOpt taskOptions, - prTick target, - prTick interval + kfTick target, + kfTick interval ) { static size_t _id = 1; diff --git a/src/scheduler/scheduler_currentTick.c b/src/scheduler/scheduler_currentTick.c index 04834e2..325a009 100644 --- a/src/scheduler/scheduler_currentTick.c +++ b/src/scheduler/scheduler_currentTick.c @@ -8,8 +8,8 @@ #include "kronkflow/scheduler.h" #include "scheduler.h" -PR_API -prTick kfScheduler_currentTick( +KF_API +kfTick kfScheduler_currentTick( const kfScheduler *sch ) { diff --git a/src/scheduler/scheduler_update.c b/src/scheduler/scheduler_update.c index b8254a8..a84df3c 100644 --- a/src/scheduler/scheduler_update.c +++ b/src/scheduler/scheduler_update.c @@ -12,7 +12,7 @@ #include "../minheap/minheap.h" #include "kronkflow/scheduler.h" -PR_API +KF_API size_t kfScheduler_tick( kfScheduler *sch, void *context diff --git a/src/task/task.h b/src/task/task.h index aeda37e..8000760 100644 --- a/src/task/task.h +++ b/src/task/task.h @@ -22,11 +22,11 @@ typedef struct prophecy_task_s { size_t id; //!< The id of the tasks - prHandler handler; //!< The handler (callback) + kfHandler handler; //!< The handler (callback) void* data; //!< The task data - prClearer clearer; //!< The data clearer - prTick interval; //!< The interval (0 if ponctual, > 0 else) - prTick target; //!< The tick remainings. + kfClearer clearer; //!< The data clearer + kfTick interval; //!< The interval (0 if ponctual, > 0 else) + kfTick target; //!< The tick remainings. } kfTask; /////////////////////////////////////////////////////////////////////////////// diff --git a/src/task/task_create.c b/src/task/task_create.c index d02418f..1333df5 100644 --- a/src/task/task_create.c +++ b/src/task/task_create.c @@ -10,12 +10,12 @@ #include "task.h" // DEBUG: Make sure opt is not NULL -PR_API +KF_API inline kfTask kfTask_create( kfTaskOpt *opt, - prTick delay, - prTick interval + kfTick delay, + kfTick interval ) { return (kfTask){ diff --git a/src/task/task_opt.c b/src/task/task_opt.c index 51c271d..7e97e8f 100644 --- a/src/task/task_opt.c +++ b/src/task/task_opt.c @@ -8,12 +8,12 @@ #include "kronkflow/macros/types.h" #include "kronkflow/task.h" -PR_API +KF_API inline kfTaskOpt kfTask_opt( - prHandler handler, + kfHandler handler, void *data, - prClearer clearer + kfClearer clearer ) { return (kfTaskOpt){ diff --git a/src/threadpool/queue/queue.h b/src/threadpool/queue/queue.h new file mode 100644 index 0000000..21f5610 --- /dev/null +++ b/src/threadpool/queue/queue.h @@ -0,0 +1,39 @@ +/* +** EPITECH PROJECT, 2026 +** KRONKFLOW +** File description: +** Generic linked list as a queue +*/ +#ifndef KRONKFLOW_QUEUE_H + #define KRONKFLOW_QUEUE_H + #include + #include + +typedef struct queue_node_s { + + void *data; + struct queue_node_s *next; + +} queue_node_t; + +typedef struct queue_s { + + queue_node_t *head; + queue_node_t *tail; + size_t size; + +} queue_t; + +typedef struct queue_s kfQueue; + +queue_t *queue_create(void); +int queue_init(kfQueue *q); +bool queue_push(queue_t *q, void *data); +void *queue_pop(queue_t *q); +void *queue_front(const queue_t *q); +bool queue_empty(const queue_t *q); +size_t queue_size(const queue_t *q); +void queue_clear(queue_t *q, void (*free_func)(void *)); +void queue_destroy(queue_t *q, void (*free_func)(void *)); + +#endif /* GUL_QUEUE_H */ diff --git a/src/threadpool/queue/queue_create.c b/src/threadpool/queue/queue_create.c new file mode 100644 index 0000000..38c6315 --- /dev/null +++ b/src/threadpool/queue/queue_create.c @@ -0,0 +1,29 @@ +/* +** EPITECH PROJECT, 2026 +** PANORAMIX +** File description: +** Create queue +*/ +#include "queue.h" +#include + +queue_t *queue_create(void) +{ + queue_t *q = calloc(1, sizeof(queue_t)); + + if (queue_init(q) == -1) { + free(q); + return NULL; + } + return q; +} + +int queue_init(kfQueue *q) +{ + if (!q) + return -1; + q->head = NULL; + q->tail = NULL; + q->size = 0; + return 0; +} diff --git a/src/threadpool/queue/queue_destroy.c b/src/threadpool/queue/queue_destroy.c new file mode 100644 index 0000000..ef62bf4 --- /dev/null +++ b/src/threadpool/queue/queue_destroy.c @@ -0,0 +1,33 @@ +/* +** EPITECH PROJECT, 2026 +** PANORAMIX +** File description: +** Destroy queue +*/ +#include "queue.h" +#include + +void queue_destroy(queue_t *q, void (*free_func)(void *)) +{ + if (!q) + return; + queue_clear(q, free_func); + free(q); +} + +void queue_clear( + queue_t *q, + void (*free_func)(void *) +) +{ + void *data; + + if (!q) + return; + while (!queue_empty(q)) { + data = queue_pop(q); + if (free_func && data) { + free_func(data); + } + } +} diff --git a/src/threadpool/queue/queue_empty.c b/src/threadpool/queue/queue_empty.c new file mode 100644 index 0000000..0099b1a --- /dev/null +++ b/src/threadpool/queue/queue_empty.c @@ -0,0 +1,13 @@ +/* +** EPITECH PROJECT, 2026 +** PANORAMIX +** File description: +** Queue is empty +*/ +#include "queue.h" +#include + +bool queue_empty(const queue_t *q) +{ + return (q == NULL || q->size == 0); +} diff --git a/src/threadpool/queue/queue_front.c b/src/threadpool/queue/queue_front.c new file mode 100644 index 0000000..34799f7 --- /dev/null +++ b/src/threadpool/queue/queue_front.c @@ -0,0 +1,15 @@ +/* +** EPITECH PROJECT, 2026 +** PANORAMIX +** File description: +** Queue front +*/ +#include "queue.h" +#include + +void *queue_front(const queue_t *q) +{ + if (!q || q->size == 0 || !q->head) + return NULL; + return q->head->data; +} diff --git a/src/threadpool/queue/queue_pop.c b/src/threadpool/queue/queue_pop.c new file mode 100644 index 0000000..b9711f8 --- /dev/null +++ b/src/threadpool/queue/queue_pop.c @@ -0,0 +1,26 @@ +/* +** EPITECH PROJECT, 2026 +** PANORAMIX +** File description: +** Queue pop +*/ +#include "queue.h" +#include +#include + +void *queue_pop(queue_t *q) +{ + queue_node_t *temp; + void *data; + + if (queue_empty(q)) + return NULL; + temp = q->head; + data = temp->data; + q->head = q->head->next; + if (q->head == NULL) + q->tail = NULL; + free(temp); + q->size--; + return data; +} diff --git a/src/threadpool/queue/queue_push.c b/src/threadpool/queue/queue_push.c new file mode 100644 index 0000000..f5631ce --- /dev/null +++ b/src/threadpool/queue/queue_push.c @@ -0,0 +1,31 @@ +/* +** EPITECH PROJECT, 2026 +** PANORAMIX +** File description: +** Queue push +*/ +#include "queue.h" +#include +#include + +bool queue_push(queue_t *q, void *data) +{ + queue_node_t *new_node; + + if (!q) + return false; + new_node = malloc(sizeof(queue_node_t)); + if (!new_node) + return false; + new_node->data = data; + new_node->next = NULL; + if (queue_empty(q)) { + q->head = new_node; + q->tail = new_node; + } else { + q->tail->next = new_node; + q->tail = new_node; + } + q->size++; + return true; +} diff --git a/src/threadpool/queue/queue_size.c b/src/threadpool/queue/queue_size.c new file mode 100644 index 0000000..6ff1aed --- /dev/null +++ b/src/threadpool/queue/queue_size.c @@ -0,0 +1,15 @@ +/* +** EPITECH PROJECT, 2026 +** PANORAMIX +** File description: +** Queue is empty +*/ +#include "queue.h" +#include + +size_t queue_size(const queue_t *q) +{ + if (!q->size) + return 0; + return q->size; +} diff --git a/src/threadpool/threadpool.h b/src/threadpool/threadpool.h new file mode 100644 index 0000000..1e3291d --- /dev/null +++ b/src/threadpool/threadpool.h @@ -0,0 +1,48 @@ +/* +** FREE PROJECT, 2026 +** Kronkflow +** File description: +** threadpool +*/ +#ifndef KRONKFLOW_THREADPOOL_PRIVATE_H + #define KRONKFLOW_THREADPOOL_PRIVATE_H + #include + #include + #include + #include +#include "kronkflow/macros/types.h" + #include "queue/queue.h" + +/////////////////////////////////////////////////////////////////////////////// +/** + * @struct kronkflow_threadpool_s + * + * @brief This struct is a dedicated threadpool for kronkflow multithreading features + */ +/////////////////////////////////////////////////////////////////////////////// +typedef struct kronkflow_threadpool_s { + + atomic_size_t pendings; //!< The number of task pendings + atomic_size_t runnings; //!< The number of thread up and runnings + atomic_size_t workers; //!< The number of threads + pthread_t* threads; //!< The threads (array) + pthread_cond_t cond; //!< The conditionnal variable + pthread_mutex_t mutex; //!< Mutex + kfQueue queue; //!< Queue for tasks + kfBool stop; //!< Does the pool should stop + void* ctx; //!< The ctx to give... + +} kfThreadPool; +/////////////////////////////////////////////////////////////////////////////// + +// TODO: Documentation +kfThreadPool *kfThreadPool_create(ssize_t nthreads, void *ctx); +int kfThreadPool_init(kfThreadPool *pool, size_t nthreads, void *ctx); +void kfThreadPool_destroy(kfThreadPool *pool); +void kfThreadPool_clear(kfThreadPool *pool); +void kfThreadPool_stop(kfThreadPool *pool); +size_t kfThreadPool_running(const kfThreadPool *pool); +size_t kfThreadPool_remaining(const kfThreadPool *pool); +int kfThreadPool_pushTask(kfThreadPool *pool, kfHandler task, void *data); + +#endif /* KRONKFLOW_THREADPOOL_PRIVATE_H */ diff --git a/src/threadpool/threadpool_create.c b/src/threadpool/threadpool_create.c new file mode 100644 index 0000000..423bab0 --- /dev/null +++ b/src/threadpool/threadpool_create.c @@ -0,0 +1,108 @@ +/* +** FREE PROJECT, 2026 +** Kronkflow +** File description: +** threadpool +*/ +#include "kronkflow/macros/optimization.h" +#include "kronkflow/macros/types.h" +#include "queue/queue.h" +#include "threadpool.h" +#include "kronkflow/threadpool.h" +#include +#include +#include +#include +#include +#include +#include + +kfThreadPool *kfThreadPool_create( + ssize_t nthreads, + void *ctx +) +{ + kfThreadPool *pool = NULL; + + if (nthreads == 0) { + return NULL; + } else if (nthreads == -1) { + nthreads = sysconf(_SC_NPROCESSORS_ONLN); + if (nthreads <= 0) { + return NULL; + } + } + pool = calloc(1, sizeof(kfThreadPool)); + if (!pool) { + return NULL; + } + // NOTE: Nthread is already good, so init won't check... + if (kfThreadPool_init(pool, nthreads, ctx) == -1) { + free(pool); + return NULL; + } + return pool; +} + +static void *__routine( + void *arg +) +{ + kfThreadPool *pool = (kfThreadPool *)arg; + + if (!pool) { + return NULL; + } + while (true) { + kfHandler handler; + void *d; + pthread_mutex_lock(&pool->mutex); + while (!pool->stop && queue_empty(&pool->queue)) { + pthread_cond_wait(&pool->cond, &pool->mutex); + } + // FIXME: Queue empty really necessary ?? + if (pool->stop && queue_empty(&pool->queue)) { + return NULL; + } + d = queue_front(&pool->queue); + handler = (void *)((long int)d >> 8); + queue_pop(&pool->queue); + pool->pendings--; + pthread_mutex_unlock(&pool->mutex); + pool->runnings++; + // NOTE: Should call handler... with ctx + handler(pool->ctx, (void *)((long int)d & 0xff)); + pool->runnings--; + } +} + +int kfThreadPool_init( + kfThreadPool *pool, + size_t nthreads, + void *ctx +) +{ + if (!pool) { + return -1; + } + if (pthread_mutex_init(&pool->mutex, NULL) != 0) { + return -1; + } + if (pthread_cond_init(&pool->cond, NULL) != 0) { + return -1; + } + pool->pendings = 0; + pool->runnings = 0; + pool->ctx = ctx; + pool->workers = nthreads; + pool->stop = kfFalse; + pool->threads = calloc(nthreads, sizeof(pthread_t)); + if (!pool->threads) { + return -1; + } + for (size_t i = 0; i < nthreads; ++i) { + pthread_create(&pool->threads[i], NULL, &__routine, pool); + } + queue_init(&pool->queue); + return 0; +} From e0976e953de9bda75a3c562612f9bd382cf57fe2 Mon Sep 17 00:00:00 2001 From: GauthierMalfilatre Date: Fri, 31 Jul 2026 17:21:32 +0200 Subject: [PATCH 2/5] [WIP] Threadpool... --- src/threadpool/threadpool.h | 7 ++++++ src/threadpool/threadpool_create.c | 12 ++++------ src/threadpool/threadpool_destroy.c | 34 +++++++++++++++++++++++++++ src/threadpool/threadpool_pushTask.c | 35 ++++++++++++++++++++++++++++ 4 files changed, 81 insertions(+), 7 deletions(-) create mode 100644 src/threadpool/threadpool_destroy.c create mode 100644 src/threadpool/threadpool_pushTask.c diff --git a/src/threadpool/threadpool.h b/src/threadpool/threadpool.h index 1e3291d..df9ee13 100644 --- a/src/threadpool/threadpool.h +++ b/src/threadpool/threadpool.h @@ -35,6 +35,13 @@ typedef struct kronkflow_threadpool_s { } kfThreadPool; /////////////////////////////////////////////////////////////////////////////// +typedef struct kronkflow_thread_task_s { + + kfHandler handler; + void *data; + +} kfThreadTask; + // TODO: Documentation kfThreadPool *kfThreadPool_create(ssize_t nthreads, void *ctx); int kfThreadPool_init(kfThreadPool *pool, size_t nthreads, void *ctx); diff --git a/src/threadpool/threadpool_create.c b/src/threadpool/threadpool_create.c index 423bab0..e6c943a 100644 --- a/src/threadpool/threadpool_create.c +++ b/src/threadpool/threadpool_create.c @@ -8,7 +8,6 @@ #include "kronkflow/macros/types.h" #include "queue/queue.h" #include "threadpool.h" -#include "kronkflow/threadpool.h" #include #include #include @@ -54,8 +53,8 @@ static void *__routine( return NULL; } while (true) { - kfHandler handler; - void *d; + kfThreadTask *task; + // void *d; pthread_mutex_lock(&pool->mutex); while (!pool->stop && queue_empty(&pool->queue)) { pthread_cond_wait(&pool->cond, &pool->mutex); @@ -64,14 +63,13 @@ static void *__routine( if (pool->stop && queue_empty(&pool->queue)) { return NULL; } - d = queue_front(&pool->queue); - handler = (void *)((long int)d >> 8); + task = queue_front(&pool->queue); queue_pop(&pool->queue); pool->pendings--; pthread_mutex_unlock(&pool->mutex); pool->runnings++; - // NOTE: Should call handler... with ctx - handler(pool->ctx, (void *)((long int)d & 0xff)); + // NOTE: Should call task... with ctx + task->handler(pool->ctx, task->data); pool->runnings--; } } diff --git a/src/threadpool/threadpool_destroy.c b/src/threadpool/threadpool_destroy.c new file mode 100644 index 0000000..7539ad6 --- /dev/null +++ b/src/threadpool/threadpool_destroy.c @@ -0,0 +1,34 @@ +/* +** FREE PROJECT, 2026 +** Kronkflow +** File description: +** threadpool destroy +*/ +#include "queue/queue.h" +#include "threadpool.h" +#include +#include +#include + +void kfThreadPool_destroy( + kfThreadPool *pool +) +{ + if (!pool) + return; + kfThreadPool_clear(pool); + free(pool); +} + +void kfThreadPool_clear( + kfThreadPool *pool +) +{ + pthread_mutex_destroy(&pool->mutex); + pthread_cond_destroy(&pool->cond); + queue_destroy(&pool->queue, NULL); + for (size_t i = 0; i < pool->workers; ++i) { + pthread_join(pool->threads[i], NULL); + } + free(pool->threads); +} diff --git a/src/threadpool/threadpool_pushTask.c b/src/threadpool/threadpool_pushTask.c new file mode 100644 index 0000000..ef4da67 --- /dev/null +++ b/src/threadpool/threadpool_pushTask.c @@ -0,0 +1,35 @@ +/* +** FREE PROJECT, 2026 +** Kronkflow +** File description: +** threadpool +*/ +#include "queue/queue.h" +#include "threadpool.h" +#include +#include + +int kfThreadPool_pushTask( + kfThreadPool *pool, + kfHandler task, + void *data +) +{ + kfThreadTask *p = NULL; + + if (!pool) { + return -1; + } + p = calloc(1, sizeof(kfThreadTask)); + if (!p) { + return -1; + } + p->data = data; + p->handler = task; + pthread_mutex_lock(&pool->mutex); + queue_push(&pool->queue, p); + ++pool->pendings; + pthread_cond_signal(&pool->cond); + pthread_mutex_unlock(&pool->mutex); + return 0; +} From f44dc961ed88d21eee712ef229302db3a244c711 Mon Sep 17 00:00:00 2001 From: GauthierMalfilatre Date: Tue, 11 Aug 2026 14:58:58 +0200 Subject: [PATCH 3/5] [WIP] Threadpool tests and functionnalities --- CMakeLists.txt | 4 ++++ cmake/CompilerOptions.cmake | 19 +++++++++++++++++ cmake/ProjectOptions.cmake | 16 ++++++++++++++ cmake/Utils.cmake | 3 +++ include/kronkflow/utils/threadpool.h | 10 +++++++++ src/threadpool/test.c | 32 ++++++++++++++++++++++++++++ src/threadpool/threadpool.h | 10 +++++---- src/threadpool/threadpool_create.c | 4 ++-- src/threadpool/threadpool_destroy.c | 10 ++++++--- src/threadpool/threadpool_pushTask.c | 2 +- 10 files changed, 100 insertions(+), 10 deletions(-) create mode 100644 cmake/CompilerOptions.cmake create mode 100644 cmake/ProjectOptions.cmake create mode 100644 cmake/Utils.cmake create mode 100644 include/kronkflow/utils/threadpool.h create mode 100644 src/threadpool/test.c diff --git a/CMakeLists.txt b/CMakeLists.txt index 70d93ee..243779c 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -1,6 +1,10 @@ cmake_minimum_required(VERSION 3.15) project(kronkflow LANGUAGES C) +include(cmake/ProjectOptions.cmake) +include(cmake/CompilerOptions.cmake) +include(cmake/Utils.cmake) + if(CMAKE_CURRENT_SOURCE_DIR STREQUAL CMAKE_CURRENT_BINARY_DIR) message(FATAL_ERROR "Les builds dans le dossier source sont interdits. Créez un dossier de build (ex: 'build/').") endif() diff --git a/cmake/CompilerOptions.cmake b/cmake/CompilerOptions.cmake new file mode 100644 index 0000000..46c1cdc --- /dev/null +++ b/cmake/CompilerOptions.cmake @@ -0,0 +1,19 @@ +add_library(project_options INTERFACE) + +target_compile_options(project_options INTERFACE + ${PROJECT_WARNINGS} + $<$:${PROJECT_DEBUG_FLAGS}> + $<$:${PROJECT_RELEASE_FLAGS}> +) + +target_link_libraries(project_options INTERFACE + stdc++exp +) + +target_link_options(project_options INTERFACE + -rdynamic +) + +target_compile_definitions(project_options INTERFACE + $<$:_DEBUG> +) diff --git a/cmake/ProjectOptions.cmake b/cmake/ProjectOptions.cmake new file mode 100644 index 0000000..e5d9eed --- /dev/null +++ b/cmake/ProjectOptions.cmake @@ -0,0 +1,16 @@ +# Global settings +set(CMAKE_CXX_STANDARD 23) +set(CMAKE_CXX_STANDARD_REQUIRED ON) +set(CMAKE_EXPORT_COMPILE_COMMANDS ON) + +# Warnings +set(PROJECT_WARNINGS + -Wall -Wextra -Wpedantic + -Wshadow -Wnull-dereference + -Wcast-align -Wmissing-declarations + -Wundef -Wunreachable-code +) + +# Build flags +set(PROJECT_DEBUG_FLAGS -g -O0) +set(PROJECT_RELEASE_FLAGS -O3) diff --git a/cmake/Utils.cmake b/cmake/Utils.cmake new file mode 100644 index 0000000..9303062 --- /dev/null +++ b/cmake/Utils.cmake @@ -0,0 +1,3 @@ +function(apply_project_settings target) + target_link_libraries(${target} PRIVATE project_options) +endfunction() diff --git a/include/kronkflow/utils/threadpool.h b/include/kronkflow/utils/threadpool.h new file mode 100644 index 0000000..49cbe28 --- /dev/null +++ b/include/kronkflow/utils/threadpool.h @@ -0,0 +1,10 @@ +// Give access to only one function that test the threadpool + +#ifndef KRONKFLOW_THREADPOOL_TEST + #define KRONKFLOW_THREADPOOL_TEST + + #include "kronkflow/macros/optimization.h" + +KF_API void kfThreadPool_test(void); + +#endif /* KRONKFLOW_THREADPOOL_TEST */ diff --git a/src/threadpool/test.c b/src/threadpool/test.c new file mode 100644 index 0000000..c394c61 --- /dev/null +++ b/src/threadpool/test.c @@ -0,0 +1,32 @@ +#include "kronkflow/macros/types.h" +#include "kronkflow/utils/threadpool.h" +#include +#include +#include +#include +#include "threadpool.h" + +static void *__test( + void *arg +) +{ + printf("%ld -> %d\n", (long int)arg, rand()); + return NULL; +} + +[[gnu::constructor]] +static void seedrand(void) +{ + srand(time(NULL)); +} + +KF_API +void kfThreadPool_test(void) +{ + kfThreadPool *pool = kfThreadPool_create(8, NULL); + + kfThreadPool_pushTask(pool, __test, (void *)10); + // kfThreadPool_join(pool); + printf("Test threadpool\n"); + kfThreadPool_destroy(pool); +} diff --git a/src/threadpool/threadpool.h b/src/threadpool/threadpool.h index df9ee13..59aa906 100644 --- a/src/threadpool/threadpool.h +++ b/src/threadpool/threadpool.h @@ -10,9 +10,11 @@ #include #include #include -#include "kronkflow/macros/types.h" + #include "kronkflow/macros/types.h" #include "queue/queue.h" +typedef void *(*kfThreadPoolHandler)(void *); + /////////////////////////////////////////////////////////////////////////////// /** * @struct kronkflow_threadpool_s @@ -37,8 +39,8 @@ typedef struct kronkflow_threadpool_s { typedef struct kronkflow_thread_task_s { - kfHandler handler; - void *data; + kfThreadPoolHandler handler; + void* data; } kfThreadTask; @@ -50,6 +52,6 @@ void kfThreadPool_clear(kfThreadPool *pool); void kfThreadPool_stop(kfThreadPool *pool); size_t kfThreadPool_running(const kfThreadPool *pool); size_t kfThreadPool_remaining(const kfThreadPool *pool); -int kfThreadPool_pushTask(kfThreadPool *pool, kfHandler task, void *data); +int kfThreadPool_pushTask(kfThreadPool *pool, kfThreadPoolHandler task, void *data); #endif /* KRONKFLOW_THREADPOOL_PRIVATE_H */ diff --git a/src/threadpool/threadpool_create.c b/src/threadpool/threadpool_create.c index e6c943a..b049f82 100644 --- a/src/threadpool/threadpool_create.c +++ b/src/threadpool/threadpool_create.c @@ -69,7 +69,7 @@ static void *__routine( pthread_mutex_unlock(&pool->mutex); pool->runnings++; // NOTE: Should call task... with ctx - task->handler(pool->ctx, task->data); + task->handler(task->data); pool->runnings--; } } @@ -89,6 +89,7 @@ int kfThreadPool_init( if (pthread_cond_init(&pool->cond, NULL) != 0) { return -1; } + queue_init(&pool->queue); pool->pendings = 0; pool->runnings = 0; pool->ctx = ctx; @@ -101,6 +102,5 @@ int kfThreadPool_init( for (size_t i = 0; i < nthreads; ++i) { pthread_create(&pool->threads[i], NULL, &__routine, pool); } - queue_init(&pool->queue); return 0; } diff --git a/src/threadpool/threadpool_destroy.c b/src/threadpool/threadpool_destroy.c index 7539ad6..c12f66b 100644 --- a/src/threadpool/threadpool_destroy.c +++ b/src/threadpool/threadpool_destroy.c @@ -7,6 +7,7 @@ #include "queue/queue.h" #include "threadpool.h" #include +#include #include #include @@ -24,11 +25,14 @@ void kfThreadPool_clear( kfThreadPool *pool ) { - pthread_mutex_destroy(&pool->mutex); - pthread_cond_destroy(&pool->cond); - queue_destroy(&pool->queue, NULL); + // pthread_mutex_lock(&pool->mutex); + pool->stop = true; + // pthread_mutex_unlock(&pool->mutex); for (size_t i = 0; i < pool->workers; ++i) { pthread_join(pool->threads[i], NULL); } + pthread_mutex_destroy(&pool->mutex); + pthread_cond_destroy(&pool->cond); + queue_destroy(&pool->queue, NULL); free(pool->threads); } diff --git a/src/threadpool/threadpool_pushTask.c b/src/threadpool/threadpool_pushTask.c index ef4da67..a35330a 100644 --- a/src/threadpool/threadpool_pushTask.c +++ b/src/threadpool/threadpool_pushTask.c @@ -11,7 +11,7 @@ int kfThreadPool_pushTask( kfThreadPool *pool, - kfHandler task, + void *(*task)(void *), void *data ) { From d94ea480dc737d1c0e1d45d6d3fc72dac401f6b7 Mon Sep 17 00:00:00 2001 From: GauthierMalfilatre Date: Tue, 8 Sep 2026 00:03:46 +0200 Subject: [PATCH 4/5] [UP] Update gitignore --- .gitignore | 1 + 1 file changed, 1 insertion(+) diff --git a/.gitignore b/.gitignore index 1660085..826b27f 100644 --- a/.gitignore +++ b/.gitignore @@ -19,6 +19,7 @@ libargot.so # Build Release/ +build/ Debug/ # Test From a94afa26b102c42d5d6bde73b4520899f35c7914 Mon Sep 17 00:00:00 2001 From: GauthierMalfilatre Date: Tue, 8 Sep 2026 00:07:27 +0200 Subject: [PATCH 5/5] [DEL] Delete multithreading folder as it will be provide by kronkpool. --- src/threadpool/queue/queue.h | 39 ---------- src/threadpool/queue/queue_create.c | 29 -------- src/threadpool/queue/queue_destroy.c | 33 --------- src/threadpool/queue/queue_empty.c | 13 ---- src/threadpool/queue/queue_front.c | 15 ---- src/threadpool/queue/queue_pop.c | 26 ------- src/threadpool/queue/queue_push.c | 31 -------- src/threadpool/queue/queue_size.c | 15 ---- src/threadpool/test.c | 32 -------- src/threadpool/threadpool.h | 57 -------------- src/threadpool/threadpool_create.c | 106 --------------------------- src/threadpool/threadpool_destroy.c | 38 ---------- src/threadpool/threadpool_pushTask.c | 35 --------- 13 files changed, 469 deletions(-) delete mode 100644 src/threadpool/queue/queue.h delete mode 100644 src/threadpool/queue/queue_create.c delete mode 100644 src/threadpool/queue/queue_destroy.c delete mode 100644 src/threadpool/queue/queue_empty.c delete mode 100644 src/threadpool/queue/queue_front.c delete mode 100644 src/threadpool/queue/queue_pop.c delete mode 100644 src/threadpool/queue/queue_push.c delete mode 100644 src/threadpool/queue/queue_size.c delete mode 100644 src/threadpool/test.c delete mode 100644 src/threadpool/threadpool.h delete mode 100644 src/threadpool/threadpool_create.c delete mode 100644 src/threadpool/threadpool_destroy.c delete mode 100644 src/threadpool/threadpool_pushTask.c diff --git a/src/threadpool/queue/queue.h b/src/threadpool/queue/queue.h deleted file mode 100644 index 21f5610..0000000 --- a/src/threadpool/queue/queue.h +++ /dev/null @@ -1,39 +0,0 @@ -/* -** EPITECH PROJECT, 2026 -** KRONKFLOW -** File description: -** Generic linked list as a queue -*/ -#ifndef KRONKFLOW_QUEUE_H - #define KRONKFLOW_QUEUE_H - #include - #include - -typedef struct queue_node_s { - - void *data; - struct queue_node_s *next; - -} queue_node_t; - -typedef struct queue_s { - - queue_node_t *head; - queue_node_t *tail; - size_t size; - -} queue_t; - -typedef struct queue_s kfQueue; - -queue_t *queue_create(void); -int queue_init(kfQueue *q); -bool queue_push(queue_t *q, void *data); -void *queue_pop(queue_t *q); -void *queue_front(const queue_t *q); -bool queue_empty(const queue_t *q); -size_t queue_size(const queue_t *q); -void queue_clear(queue_t *q, void (*free_func)(void *)); -void queue_destroy(queue_t *q, void (*free_func)(void *)); - -#endif /* GUL_QUEUE_H */ diff --git a/src/threadpool/queue/queue_create.c b/src/threadpool/queue/queue_create.c deleted file mode 100644 index 38c6315..0000000 --- a/src/threadpool/queue/queue_create.c +++ /dev/null @@ -1,29 +0,0 @@ -/* -** EPITECH PROJECT, 2026 -** PANORAMIX -** File description: -** Create queue -*/ -#include "queue.h" -#include - -queue_t *queue_create(void) -{ - queue_t *q = calloc(1, sizeof(queue_t)); - - if (queue_init(q) == -1) { - free(q); - return NULL; - } - return q; -} - -int queue_init(kfQueue *q) -{ - if (!q) - return -1; - q->head = NULL; - q->tail = NULL; - q->size = 0; - return 0; -} diff --git a/src/threadpool/queue/queue_destroy.c b/src/threadpool/queue/queue_destroy.c deleted file mode 100644 index ef62bf4..0000000 --- a/src/threadpool/queue/queue_destroy.c +++ /dev/null @@ -1,33 +0,0 @@ -/* -** EPITECH PROJECT, 2026 -** PANORAMIX -** File description: -** Destroy queue -*/ -#include "queue.h" -#include - -void queue_destroy(queue_t *q, void (*free_func)(void *)) -{ - if (!q) - return; - queue_clear(q, free_func); - free(q); -} - -void queue_clear( - queue_t *q, - void (*free_func)(void *) -) -{ - void *data; - - if (!q) - return; - while (!queue_empty(q)) { - data = queue_pop(q); - if (free_func && data) { - free_func(data); - } - } -} diff --git a/src/threadpool/queue/queue_empty.c b/src/threadpool/queue/queue_empty.c deleted file mode 100644 index 0099b1a..0000000 --- a/src/threadpool/queue/queue_empty.c +++ /dev/null @@ -1,13 +0,0 @@ -/* -** EPITECH PROJECT, 2026 -** PANORAMIX -** File description: -** Queue is empty -*/ -#include "queue.h" -#include - -bool queue_empty(const queue_t *q) -{ - return (q == NULL || q->size == 0); -} diff --git a/src/threadpool/queue/queue_front.c b/src/threadpool/queue/queue_front.c deleted file mode 100644 index 34799f7..0000000 --- a/src/threadpool/queue/queue_front.c +++ /dev/null @@ -1,15 +0,0 @@ -/* -** EPITECH PROJECT, 2026 -** PANORAMIX -** File description: -** Queue front -*/ -#include "queue.h" -#include - -void *queue_front(const queue_t *q) -{ - if (!q || q->size == 0 || !q->head) - return NULL; - return q->head->data; -} diff --git a/src/threadpool/queue/queue_pop.c b/src/threadpool/queue/queue_pop.c deleted file mode 100644 index b9711f8..0000000 --- a/src/threadpool/queue/queue_pop.c +++ /dev/null @@ -1,26 +0,0 @@ -/* -** EPITECH PROJECT, 2026 -** PANORAMIX -** File description: -** Queue pop -*/ -#include "queue.h" -#include -#include - -void *queue_pop(queue_t *q) -{ - queue_node_t *temp; - void *data; - - if (queue_empty(q)) - return NULL; - temp = q->head; - data = temp->data; - q->head = q->head->next; - if (q->head == NULL) - q->tail = NULL; - free(temp); - q->size--; - return data; -} diff --git a/src/threadpool/queue/queue_push.c b/src/threadpool/queue/queue_push.c deleted file mode 100644 index f5631ce..0000000 --- a/src/threadpool/queue/queue_push.c +++ /dev/null @@ -1,31 +0,0 @@ -/* -** EPITECH PROJECT, 2026 -** PANORAMIX -** File description: -** Queue push -*/ -#include "queue.h" -#include -#include - -bool queue_push(queue_t *q, void *data) -{ - queue_node_t *new_node; - - if (!q) - return false; - new_node = malloc(sizeof(queue_node_t)); - if (!new_node) - return false; - new_node->data = data; - new_node->next = NULL; - if (queue_empty(q)) { - q->head = new_node; - q->tail = new_node; - } else { - q->tail->next = new_node; - q->tail = new_node; - } - q->size++; - return true; -} diff --git a/src/threadpool/queue/queue_size.c b/src/threadpool/queue/queue_size.c deleted file mode 100644 index 6ff1aed..0000000 --- a/src/threadpool/queue/queue_size.c +++ /dev/null @@ -1,15 +0,0 @@ -/* -** EPITECH PROJECT, 2026 -** PANORAMIX -** File description: -** Queue is empty -*/ -#include "queue.h" -#include - -size_t queue_size(const queue_t *q) -{ - if (!q->size) - return 0; - return q->size; -} diff --git a/src/threadpool/test.c b/src/threadpool/test.c deleted file mode 100644 index c394c61..0000000 --- a/src/threadpool/test.c +++ /dev/null @@ -1,32 +0,0 @@ -#include "kronkflow/macros/types.h" -#include "kronkflow/utils/threadpool.h" -#include -#include -#include -#include -#include "threadpool.h" - -static void *__test( - void *arg -) -{ - printf("%ld -> %d\n", (long int)arg, rand()); - return NULL; -} - -[[gnu::constructor]] -static void seedrand(void) -{ - srand(time(NULL)); -} - -KF_API -void kfThreadPool_test(void) -{ - kfThreadPool *pool = kfThreadPool_create(8, NULL); - - kfThreadPool_pushTask(pool, __test, (void *)10); - // kfThreadPool_join(pool); - printf("Test threadpool\n"); - kfThreadPool_destroy(pool); -} diff --git a/src/threadpool/threadpool.h b/src/threadpool/threadpool.h deleted file mode 100644 index 59aa906..0000000 --- a/src/threadpool/threadpool.h +++ /dev/null @@ -1,57 +0,0 @@ -/* -** FREE PROJECT, 2026 -** Kronkflow -** File description: -** threadpool -*/ -#ifndef KRONKFLOW_THREADPOOL_PRIVATE_H - #define KRONKFLOW_THREADPOOL_PRIVATE_H - #include - #include - #include - #include - #include "kronkflow/macros/types.h" - #include "queue/queue.h" - -typedef void *(*kfThreadPoolHandler)(void *); - -/////////////////////////////////////////////////////////////////////////////// -/** - * @struct kronkflow_threadpool_s - * - * @brief This struct is a dedicated threadpool for kronkflow multithreading features - */ -/////////////////////////////////////////////////////////////////////////////// -typedef struct kronkflow_threadpool_s { - - atomic_size_t pendings; //!< The number of task pendings - atomic_size_t runnings; //!< The number of thread up and runnings - atomic_size_t workers; //!< The number of threads - pthread_t* threads; //!< The threads (array) - pthread_cond_t cond; //!< The conditionnal variable - pthread_mutex_t mutex; //!< Mutex - kfQueue queue; //!< Queue for tasks - kfBool stop; //!< Does the pool should stop - void* ctx; //!< The ctx to give... - -} kfThreadPool; -/////////////////////////////////////////////////////////////////////////////// - -typedef struct kronkflow_thread_task_s { - - kfThreadPoolHandler handler; - void* data; - -} kfThreadTask; - -// TODO: Documentation -kfThreadPool *kfThreadPool_create(ssize_t nthreads, void *ctx); -int kfThreadPool_init(kfThreadPool *pool, size_t nthreads, void *ctx); -void kfThreadPool_destroy(kfThreadPool *pool); -void kfThreadPool_clear(kfThreadPool *pool); -void kfThreadPool_stop(kfThreadPool *pool); -size_t kfThreadPool_running(const kfThreadPool *pool); -size_t kfThreadPool_remaining(const kfThreadPool *pool); -int kfThreadPool_pushTask(kfThreadPool *pool, kfThreadPoolHandler task, void *data); - -#endif /* KRONKFLOW_THREADPOOL_PRIVATE_H */ diff --git a/src/threadpool/threadpool_create.c b/src/threadpool/threadpool_create.c deleted file mode 100644 index b049f82..0000000 --- a/src/threadpool/threadpool_create.c +++ /dev/null @@ -1,106 +0,0 @@ -/* -** FREE PROJECT, 2026 -** Kronkflow -** File description: -** threadpool -*/ -#include "kronkflow/macros/optimization.h" -#include "kronkflow/macros/types.h" -#include "queue/queue.h" -#include "threadpool.h" -#include -#include -#include -#include -#include -#include -#include - -kfThreadPool *kfThreadPool_create( - ssize_t nthreads, - void *ctx -) -{ - kfThreadPool *pool = NULL; - - if (nthreads == 0) { - return NULL; - } else if (nthreads == -1) { - nthreads = sysconf(_SC_NPROCESSORS_ONLN); - if (nthreads <= 0) { - return NULL; - } - } - pool = calloc(1, sizeof(kfThreadPool)); - if (!pool) { - return NULL; - } - // NOTE: Nthread is already good, so init won't check... - if (kfThreadPool_init(pool, nthreads, ctx) == -1) { - free(pool); - return NULL; - } - return pool; -} - -static void *__routine( - void *arg -) -{ - kfThreadPool *pool = (kfThreadPool *)arg; - - if (!pool) { - return NULL; - } - while (true) { - kfThreadTask *task; - // void *d; - pthread_mutex_lock(&pool->mutex); - while (!pool->stop && queue_empty(&pool->queue)) { - pthread_cond_wait(&pool->cond, &pool->mutex); - } - // FIXME: Queue empty really necessary ?? - if (pool->stop && queue_empty(&pool->queue)) { - return NULL; - } - task = queue_front(&pool->queue); - queue_pop(&pool->queue); - pool->pendings--; - pthread_mutex_unlock(&pool->mutex); - pool->runnings++; - // NOTE: Should call task... with ctx - task->handler(task->data); - pool->runnings--; - } -} - -int kfThreadPool_init( - kfThreadPool *pool, - size_t nthreads, - void *ctx -) -{ - if (!pool) { - return -1; - } - if (pthread_mutex_init(&pool->mutex, NULL) != 0) { - return -1; - } - if (pthread_cond_init(&pool->cond, NULL) != 0) { - return -1; - } - queue_init(&pool->queue); - pool->pendings = 0; - pool->runnings = 0; - pool->ctx = ctx; - pool->workers = nthreads; - pool->stop = kfFalse; - pool->threads = calloc(nthreads, sizeof(pthread_t)); - if (!pool->threads) { - return -1; - } - for (size_t i = 0; i < nthreads; ++i) { - pthread_create(&pool->threads[i], NULL, &__routine, pool); - } - return 0; -} diff --git a/src/threadpool/threadpool_destroy.c b/src/threadpool/threadpool_destroy.c deleted file mode 100644 index c12f66b..0000000 --- a/src/threadpool/threadpool_destroy.c +++ /dev/null @@ -1,38 +0,0 @@ -/* -** FREE PROJECT, 2026 -** Kronkflow -** File description: -** threadpool destroy -*/ -#include "queue/queue.h" -#include "threadpool.h" -#include -#include -#include -#include - -void kfThreadPool_destroy( - kfThreadPool *pool -) -{ - if (!pool) - return; - kfThreadPool_clear(pool); - free(pool); -} - -void kfThreadPool_clear( - kfThreadPool *pool -) -{ - // pthread_mutex_lock(&pool->mutex); - pool->stop = true; - // pthread_mutex_unlock(&pool->mutex); - for (size_t i = 0; i < pool->workers; ++i) { - pthread_join(pool->threads[i], NULL); - } - pthread_mutex_destroy(&pool->mutex); - pthread_cond_destroy(&pool->cond); - queue_destroy(&pool->queue, NULL); - free(pool->threads); -} diff --git a/src/threadpool/threadpool_pushTask.c b/src/threadpool/threadpool_pushTask.c deleted file mode 100644 index a35330a..0000000 --- a/src/threadpool/threadpool_pushTask.c +++ /dev/null @@ -1,35 +0,0 @@ -/* -** FREE PROJECT, 2026 -** Kronkflow -** File description: -** threadpool -*/ -#include "queue/queue.h" -#include "threadpool.h" -#include -#include - -int kfThreadPool_pushTask( - kfThreadPool *pool, - void *(*task)(void *), - void *data -) -{ - kfThreadTask *p = NULL; - - if (!pool) { - return -1; - } - p = calloc(1, sizeof(kfThreadTask)); - if (!p) { - return -1; - } - p->data = data; - p->handler = task; - pthread_mutex_lock(&pool->mutex); - queue_push(&pool->queue, p); - ++pool->pendings; - pthread_cond_signal(&pool->cond); - pthread_mutex_unlock(&pool->mutex); - return 0; -}