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]

Reply via email to