Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 14 additions & 3 deletions src/unix/sock.c
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
126 changes: 126 additions & 0 deletions test/test.c
Original file line number Diff line number Diff line change
Expand Up @@ -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];
Expand Down Expand Up @@ -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) {}
Expand All @@ -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) {}

Expand Down