From d6cf42a9cab969ea1f6f5f97526ac206da456ca1 Mon Sep 17 00:00:00 2001 From: Lee Yunjin Date: Sun, 4 Oct 2026 21:02:58 +0900 Subject: [PATCH] tests: add regression tests for recent priority queue and ringbuf fixes - test_sharded_sched: test_pool_priority_zero_no_starvation submits 64 tasks all at priority 0 through the thread pool and asserts every one completes via futures, covering the slot-0 NULL-element drop fixed in 95854ca8 that deadlocked submitters in ttak_future_get. - test_ringbuf: test_ringbuf_concurrent_query_race hammers push/pop from producer/consumer threads while two querier threads call ttak_ringbuf_count/is_empty/is_full, covering the data race fixed in 439a7ff9. Asserts the capacity bound and pushed-minus-popped sum sanity. - test_priority: test_priority_queue_zero_priority_slot0 reproduces the heap slot-0 NULL encoding directly at queue level, and test_priority_queue_heap_stress pushes 5000 items with mixed priorities (half of them 0, many duplicates) and verifies pop order is monotonically non-increasing against a max-extraction model, covering the heap restructure in cb804a94/95854ca8. --- tests/test_priority.c | 88 +++++++++++++++++++++++++++++++++++++- tests/test_ringbuf.c | 79 ++++++++++++++++++++++++++++++++++ tests/test_sharded_sched.c | 41 ++++++++++++++++++ 3 files changed, 206 insertions(+), 2 deletions(-) diff --git a/tests/test_priority.c b/tests/test_priority.c index 3b2aabb4..73563cb3 100644 --- a/tests/test_priority.c +++ b/tests/test_priority.c @@ -1,10 +1,11 @@ #include #include +#include #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); @@ -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; } diff --git a/tests/test_ringbuf.c b/tests/test_ringbuf.c index 118ef87f..a18080ab 100644 --- a/tests/test_ringbuf.c +++ b/tests/test_ringbuf.c @@ -1,4 +1,6 @@ #include +#include +#include #include "test_macros.h" static void test_ringbuf_push_pop_cycle(void) { @@ -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; } diff --git a/tests/test_sharded_sched.c b/tests/test_sharded_sched.c index f31f8d4e..e6720524 100644 --- a/tests/test_sharded_sched.c +++ b/tests/test_sharded_sched.c @@ -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); @@ -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; }