This is an automated email from the ASF dual-hosted git repository. asf-gitbox-commits pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/qpid-proton.git
commit 14018a5d74f48538a8df4e3e4ee3292d8dc2b123 Author: Andrew Stitcher <[email protected]> AuthorDate: Tue Oct 6 17:17:28 2026 -0400 PROTON-2985: Add a -perturb option to threaderciser Proactor races rarely reproduce under normal scheduling, because the locks around task scheduling are held only briefly. -perturb randomly yields or sleeps for up to 2ms immediately before and after this file's own calls to pn_proactor_connect(), pn_connection_wake(), pn_proactor_listen(), pn_listener_close(), and around the pn_proactor_wait()/pn_event_batch_next()/ pn_proactor_done() loop, varying when concurrent threads arrive at those entry points. It changes only the timing of the actions, never which ones are taken, and is off unless asked for, so the default test run is unaffected. The start banner reports perturb=on/off. Assisted-By: Claude Opus 5 <[email protected]> --- c/tests/threaderciser.c | 45 ++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 44 insertions(+), 1 deletion(-) diff --git a/c/tests/threaderciser.c b/c/tests/threaderciser.c index 311ca78d4..3d2f9cf3f 100644 --- a/c/tests/threaderciser.c +++ b/c/tests/threaderciser.c @@ -59,6 +59,7 @@ #include <proton/proactor.h> #include <inttypes.h> +#include <sched.h> #include <stdarg.h> #include <stdio.h> #include <stdlib.h> @@ -72,6 +73,32 @@ #define TIMEOUT_MAX 100 /* Milliseconds */ #define SLEEP_MAX 100 /* Milliseconds */ +/* + Scheduler-perturbation mode, enabled with -perturb. + + Proactor races rarely reproduce under normal scheduling, as the locks around + task scheduling are usually held only briefly. Rather than instrumenting the + proactor itself, perturb the timing of this file's own calls immediately + around the API entry points known to take those locks. Randomizing when + concurrent callers arrive at them is a cheap stand-in for full schedule + randomization. + + This changes only the timing of the actions, never which ones are taken. +*/ +static bool perturb_enabled = false; +#define PERTURB_MAX_NSEC 2000000L /* 2ms upper bound for randomized micro-sleeps */ + +static void perturb(void) { + if (!perturb_enabled) return; + /* rand() as elsewhere in this file: needs to vary interleavings, not be secure. */ + if (rand() % 2) { + sched_yield(); + } else { + struct timespec ts = { 0, rand() % PERTURB_MAX_NSEC }; + nanosleep(&ts, NULL); + } +} + /* Set of actions that can be enabled/disabled/counted */ typedef enum { A_LISTEN, A_CLOSE_LISTEN, A_CONNECT, A_CLOSE_CONNECT, A_WAKE, A_TIMEOUT, A_CANCEL_TIMEOUT } action; const char* action_name[] = { "listen", "close-listen", "connect", "close-connect", "wake", "timeout", "cancel-timeout" }; @@ -229,7 +256,9 @@ void cpool_connect(cpool *cp, pn_proactor_t *proactor, const char *addr) { connection_ctx *ctx = connection_ctx_new(); if (cpool_add(cp, ctx)) { debuga(A_CONNECT, ctx->pn_connection); + perturb(); pn_proactor_connect(proactor, ctx->pn_connection, addr); + perturb(); } else { pn_connection_free(ctx->pn_connection); /* Won't be freed by proactor */ connection_ctx_free(ctx); @@ -244,7 +273,9 @@ void cpool_wake(cpool *cp) { pthread_mutex_lock(&ctx->lock); if (ctx && ctx->pn_connection) { debuga(A_WAKE, ctx->pn_connection); + perturb(); pn_connection_wake(ctx->pn_connection); + perturb(); } pthread_mutex_unlock(&ctx->lock); cpool_unref(ctx); @@ -285,7 +316,9 @@ static void lpool_listen(lpool *lp, pn_proactor_t *proactor) { listener_ctx *ctx = listener_ctx_new(); if (lpool_add(lp, ctx)) { debuga(A_LISTEN, ctx->pn_listener); + perturb(); pn_proactor_listen(proactor, ctx->pn_listener, a, BACKLOG); + perturb(); } else { pn_listener_free(ctx->pn_listener); /* Won't be freed by proactor */ listener_ctx_free(ctx); @@ -312,7 +345,9 @@ void lpool_close(lpool *lp) { if (ctx) { pthread_mutex_lock(&ctx->lock); if (ctx->pn_listener) { + perturb(); pn_listener_close(ctx->pn_listener); + perturb(); debuga(A_CLOSE_LISTEN, ctx->pn_listener); } pthread_mutex_unlock(&ctx->lock); @@ -473,11 +508,14 @@ static void* proactor_thread(void* void_g) { global *g = (global*) void_g; bool ok = true; while (ok) { + perturb(); pn_event_batch_t *events = pn_proactor_wait(g->proactor); pn_event_t *e; while (ok && (e = pn_event_batch_next(events))) { ok = ok && handle(g, e); + perturb(); } + perturb(); pn_proactor_done(g->proactor, events); } debug("proactor_thread end"); @@ -494,6 +532,7 @@ void usage(const char **argv, const char **arg) { fprintf(stderr, " -time TIME: total run-time in seconds (default %d)\n", default_runtime); fprintf(stderr, " -threads THREADS: total number of threads (default %d)\n", default_threads); fprintf(stderr, " -debug: print debug messages\n"); + fprintf(stderr, " -perturb: randomly yield/sleep around proactor calls to vary thread interleavings\n"); fprintf(stderr, "Flags to enable specific actions (all enabled by default)\n"); fprintf(stderr, " "); for (int i = 0; i < (int)action_size; ++i) fprintf(stderr, " -%s", action_name[i]); @@ -534,6 +573,9 @@ int main(int argc, const char* argv[]) { else if (!strcmp(*arg, "-debug")) { debug_enable = true; } + else if (!strcmp(*arg, "-perturb")) { + perturb_enabled = true; + } else if (!strncmp(*arg, "-no-", 4)) { action_enabled[find_action((*arg) + 4, argv, arg)] = false; } @@ -552,7 +594,8 @@ int main(int argc, const char* argv[]) { /* Set up global state, start threads */ - printf("threaderciser start: threads=%d, time=%d, actions=[", threads, runtime); + printf("threaderciser start: threads=%d, time=%d, perturb=%s, actions=[", + threads, runtime, perturb_enabled ? "on" : "off"); bool comma = false; for (size_t i = 0; i < action_size; ++i) { if (action_enabled[i]) { --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
