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
88 changes: 86 additions & 2 deletions tests/test_priority.c
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
#include <ttak/priority/queue.h>
#include <ttak/async/task.h>
#include <stdint.h>
#include "test_macros.h"

void *dummy_func(void *arg) { (void)arg; return NULL; }

void test_priority_queue_basic() {
void test_priority_queue_basic(void) {
struct __internal_ttak_proc_priority_queue_t q;
ttak_priority_queue_init(&q);

Expand All @@ -29,7 +30,90 @@ void test_priority_queue_basic() {
ttak_task_destroy(t2, now);
}

int main() {
/* A priority-0 task in heap slot 0 must survive push/pop (it used to
* encode to NULL and be silently dropped). */
void test_priority_queue_zero_priority_slot0(void) {
struct __internal_ttak_proc_priority_queue_t q;
ttak_priority_queue_init(&q);

uint64_t now = 100;
ttak_task_t *t0 = ttak_task_create(dummy_func, NULL, NULL, now);
ASSERT(t0 != NULL);

q.push(&q, t0, 0, now);
ASSERT(q.get_size(&q) == 1);
ASSERT(q.pop(&q, now) == t0);
ASSERT(q.get_size(&q) == 0);

/* Slot-0 entry buried under higher priorities must still come out. */
q.push(&q, t0, 0, now);
for (int i = 0; i < 32; i++) {
q.push(&q, t0, i + 1, now);
}
ASSERT(q.get_size(&q) == 33);
for (int i = 0; i < 32; i++) {
ASSERT(q.pop(&q, now) == t0);
}
ASSERT(q.pop(&q, now) == t0);
ASSERT(q.get_size(&q) == 0);

ttak_task_destroy(t0, now);
}

void test_priority_queue_heap_stress(void) {
struct __internal_ttak_proc_priority_queue_t q;
ttak_priority_queue_init(&q);

uint64_t now = 100;
const int N = 5000;
ttak_task_t *t = ttak_task_create(dummy_func, NULL, NULL, now);
ASSERT(t != NULL);

/* Deterministic LCG; skewed toward 0 with duplicates to exercise the
* free list and the comparator. */
uint32_t rng = 0x12345678u;
int model[5000];
for (int i = 0; i < N; i++) {
rng = rng * 1664525u + 1013904223u;
int p;
switch (rng % 4) {
case 0: p = 0; break;
case 1: p = 0; break; /* half of all pushes are priority 0 */
case 2: p = 1; break;
default: p = 7; break;
}
if (i % 997 == 0) p = 42;
model[i] = p;
q.push(&q, t, p, now);
}
ASSERT(q.get_size(&q) == (size_t)N);

/* Pop all and check order against a max-extraction model. */
char removed[5000] = {0};
int prev = INT32_MAX;
for (int k = 0; k < N; k++) {
ttak_task_t *popped = q.pop(&q, now);
ASSERT(popped == t);
int best = -1;
for (int i = 0; i < N; i++) {
if (!removed[i] && (best < 0 || model[i] > model[best])) best = i;
}
ASSERT(best >= 0);
removed[best] = 1;
ASSERT_MSG(model[best] <= prev,
"pop %d: priority %d after %d (not non-increasing)",
k, model[best], prev);
prev = model[best];
}
ASSERT(q.get_size(&q) == 0);
ASSERT(q.pop(&q, now) == NULL);

ttak_task_destroy(t, now);
}

int main(void) {
RUN_TEST(test_priority_queue_basic);
RUN_TEST(test_priority_queue_zero_priority_slot0);
RUN_TEST(test_priority_queue_heap_stress);
return 0;
}
79 changes: 79 additions & 0 deletions tests/test_ringbuf.c
Original file line number Diff line number Diff line change
@@ -1,4 +1,6 @@
#include <ttak/container/ringbuf.h>
#include <pthread.h>
#include <stdatomic.h>
#include "test_macros.h"

static void test_ringbuf_push_pop_cycle(void) {
Expand All @@ -22,7 +24,84 @@ static void test_ringbuf_push_pop_cycle(void) {
ttak_ringbuf_destroy(rb);
}

/* Query functions must stay consistent under concurrent push/pop. */
#define RB_CAPACITY 64
#define RB_ROUNDS 20000

typedef struct {
ttak_ringbuf_t *rb;
_Atomic int *violations;
_Atomic int *pushed;
_Atomic int *popped;
} rb_shared_t;

static void *rb_producer(void *arg) {
rb_shared_t *s = (rb_shared_t *)arg;
for (int i = 0; i < RB_ROUNDS; ++i) {
if (ttak_ringbuf_push(s->rb, &i)) {
atomic_fetch_add(s->pushed, 1);
}
}
return NULL;
}

static void *rb_consumer(void *arg) {
rb_shared_t *s = (rb_shared_t *)arg;
int out = 0;
for (int i = 0; i < RB_ROUNDS; ++i) {
if (ttak_ringbuf_pop(s->rb, &out)) {
atomic_fetch_add(s->popped, 1);
}
}
return NULL;
}

static void *rb_querier(void *arg) {
rb_shared_t *s = (rb_shared_t *)arg;
for (int i = 0; i < RB_ROUNDS * 4; ++i) {
size_t count = ttak_ringbuf_count(s->rb);
/* Each query call is a separate locked snapshot, so cross-checks
* between them would race; only per-snapshot bounds are asserted. */
(void)ttak_ringbuf_is_empty(s->rb);
(void)ttak_ringbuf_is_full(s->rb);
if (count > RB_CAPACITY) {
atomic_fetch_add(s->violations, 1);
return NULL;
}
}
return NULL;
}

static void test_ringbuf_concurrent_query_race(void) {
ttak_ringbuf_t *rb = ttak_ringbuf_create(RB_CAPACITY, sizeof(int));
_Atomic int violations = 0;
_Atomic int pushed = 0;
_Atomic int popped = 0;
rb_shared_t shared = { rb, &violations, &pushed, &popped };

pthread_t producer, consumer, querier1, querier2;
ASSERT(pthread_create(&producer, NULL, rb_producer, &shared) == 0);
ASSERT(pthread_create(&consumer, NULL, rb_consumer, &shared) == 0);
ASSERT(pthread_create(&querier1, NULL, rb_querier, &shared) == 0);
ASSERT(pthread_create(&querier2, NULL, rb_querier, &shared) == 0);

pthread_join(producer, NULL);
pthread_join(consumer, NULL);
pthread_join(querier1, NULL);
pthread_join(querier2, NULL);

ASSERT(violations == 0);

/* Buffered items must equal pushes minus pops. */
size_t count = ttak_ringbuf_count(rb);
ASSERT(count <= RB_CAPACITY);
ASSERT(count == (size_t)(pushed - popped));

ttak_ringbuf_destroy(rb);
}

int main(void) {
RUN_TEST(test_ringbuf_push_pop_cycle);
RUN_TEST(test_ringbuf_concurrent_query_race);
return 0;
}
41 changes: 41 additions & 0 deletions tests/test_sharded_sched.c
Original file line number Diff line number Diff line change
Expand Up @@ -262,6 +262,46 @@ void test_worker_shard_affinity(void) {
ttak_thread_pool_destroy(pool);
}

/* -------------------------------------------------------------------------
* 8. Priority-0 tasks must never be dropped or starved, even when one lands
* in heap slot 0 of the task queue.
* ---------------------------------------------------------------------- */
static _Atomic int prio_zero_counter = 0;

void *prio_zero_task(void *arg) {
(void)arg;
prio_zero_counter++;
return NULL;
}

void test_pool_priority_zero_no_starvation(void) {
uint64_t now = ttak_get_tick_count();
ttak_thread_pool_t *pool = ttak_thread_pool_create(4, 0, now);
ASSERT(pool != NULL);

prio_zero_counter = 0;
const int N = 64;
ttak_future_t *futures[64];

for (int i = 0; i < N; i++) {
futures[i] = ttak_thread_pool_submit_task(pool, prio_zero_task, NULL, 0, now);
ASSERT(futures[i] != NULL);
}

for (int i = 0; i < N; i++) {
ttak_future_get(futures[i]);
}

ASSERT_MSG(prio_zero_counter == N,
"only %d of %d priority-0 tasks completed", prio_zero_counter, N);

uint64_t elapsed = ttak_get_tick_count() - now;
ASSERT_MSG(elapsed < 10000, "priority-0 burst took %llu ms (expected < 10000)",
(unsigned long long)elapsed);

ttak_thread_pool_destroy(pool);
}

int main(void) {
RUN_TEST(test_route_table_bounds);
RUN_TEST(test_shard_mapping_deterministic);
Expand All @@ -270,6 +310,7 @@ int main(void) {
RUN_TEST(test_sharded_scheduler_history);
RUN_TEST(test_pool_sharded_routing);
RUN_TEST(test_pool_sharded_routing_net_urgency);
RUN_TEST(test_pool_priority_zero_no_starvation);
RUN_TEST(test_work_stealing);
return 0;
}
Loading