From 3472e1a6f003341917315f61c5fed027023ea6b8 Mon Sep 17 00:00:00 2001 From: Groene AI <270696204+groeneai@users.noreply.github.com> Date: Mon, 14 Sep 2026 23:05:22 +0000 Subject: [PATCH] Retry the socket wait when poll() is interrupted by a signal `natsSock_WaitReady` reported every `poll()` failure as `NATS_IO_ERROR`, including `EINTR`. `poll` is never restarted after a signal handler runs, whatever `SA_RESTART` says (`signal(7)`), so `EINTR` there says nothing about the socket: the wait was simply cut short. Reporting it as an I/O error makes `natsSock_ConnectTcp` close the fd and move to the next `addrinfo`, so a single asynchronous signal delivered during a connect or a handshake read fails the whole connection attempt. Any embedder that delivers periodic per-thread signals reaches this. ClickHouse does so by default: it arms a 10 s profiler timer (SIGUSR1 with `SIGEV_THREAD_ID`) on every thread it takes from its global thread pool, with the first fire placed at a uniformly random point inside that first period so that short-lived work is still sampled, and it runs the NATS event loop on such a thread while posting the connect onto it milliseconds later. Its `NATS` table engine then fails to be created with Cannot connect to Nats last error: (unix/sock.c:56): poll error: 4 where `errno 4` is `EINTR`. ClickHouse CI hit that leaf three times inside 3 h 35 min on 2026-09-14, on three unrelated carriers and three different build flavours: on master in `Integration tests (amd_asan_ubsan, db disk, old analyzer, 4/8)` at 19:30:23Z, and on two unrelated pull requests in `Integration tests (arm_binary, distributed plan, 1/4)` at 20:53:58Z and in `Integration tests (amd_llvm_coverage, 8/8)` at 23:05:25Z. The third check is green only because the integration runner's retry passed; its recorded context still carries the leaf. Both test functions of a new integration module are represented. `arm_binary` is a plain aarch64 build with no sanitizer, so the defect is neither sanitizer- nor x86-specific. That module sets the engine's connect attempts to 1, which is what turns one interrupted `poll` into a failed DDL rather than a retried connect. The `poll` now sits in a loop that recomputes `natsDeadline_GetTimeout(deadline)` on each iteration and continues when `poll` fails with `EINTR`. Recomputing is what bounds the loop: `natsDeadline_GetTimeout` returns -1 for an inactive deadline, which is the requested wait-forever semantics, and otherwise the remaining milliseconds clamped at 0, so an expired deadline gives `poll` a timeout of 0, `poll` returns 0, and the pre-existing `NATS_TIMEOUT` arm fires exactly as before. Re-arming the full timeout instead would be the one way to get this wrong, which is why the new test arms bound the elapsed time as well as the status. The three result arms, the `pfd` setup and the `waitMode` switch are unchanged, so nothing changes on any path that does not see `EINTR`, and there is no new symbol, no signature change and no lock. This is the only `poll` or `select` in the non-Windows sources, and it is reached both for connect completion and for the handshake read, so the two call shapes are covered by the one change. `natsSock_Read` and `natsSock_Write` report an `EINTR` from `recv`/`send` the same fatal way and are left alone as unreproduced: their SSL arms cannot see one, because the TLS blocking window from `natsSock_SetBlocking(fd, true)` in `_makeTLSConn` to the matching restore around `SSL_do_handshake` contains no such call, and their plain arms see a blocking fd only without an external event loop (`_processConnInit`, `conn.c:1983`, when `opts->writeDeadline <= 0`, after which `_spinUpSocketWatchers` runs `_readLoop` on it), which is not a mode ClickHouse uses, since it always sets an event loop and `conn.c:1992` then restores non-blocking. `SSL_do_handshake` itself can report a signal-interrupted blocking handshake, but it is a different call site with a different retry contract, so it too is deliberately left alone. Test: two arms in the existing `test_natsWaitReady`, both run under a storm thread that `pthread_kill`s SIGALRM at the waiting thread about every millisecond, and both asserting a receipt counter the handler increments, so that a failed `sigaction`, a failed mutex or thread create, or undelivered signals fail the case instead of quietly degrading it to an unsignalled wait. The first arm wraps the existing no-deadline case, where the wait survives some 460 delivered signals and still returns `NATS_OK` when the fake server's byte arrives, inside the same 450-600 ms bound; that is what pins progress. The second keeps a 50 ms deadline and asserts `NATS_TIMEOUT` inside 40-100 ms under some 46 signals; that is what pins the deadline being recomputed rather than re-armed. The stop flag is guarded by a `natsMutex` instead of being a plain `volatile`, because the suite's own test properties set `TSAN_OPTIONS=...:halt_on_error=1` whenever `NATS_SANITIZE` is on, and an unsynchronised flag aborts this very case under `-fsanitize=thread`. The storm also stops by itself after 1500 ms, far past both arms' bounds, so that an implementation re-arming the full 50 ms on each retry fails the second arm's duration bound at about 1.55 s rather than stalling in the wait until the fake server tears the connection down. Measured on this branch: both arms fail before this change with the `poll error: 4` message above and pass 10 of 10 after; with the `pthread_kill` call removed both fail on the receipt counter alone; with the deadline recompute hoisted back out of the loop the second arm fails on duration at 1549 ms; and `natsWaitReady` is clean under `-fsanitize=thread`, where the unsynchronised flag is reported as a data race. Delivery is targeted at one thread, so no helper thread's `nats_Sleep` is cut short; the arms are guarded for non-Windows, and the no-op handler is left installed rather than restored, so a signal still in flight cannot terminate the process. The defect is verbatim in this fork's `ClickHouse/v3.13.0` branch and in upstream `nats-io/nats.c` `main`. CI reports: https://s3.amazonaws.com/clickhouse-test-reports/praktika.html?REF=master&sha=3fba61b4895078ed184a939a9012caba51022398&name_0=MasterCI&name_1=Integration%20tests%20%28amd_asan_ubsan%2C%20db%20disk%2C%20old%20analyzer%2C%204%2F8%29 https://s3.amazonaws.com/clickhouse-test-reports/praktika.html?PR=116234&sha=a9aa67f7a30ee303d3f001a812db3f712411414e&name_0=PR&name_1=Integration%20tests%20%28arm_binary%2C%20distributed%20plan%2C%201%2F4%29 https://s3.amazonaws.com/clickhouse-test-reports/praktika.html?PR=96844&sha=4631dd43b10e7028e31e10a26564556c2ec95aa0&name_0=PR&name_1=Integration%20tests%20%28amd_llvm_coverage%2C%208%2F8%29 Co-Authored-By: Claude Opus 5 --- src/unix/sock.c | 17 +++++-- test/test.c | 126 ++++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 140 insertions(+), 3 deletions(-) diff --git a/src/unix/sock.c b/src/unix/sock.c index 4bfcfa39c..cca4032a6 100644 --- a/src/unix/sock.c +++ b/src/unix/sock.c @@ -48,10 +48,21 @@ natsSock_WaitReady(int waitMode, natsSockCtx *ctx) abort(); } - if (deadline != NULL) - timeout = natsDeadline_GetTimeout(deadline); + for (;;) + { + if (deadline != NULL) + timeout = natsDeadline_GetTimeout(deadline); + + res = poll(&pfd, 1, timeout); + + // poll() is not restarted after a signal handler runs, so EINTR says nothing + // about the socket. Continue the wait with what is left of the deadline. + if ((res == NATS_SOCK_ERROR) && (NATS_SOCK_GET_ERROR == EINTR)) + continue; + + break; + } - res = poll(&pfd, 1, timeout); if (res == NATS_SOCK_ERROR) return nats_setError(NATS_IO_ERROR, "poll error: %d", NATS_SOCK_GET_ERROR); else if (res == 0) diff --git a/test/test.c b/test/test.c index d55d33ec6..fb70e5059 100644 --- a/test/test.c +++ b/test/test.c @@ -4890,11 +4890,105 @@ _testSockShutdownThread(void *closure) natsSock_Shutdown(ctx->fd); } +#ifndef _WIN32 +// The storm stops on its own after this long, well past every arm's bound, so an +// implementation that keeps waiting fails an arm instead of stalling inside it. +#define _TEST_STORM_MAX_MS (1500) + +static volatile sig_atomic_t _testSignalsReceived = 0; +static natsMutex *_testStormMu = NULL; +static bool _testStormStop = false; +static pthread_t _testStormTarget; + +static void +_testNoOpSignalHandler(int sig) +{ + (void) sig; + _testSignalsReceived++; +} + +static void +_testSignalStormThread(void *closure) +{ + int64_t stormEnd = nats_Now() + _TEST_STORM_MAX_MS; + bool stop = false; + + (void) closure; + + while (!stop && (nats_Now() < stormEnd)) + { + pthread_kill(_testStormTarget, SIGALRM); + nats_Sleep(1); + + natsMutex_Lock(_testStormMu); + stop = _testStormStop; + natsMutex_Unlock(_testStormMu); + } +} + +// Delivers SIGALRM to the CALLING thread about every millisecond until stopped. +// Delivery is targeted, so no other thread's nats_Sleep() is cut short by it. +static natsStatus +_testStartSignalStorm(natsThread **storm) +{ + natsStatus s; + struct sigaction sa; + + _testSignalsReceived = 0; + _testStormTarget = pthread_self(); + + s = natsMutex_Create(&_testStormMu); + if (s != NATS_OK) + return s; + + natsMutex_Lock(_testStormMu); + _testStormStop = false; + natsMutex_Unlock(_testStormMu); + + memset(&sa, 0, sizeof(sa)); + // A handler that does nothing is enough: poll() is not restarted after any + // handler has run, whatever sa_flags says. + sa.sa_handler = _testNoOpSignalHandler; + sigemptyset(&sa.sa_mask); + sa.sa_flags = 0; + if (sigaction(SIGALRM, &sa, NULL) != 0) + return NATS_SYS_ERROR; + + return natsThread_Create(storm, _testSignalStormThread, NULL); +} + +static void +_testStopSignalStorm(natsThread **storm) +{ + // NULL when natsMutex_Create() itself failed, in which case no storm thread + // was started either and the arm has already reported that status. + if (_testStormMu != NULL) + { + natsMutex_Lock(_testStormMu); + _testStormStop = true; + natsMutex_Unlock(_testStormMu); + } + if (*storm != NULL) + { + natsThread_Join(*storm); + natsThread_Destroy(*storm); + *storm = NULL; + } + natsMutex_Destroy(_testStormMu); + _testStormMu = NULL; + // The handler stays installed on purpose: restoring SIG_DFL here would let a + // SIGALRM still in flight terminate the process. +} +#endif + void test_natsWaitReady(void) { natsStatus s = NATS_OK; natsThread *t = NULL; natsThread *t2 = NULL; +#ifndef _WIN32 + natsThread *storm = NULL; +#endif natsSockCtx ctx; int64_t start, dur; char buffer[1]; @@ -4932,12 +5026,26 @@ void test_natsWaitReady(void) // Ensure that we get a would_block on read.. while (recv(ctx.fd, buffer, 1, 0) != -1) {} +#ifndef _WIN32 + test("WaitReady no deadline while signals are delivered: "); + natsSock_ClearDeadline(&ctx); + s = _testStartSignalStorm(&storm); + start = nats_Now(); + IFOK(s, natsSock_WaitReady(WAIT_FOR_READ, &ctx)); + dur = nats_Now()-start; + _testStopSignalStorm(&storm); + // An interrupted wait must keep waiting: the socket only becomes readable + // after ~500ms. The receipt count is what makes this arm about signals. + testCond((s == NATS_OK) && (dur >= 450) && (dur <= 600) + && (_testSignalsReceived > 0)); +#else test("WaitReady no deadline: "); natsSock_ClearDeadline(&ctx); start = nats_Now(); s = natsSock_WaitReady(WAIT_FOR_READ, &ctx); dur = nats_Now()-start; testCond((s == NATS_OK) && (dur >= 450) && (dur <= 600)); +#endif // Ensure that we get a would_block on read.. while (recv(ctx.fd, buffer, 1, 0) != -1) {} @@ -4949,6 +5057,24 @@ void test_natsWaitReady(void) dur = nats_Now()-start; testCond((s == NATS_TIMEOUT) && (dur >= 40) && (dur <= 100)); +#ifndef _WIN32 + // Ensure that we get a would_block on read.. + while (recv(ctx.fd, buffer, 1, 0) != -1) {} + + test("WaitReady deadline timeout while signals are delivered: "); + s = _testStartSignalStorm(&storm); + natsSock_InitDeadline(&ctx, 50); + start = nats_Now(); + IFOK(s, natsSock_WaitReady(WAIT_FOR_READ, &ctx)); + dur = nats_Now()-start; + _testStopSignalStorm(&storm); + // The deadline is what bounds the wait, so it has to be recomputed on each + // retry: a retry re-arming the full 50ms instead runs until the storm stops, + // failing this bound at about 1.55s. + testCond((s == NATS_TIMEOUT) && (dur >= 40) && (dur <= 100) + && (_testSignalsReceived > 0)); +#endif + // Ensure that we get a would_block on read.. while (recv(ctx.fd, buffer, 1, 0) != -1) {}