Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 5 additions & 5 deletions examples/texec_example.c
Original file line number Diff line number Diff line change
Expand Up @@ -10,14 +10,14 @@ static int hello_task(void* user) {

int main(void) {
const texec_executor_create_thread_pool_info_t tpci = {
.header = {.type = TEXEC_STRUCTURE_TYPE_EXECUTOR_CREATE_THREAD_POOL_INFO, .next = NULL},
.header = {.type = TEXEC_STRUCT_TYPE_EXECUTOR_CREATE_THREAD_POOL_INFO, .next = NULL},
.thread_count = 2,
.queue_capacity = 32,
.backpressure = TEXEC_BACKPRESSURE_BLOCK,
};

const texec_executor_create_info_t eci = {
.header = {.type = TEXEC_STRUCTURE_TYPE_EXECUTOR_CREATE_INFO, .next = &tpci},
.header = {.type = TEXEC_STRUCT_TYPE_EXECUTOR_CREATE_INFO, .next = &tpci},
.kind = TEXEC_EXECUTOR_KIND_THREAD_POOL,
};

Expand All @@ -29,10 +29,10 @@ int main(void) {
return 1;
}

const texec_executor_submit_info_t submit = {
.header = {.type = TEXEC_STRUCTURE_TYPE_EXECUTOR_SUBMIT_INFO, .next = NULL},
const texec_submit_info_t submit = {
.header = {.type = TEXEC_STRUCT_TYPE_SUBMIT_INFO, .next = NULL},
.task = {
.fn = hello_task,
.run = hello_task,
.ctx = "work item",
},
};
Expand Down
32 changes: 16 additions & 16 deletions include/texec/base.h
Original file line number Diff line number Diff line change
Expand Up @@ -18,23 +18,23 @@ typedef enum texec_status {
TEXEC_STATUS_INTERNAL_ERROR
} texec_status_t;

typedef enum texec_structure_type {
TEXEC_STRUCTURE_TYPE_EXECUTOR_CREATE_INFO = 0x1000,
TEXEC_STRUCTURE_TYPE_EXECUTOR_SUBMIT_INFO = 0x2000,
TEXEC_STRUCTURE_TYPE_TASK_GROUP_CREATE_INFO = 0x3000,
TEXEC_STRUCTURE_TYPE_QUEUE_CREATE_INFO = 0x4000,
typedef enum texec_struct_type {
TEXEC_STRUCT_TYPE_EXECUTOR_CREATE_INFO = 0x1000,
TEXEC_STRUCT_TYPE_SUBMIT_INFO = 0x2000,
TEXEC_STRUCT_TYPE_TASK_GROUP_CREATE_INFO = 0x3000,
TEXEC_STRUCT_TYPE_QUEUE_CREATE_INFO = 0x4000,

TEXEC_STRUCTURE_TYPE_EXECUTOR_CREATE_INLINE_INFO = 0x1001,
TEXEC_STRUCTURE_TYPE_EXECUTOR_CREATE_THREAD_POOL_INFO = 0x1002,
TEXEC_STRUCTURE_TYPE_EXECUTOR_CREATE_DIAGNOSTICS_INFO = 0x1003,
TEXEC_STRUCT_TYPE_EXECUTOR_CREATE_INLINE_INFO = 0x1001,
TEXEC_STRUCT_TYPE_EXECUTOR_CREATE_THREAD_POOL_INFO = 0x1002,
TEXEC_STRUCT_TYPE_EXECUTOR_CREATE_DIAGNOSTICS_INFO = 0x1003,

TEXEC_STRUCTURE_TYPE_EXECUTOR_SUBMIT_PRIORITY = 0x2001,
TEXEC_STRUCTURE_TYPE_EXECUTOR_SUBMIT_DEADLINE = 0x2002,
TEXEC_STRUCTURE_TYPE_EXECUTOR_SUBMIT_TRACE_CONTEXT = 0x2003,
TEXEC_STRUCTURE_TYPE_EXECUTOR_SUBMIT_BACKPRESSURE = 0x2004,
TEXEC_STRUCT_TYPE_SUBMIT_PRIORITY = 0x2001,
TEXEC_STRUCT_TYPE_SUBMIT_DEADLINE = 0x2002,
TEXEC_STRUCT_TYPE_SUBMIT_TRACE_CONTEXT = 0x2003,
TEXEC_STRUCT_TYPE_SUBMIT_BACKPRESSURE = 0x2004,

TEXEC_STRUCTURE_TYPE_QUEUE_CREATE_FULL_POLICY_INFO = 0x4001,
} texec_structure_type_t;
TEXEC_STRUCT_TYPE_QUEUE_CREATE_FULL_POLICY_INFO = 0x4001,
} texec_struct_type_t;

typedef enum texec_backpressure_policy {
TEXEC_BACKPRESSURE_REJECT = 0,
Expand All @@ -54,11 +54,11 @@ typedef struct texec_allocator {
void texec_set_default_allocator(const texec_allocator_t* allocator);

typedef struct texec_structure_header {
texec_structure_type_t type;
texec_struct_type_t type;
const void* next;
} texec_structure_header_t;

static inline const void* texec_structure_find(const void* first, texec_structure_type_t type) {
static inline const void* texec_structure_find(const void* first, texec_struct_type_t type) {
const texec_structure_header_t* it = (const texec_structure_header_t*)first;
while (it && it->type != type) {
it = (const texec_structure_header_t*)it->next;
Expand Down
23 changes: 23 additions & 0 deletions include/texec/diagnostics.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
#pragma once

#ifdef __cplusplus
extern "C" {
#endif

struct texec_submit_info_t;
struct texec_task;

typedef void (*texec_on_submit_fn_t)(void* user, const struct texec_submit_info_t* submit_info);
typedef void (*texec_on_task_begin_fn_t)(void* user, const struct texec_task* task, const void* trace_context);
typedef void (*texec_on_task_end_fn_t)(void* user, const struct texec_task* task, const void* trace_context, int task_result);

typedef struct texec_diagnostics {
void* user;
texec_on_submit_fn_t on_submit;
texec_on_task_begin_fn_t on_task_begin;
texec_on_task_end_fn_t on_task_end;
} texec_diagnostics_t;

#ifdef __cplusplus
}
#endif
4 changes: 2 additions & 2 deletions include/texec/executor.h
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,8 @@ typedef struct texec_executor texec_executor_t;

texec_status_t texec_executor_create(const texec_executor_create_info_t* info, const texec_allocator_t* allocator, texec_executor_t** out_executor);
texec_status_t texec_executor_destroy(texec_executor_t* ex);
texec_status_t texec_executor_submit(texec_executor_t* ex, const texec_executor_submit_info_t* info, texec_task_handle_t** out_handle);
texec_status_t texec_executor_submit_many(texec_executor_t* ex, const texec_executor_submit_info_t* infos, size_t count, texec_task_group_t** out_group);
texec_status_t texec_executor_submit(texec_executor_t* ex, const texec_submit_info_t* info, texec_task_handle_t** out_handle);
texec_status_t texec_executor_submit_many(texec_executor_t* ex, const texec_submit_info_t* infos, size_t count, texec_task_group_t** out_group);
void texec_executor_close(texec_executor_t* ex);
void texec_executor_join(texec_executor_t* ex);

Expand Down
16 changes: 2 additions & 14 deletions include/texec/executor_create_info.h
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
#include <stddef.h>

#include "texec/base.h"
#include "texec/task.h"
#include "texec/diagnostics.h"

#ifdef __cplusplus
extern "C" {
Expand All @@ -28,21 +28,9 @@ typedef struct texec_executor_create_thread_pool_info {
texec_backpressure_policy_t backpressure;
} texec_executor_create_thread_pool_info_t;

struct texec_executor_submit_info;
typedef void (*texec_executor_on_submit_fn_t)(void* user, const struct texec_executor_submit_info* submit_info);
typedef void (*texec_executor_on_task_begin_fn_t)(void* user, const texec_task_t* task, const void* trace_context);
typedef void (*texec_executor_on_task_end_fn_t)(void* user, const texec_task_t* task, const void* trace_context, int task_result);

typedef struct texec_executor_diagnostics {
void* user;
texec_executor_on_submit_fn_t on_submit;
texec_executor_on_task_begin_fn_t on_task_begin;
texec_executor_on_task_end_fn_t on_task_end;
} texec_executor_diagnostics_t;

typedef struct texec_executor_create_diagnostics_info {
texec_structure_header_t header;
const texec_executor_diagnostics_t* diag;
const texec_diagnostics_t* diag;
} texec_executor_create_diagnostics_info_t;

#ifdef __cplusplus
Expand Down
32 changes: 16 additions & 16 deletions include/texec/executor_submit_info.h
Original file line number Diff line number Diff line change
Expand Up @@ -9,38 +9,38 @@
extern "C" {
#endif

typedef struct texec_executor_submit_info {
typedef struct texec_submit_info {
texec_structure_header_t header;
texec_task_t task;
} texec_executor_submit_info_t;
} texec_submit_info_t;

// --- Submit Extensions ---

typedef enum texec_executor_submit_priority {
TEXEC_EXECUTOR_SUBMIT_PRIORITY_LOW = -1,
TEXEC_EXECUTOR_SUBMIT_PRIORITY_NORMAL = 0,
TEXEC_EXECUTOR_SUBMIT_PRIORITY_HIGH = 1
} texec_executor_submit_priority_t;
typedef enum texec_submit_priority {
TEXEC_SUBMIT_PRIORITY_LOW = -1,
TEXEC_SUBMIT_PRIORITY_NORMAL = 0,
TEXEC_SUBMIT_PRIORITY_HIGH = 1
} texec_submit_priority_t;

typedef struct texec_executor_submit_priority_info {
typedef struct texec_submit_priority_info {
texec_structure_header_t header;
texec_executor_submit_priority_t priority;
} texec_executor_submit_priority_info_t;
texec_submit_priority_t priority;
} texec_submit_priority_info_t;

typedef struct texec_executor_submit_deadline_info {
typedef struct texec_submit_deadline_info {
texec_structure_header_t header;
uint64_t deadline_ns;
} texec_executor_submit_deadline_info_t;
} texec_submit_deadline_info_t;

typedef struct texec_executor_submit_trace_context_info {
typedef struct texec_submit_trace_context_info {
texec_structure_header_t header;
const void* trace_context;
} texec_executor_submit_trace_context_info_t;
} texec_submit_trace_context_info_t;

typedef struct texec_executor_submit_backpressure_info {
typedef struct texec_submit_backpressure_info {
texec_structure_header_t header;
texec_backpressure_policy_t backpressure;
} texec_executor_submit_backpressure_info_t;
} texec_submit_backpressure_info_t;

#ifdef __cplusplus
}
Expand Down
8 changes: 4 additions & 4 deletions include/texec/task.h
Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,13 @@
extern "C" {
#endif

typedef int (*texec_task_fn_t)(void* ctx);
typedef void (*texec_task_cleanup_fn_t)(void* ctx);
typedef int (*texec_task_run_t)(void* ctx);
typedef void (*texec_task_on_complete_fn_t)(void* ctx);

typedef struct texec_task {
texec_task_fn_t fn;
texec_task_run_t run;
void* ctx;
texec_task_cleanup_fn_t cleanup; // optional; called after fn, on the executing thread
texec_task_on_complete_fn_t on_complete; // optional; called after run, on the executing thread
} texec_task_t;

#ifdef __cplusplus
Expand Down
18 changes: 9 additions & 9 deletions src/executor.c
Original file line number Diff line number Diff line change
Expand Up @@ -9,20 +9,20 @@ static const size_t TP_EXECUTOR_DEFAULT_THREAD_COUNT = 1;
static const size_t TP_EXECUTOR_DEFAULT_QUEUE_CAPACITY = 1024;

static inline const texec_executor_create_thread_pool_info_t*
find_executor_thread_pool_info(const texec_executor_create_info_t* info) {
return texec_structure_find(info->header.next, TEXEC_STRUCTURE_TYPE_EXECUTOR_CREATE_THREAD_POOL_INFO);
find_executor_thread_pool_create_info(const texec_executor_create_info_t* info) {
return texec_structure_find(info->header.next, TEXEC_STRUCT_TYPE_EXECUTOR_CREATE_THREAD_POOL_INFO);
}

static inline const texec_executor_create_diagnostics_info_t*
find_executor_diag_info(const texec_executor_create_info_t* info) {
return texec_structure_find(info->header.next, TEXEC_STRUCTURE_TYPE_EXECUTOR_CREATE_DIAGNOSTICS_INFO);
return texec_structure_find(info->header.next, TEXEC_STRUCT_TYPE_EXECUTOR_CREATE_DIAGNOSTICS_INFO);
}

static inline texec_status_t executor_create_thread_pool(const texec_allocator_t* alloc,
const texec_executor_diagnostics_t* diag,
const texec_diagnostics_t* diag,
const texec_executor_create_info_t* info,
texec_executor_t** out_ex) {
const texec_executor_create_thread_pool_info_t* tp_info = find_executor_thread_pool_info(info);
const texec_executor_create_thread_pool_info_t* tp_info = find_executor_thread_pool_create_info(info);
if (!tp_info) return TEXEC_STATUS_INVALID_ARGUMENT;

const texec_thread_pool_executor_config_t cfg = {
Expand Down Expand Up @@ -53,12 +53,12 @@ texec_status_t texec_executor_create(const texec_executor_create_info_t* info, c
if (!out_executor) return TEXEC_STATUS_INVALID_ARGUMENT;
*out_executor = NULL;

if (!info || info->header.type != TEXEC_STRUCTURE_TYPE_EXECUTOR_CREATE_INFO) {
if (!info || info->header.type != TEXEC_STRUCT_TYPE_EXECUTOR_CREATE_INFO) {
return TEXEC_STATUS_INVALID_ARGUMENT;
}

const texec_executor_create_diagnostics_info_t* diag_info = find_executor_diag_info(info);
const texec_executor_diagnostics_t* diag = diag_info ? diag_info->diag : NULL;
const texec_diagnostics_t* diag = diag_info ? diag_info->diag : NULL;

if (!alloc) {
alloc = texec_get_default_allocator();
Expand Down Expand Up @@ -90,12 +90,12 @@ texec_status_t texec_executor_destroy(texec_executor_t* ex) {
return ex->vtbl->destroy(ex);
}

texec_status_t texec_executor_submit(texec_executor_t* ex, const texec_executor_submit_info_t* info, texec_task_handle_t** out_handle) {
texec_status_t texec_executor_submit(texec_executor_t* ex, const texec_submit_info_t* info, texec_task_handle_t** out_handle) {
if (!ex || !out_handle) return TEXEC_STATUS_INVALID_ARGUMENT;
return ex->vtbl->submit(ex, info, out_handle);
}

texec_status_t texec_executor_submit_many(texec_executor_t* ex, const texec_executor_submit_info_t* infos, size_t count, texec_task_group_t** out_group) {
texec_status_t texec_executor_submit_many(texec_executor_t* ex, const texec_submit_info_t* infos, size_t count, texec_task_group_t** out_group) {
if (!ex || !out_group) return TEXEC_STATUS_INVALID_ARGUMENT;
return ex->vtbl->submit_many(ex, infos, count, out_group);
}
Expand Down
18 changes: 18 additions & 0 deletions src/internal/diagnostics.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
#pragma once

#include "texec/diagnostics.h"

static inline void texec_diagnostics_on_submit(const texec_diagnostics_t* diag, const struct texec_submit_info_t* submit_info) {
if (!diag) return;
diag->on_submit(diag->user, submit_info);
}

static inline void texec_diagnostics_on_task_begin(const texec_diagnostics_t* diag, const struct texec_task* task, const void* trace_context) {
if (!diag) return;
diag->on_task_begin(diag->user, task, trace_context);
}

static inline void texec_diagnostics_on_task_end(const texec_diagnostics_t* diag, const struct texec_task* task, const void* trace_context, int task_result) {
if (!diag) return;
diag->on_task_end(diag->user, task, trace_context, task_result);
}
34 changes: 12 additions & 22 deletions src/internal/executor.h
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
#include "texec/task_handle.h"

#include "internal/allocator.h"
#include "internal/diagnostics.h"
#include "internal/task_handle.h"
#include "internal/work_item.h"

Expand All @@ -14,8 +15,8 @@ typedef enum texec_executor_state {
TEXEC_EXECUTOR_STATE_CLOSED,
} texec_executor_state_t;

typedef texec_status_t (*texec_executor_submit_fn_t)(texec_executor_t* ex, const texec_executor_submit_info_t* info, texec_task_handle_t** out_handle);
typedef texec_status_t (*texec_executor_submit_many_fn_t)(texec_executor_t* ex, const texec_executor_submit_info_t* infos, size_t count, texec_task_group_t** out_group);
typedef texec_status_t (*texec_executor_submit_fn_t)(texec_executor_t* ex, const texec_submit_info_t* info, texec_task_handle_t** out_handle);
typedef texec_status_t (*texec_executor_submit_many_fn_t)(texec_executor_t* ex, const texec_submit_info_t* infos, size_t count, texec_task_group_t** out_group);
typedef void (*texec_executor_close_fn_t)(texec_executor_t* ex);
typedef void (*texec_executor_join_fn_t)(texec_executor_t* ex);
typedef texec_status_t (*texec_executor_destroy_fn_t)(texec_executor_t* ex);
Expand All @@ -33,42 +34,31 @@ typedef struct texec_executor_vtable {
struct texec_executor {
const texec_executor_vtable_t* vtbl;
const texec_allocator_t* alloc;
const texec_executor_diagnostics_t* diag;
const texec_diagnostics_t* diag;
texec_executor_kind_t kind;
texec_executor_state_t state;
};

typedef struct texec_thread_pool_executor_config {
const texec_allocator_t* alloc;
const texec_executor_diagnostics_t* diag;
const texec_diagnostics_t* diag;
size_t thread_count;
size_t queue_capacity;
texec_backpressure_policy_t backpressure;
} texec_thread_pool_executor_config_t;

texec_status_t texec_executor_create_thread_pool(const texec_thread_pool_executor_config_t* cfg, texec_executor_t** out_ex);

static inline void texec_executor_diagnostics_on_submit(const texec_executor_diagnostics_t* diag, const struct texec_executor_submit_info* submit_info) {
if (diag) diag->on_submit(diag->user, submit_info);
}

static inline void texec_executor_diagnostics_on_task_begin(const texec_executor_diagnostics_t* diag, const texec_task_t* task, const void* trace_context) {
if (diag) diag->on_task_begin(diag->user, task, trace_context);
}

static inline void texec_executor_diagnostics_on_task_end(const texec_executor_diagnostics_t* diag, const texec_task_t* task, const void* trace_context, int task_result) {
if (diag) diag->on_task_end(diag->user, task, trace_context, task_result);
}

static inline void texec_task_cleanup(const texec_task_t* t) {
if (t->cleanup) t->cleanup(t->ctx);
static inline void texec_task_on_complete(const texec_task_t* t) {
if (!t->on_complete) return;
t->on_complete(t->ctx);
}

static inline void texec_executor_consume_work_item(const texec_executor_t* ex, texec_work_item_t* wi) {
texec_executor_diagnostics_on_task_begin(ex->diag, &wi->task, wi->trace_context);
const int result = wi->task.fn(wi->task.ctx);
texec_executor_diagnostics_on_task_end(ex->diag, &wi->task, wi->trace_context, result);
texec_task_cleanup(&wi->task);
texec_diagnostics_on_task_begin(ex->diag, &wi->task, wi->trace_context);
const int result = wi->task.run(wi->task.ctx);
texec_diagnostics_on_task_end(ex->diag, &wi->task, wi->trace_context, result);
texec_task_on_complete(&wi->task);
texec_task_handle_complete(wi->handle, result);
texec_work_item_destroy(wi, ex->alloc);
}
4 changes: 3 additions & 1 deletion src/queue.c
Original file line number Diff line number Diff line change
Expand Up @@ -141,7 +141,7 @@ texec_status_t texec_queue_create(const texec_queue_create_info_t* info, const t

*out_q = NULL;

if (!info || info->header.type != TEXEC_STRUCTURE_TYPE_QUEUE_CREATE_INFO || info->capacity == 0) {
if (!info || info->header.type != TEXEC_STRUCT_TYPE_QUEUE_CREATE_INFO || info->capacity == 0) {
return TEXEC_STATUS_INVALID_ARGUMENT;
}

Expand Down Expand Up @@ -176,6 +176,8 @@ texec_status_t texec_queue_destroy(texec_queue_t* q) {

texec_free(q->alloc, q->buf, q->capacity * sizeof(uintptr_t), _Alignof(uintptr_t));
texec_free(q->alloc, q, sizeof(*q), _Alignof(texec_queue_t));

return TEXEC_STATUS_OK;
}

void texec_queue_close(texec_queue_t* q) {
Expand Down
2 changes: 1 addition & 1 deletion src/task_group.c
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,7 @@ texec_status_t texec_task_group_create(const texec_task_group_create_info_t* inf
if (!out_group) return TEXEC_STATUS_INVALID_ARGUMENT;
*out_group = NULL;

if (!info || info->header.type != TEXEC_STRUCTURE_TYPE_TASK_GROUP_CREATE_INFO) {
if (!info || info->header.type != TEXEC_STRUCT_TYPE_TASK_GROUP_CREATE_INFO) {
return TEXEC_STATUS_INVALID_ARGUMENT;
}

Expand Down
Loading