#include "akerror.h" #include "err_capture.h" #include "err_threads.h" #include /* * Handing an error context from one thread to another. * * A context is not thread state. It lives in AKERR_ARRAY_ERROR, which is * process-global, and the only field of it the library reads across an * ownership boundary is the reference count -- which is only ever touched under * the pool lock. So a context can be raised on one thread, handed to another, * and handled and released there, and akerr_release_error() does not care which * thread checked the slot out. That is a property this library promises, and it * is what makes the worker/collector shape usable at all. * * The other half of the promise is what it does *not* cover: two threads inside * one context at once. Ownership moves, it does not fork. This test asserts the * supported half; the unsupported half cannot be asserted without deliberately * racing, which ThreadSanitizer would then correctly fail. * * The queue below is a plain mutex and two condition variables rather than the * __atomic builtins the rest of these tests use, and that is deliberate: the * mutex *is* the thing under test. Context content is written with no lock at * all, so the handoff itself is what publishes those writes to the receiver. * * The claim is proved from four directions: * * 1. The context is still a live pool slot after it crosses, holding exactly * the one reference it was checked out with. * 2. Its message and its whole stack trace -- producer frame and all -- arrive * intact, and in each producer's own order. * 3. The slot is never recycled underneath the transfer: akerr_slot_owner[] * still names the producer when the collector picks it up. * 4. A context outlives the thread that raised it (see main()). */ #define ITERATIONS 500 #define AKERR_HANDOFF_DEPTH 32 /* * The sizing rule the README gives, made executable. Every queued error is a * checked-out pool slot, and so is every producer's error in flight. Outrun the * pool and ENSURE_ERROR_READY exits the process from inside FAIL, with no slot * left to raise the failure from. */ typedef char akerr_assert_handoff_fits_pool[ (AKERR_HANDOFF_DEPTH + AKERR_TEST_THREADS < AKERR_MAX_ARRAY_ERROR) ? 1 : -1]; static struct { pthread_mutex_t lock; pthread_cond_t not_full; pthread_cond_t not_empty; akerr_ErrorContext *slot[AKERR_HANDOFF_DEPTH]; int head; int count; } queue; /* Bounded on purpose: an unbounded queue of errors is an unbounded number of * checked-out pool slots. Blocking the producer is the backpressure. */ static void queue_push(akerr_ErrorContext *errctx) { pthread_mutex_lock(&queue.lock); while ( queue.count == AKERR_HANDOFF_DEPTH ) { pthread_cond_wait(&queue.not_full, &queue.lock); } queue.slot[(queue.head + queue.count) % AKERR_HANDOFF_DEPTH] = errctx; queue.count += 1; pthread_cond_signal(&queue.not_empty); pthread_mutex_unlock(&queue.lock); } static akerr_ErrorContext *queue_pop(void) { akerr_ErrorContext *errctx; pthread_mutex_lock(&queue.lock); while ( queue.count == 0 ) { pthread_cond_wait(&queue.not_empty, &queue.lock); } errctx = queue.slot[queue.head]; queue.head = (queue.head + 1) % AKERR_HANDOFF_DEPTH; queue.count -= 1; pthread_cond_signal(&queue.not_full); pthread_mutex_unlock(&queue.lock); return errctx; } /* * Raise an error and give it away. The stack-trace frame is appended before the * push so the trace records the crossing, and it is the last thing this thread * does to the context: after queue_push() returns, `e` belongs to the collector * and reading even e->status here would be the unsupported half of the rule. */ static void produce_one(akerr_ThreadArg *arg, int seq) { PREPARE_ERROR(e); FAIL(e, AKERR_VALUE, "thread %d seq %d", arg->id, seq); AKERR_TCHECK(arg, akerr_slot_claim(e->arrayid, arg->id) == 0); AKERR_STACKTRACE_APPEND(e, "queued by thread %d\n", arg->id); queue_push(e); } /* * One received error, handled and released on a thread that never called * akerr_next_error(). That release is the whole claim. * * `seen` is the collector's own per-producer sequence counter. Collector-local * means no atomics: keeping the producers in order is the queue's job, and * checking it is this thread's. */ static void collect_one(akerr_ThreadArg *arg, akerr_ErrorContext *e, int *seen) { char expected[64]; int producer = 0; int seq = 0; AKERR_TCHECK(arg, akerr_valid_error_address(e) == 1); /* It crossed holding exactly the reference it was checked out with. */ AKERR_TCHECK(arg, e->refcount == 1); AKERR_TCHECK(arg, sscanf(e->message, "thread %d seq %d", &producer, &seq) == 2); AKERR_TCHECK(arg, producer >= 2 && producer <= AKERR_TEST_THREADS); if ( producer >= 2 && producer <= AKERR_TEST_THREADS ) { AKERR_TCHECK(arg, seq == seen[producer]); seen[producer] += 1; } /* The slot still belongs to the producer, so nothing recycled it while it * was in flight. */ AKERR_TCHECK(arg, akerr_slot_holder(e->arrayid) == producer); snprintf(expected, sizeof(expected), "thread %d seq %d", producer, seq); /* Nothing to attempt -- the error is already in hand. The blocks are here * because this is the assembly the macros require, and because a real * collector reads exactly like this. */ ATTEMPT { } CLEANUP { } PROCESS(e) { /* case 0: a handed-off error that arrives with no status means somebody * wrote over the context after the producer let it go. */ int error_was_lost = 1; AKERR_TCHECK(arg, error_was_lost == 0); } HANDLE(e, AKERR_VALUE) { /* HANDLE rewinds the cursor, but the bytes are still there: the whole * trace crossed with the context, producer frame and handoff frame. */ AKERR_TCHECK(arg, strstr(e->stacktracebuf, expected) != NULL); AKERR_TCHECK(arg, strstr(e->stacktracebuf, "queued by thread") != NULL); /* Give the slot up before FINISH releases the context: the other order * hands it back to the pool while this thread still claims it. */ akerr_slot_drop(e->arrayid); } FINISH_NORETURN(e); /* FINISH_NORETURN, not FINISH(e, false): FINISH_LOGIC decides whether to * propagate at run time, so the compiler still parses its * `return __err_context` and diagnoses it in a function returning void, * whatever __pass_up says. An error this collector did not handle takes the * process down from here, which is right -- but note it is now the * collector's thread deciding the exit status. */ } /* * Drain exactly what the producers will send. A fixed count rather than a * sentinel: a miscounted handoff should fail the test, not hang it. */ static void collect_all(akerr_ThreadArg *arg) { int seen[AKERR_TEST_THREADS + 1] = { 0 }; int total = (AKERR_TEST_THREADS - 1) * ITERATIONS; for ( int i = 0; i < total; i++ ) { collect_one(arg, queue_pop(), seen); } } static void *handoff_body(void *raw) { akerr_ThreadArg *arg = raw; pthread_barrier_wait(arg->barrier); if ( arg->id == 1 ) { collect_all(arg); } else { for ( int i = 0; i < ITERATIONS; i++ ) { produce_one(arg, i); } } return NULL; } /* * Written by the raising thread, read by main() after pthread_join(). The join * is the happens-before edge, which is the same thing the queue's mutex does * above -- a plain global needs no atomics once something orders it. */ static akerr_ErrorContext *parked; static void *raise_and_exit(void *unused) { PREPARE_ERROR(e); (void)unused; FAIL(e, AKERR_IO, "raised on a thread that exited"); parked = e; return NULL; } int main(void) { pthread_t raiser; int failures = 0; akerr_log_method = &akerr_thread_logger; akerr_init(); AKERR_CHECK(akerr_slots_in_use() == 0); AKERR_CHECK(pthread_mutex_init(&queue.lock, NULL) == 0); AKERR_CHECK(pthread_cond_init(&queue.not_full, NULL) == 0); AKERR_CHECK(pthread_cond_init(&queue.not_empty, NULL) == 0); failures = akerr_run_threads(&handoff_body); AKERR_CHECK(failures == 0); /* Every handed-off context was released by the thread that received it. */ AKERR_CHECK(akerr_slots_in_use() == 0); /* * A context outlives the thread that raised it: the pool is process-global, * not thread-local storage. By the time these checks run, the thread that * called FAIL() no longer exists. */ AKERR_CHECK(pthread_create(&raiser, NULL, &raise_and_exit, NULL) == 0); AKERR_CHECK(pthread_join(raiser, NULL) == 0); AKERR_CHECK(parked != NULL); AKERR_CHECK(akerr_valid_error_address(parked) == 1); AKERR_CHECK(parked->status == AKERR_IO); AKERR_CHECK(parked->refcount == 1); AKERR_CHECK(strstr(parked->stacktracebuf, "raised on a thread that exited") != NULL); RELEASE_ERROR(parked); AKERR_CHECK(parked == NULL); AKERR_CHECK(akerr_slots_in_use() == 0); for ( int i = 0; i < AKERR_MAX_ARRAY_ERROR; i++ ) { AKERR_CHECK(akerr_slot_holder(i) == 0); } /* Nothing here reports through the log method: a handoff is not an error. */ AKERR_CHECK(akerr_thread_logs() == 0); pthread_cond_destroy(&queue.not_empty); pthread_cond_destroy(&queue.not_full); pthread_mutex_destroy(&queue.lock); fprintf(stderr, "err_threads_handoff ok (%d producers x %d errors)\n", AKERR_TEST_THREADS - 1, ITERATIONS); return 0; }