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) {}