Skip to content
Merged
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
87 changes: 68 additions & 19 deletions src/adapters/libuv.h
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,8 @@ typedef struct __natsLibuvEvent
{
int type;
bool add;
// Only meaningful for NATS_LIBUV_ATTACH.
natsSock socket;
struct __natsLibuvEvent *next;

} natsLibuvEvent;
Expand Down Expand Up @@ -107,18 +109,9 @@ natsLibuv_SetThreadLocalLoop(uv_loop_t *loop)
}

static natsStatus
uvScheduleToEventLoop(natsLibuvEvents *nle, int eventType, bool add)
uvEnqueueEvent(natsLibuvEvents *nle, natsLibuvEvent *newEvent)
{
natsLibuvEvent *newEvent = NULL;
int res;

newEvent = (natsLibuvEvent*) malloc(sizeof(natsLibuvEvent));
if (newEvent == NULL)
return NATS_NO_MEMORY;

newEvent->type = eventType;
newEvent->add = add;
newEvent->next = NULL;
int res;

uv_mutex_lock(nle->lock);

Expand All @@ -142,6 +135,50 @@ uvScheduleToEventLoop(natsLibuvEvents *nle, int eventType, bool add)
return (res == 0 ? NATS_OK : NATS_ERR);
}

static natsLibuvEvent*
uvNewEvent(int eventType, bool add)
{
// Zero-initialized so that `socket` has a defined value for the event
// types that do not carry one.
natsLibuvEvent *newEvent = (natsLibuvEvent*) calloc(1, sizeof(natsLibuvEvent));

if (newEvent == NULL)
return NULL;

newEvent->type = eventType;
newEvent->add = add;
newEvent->next = NULL;

return newEvent;
}

static natsStatus
uvScheduleToEventLoop(natsLibuvEvents *nle, int eventType, bool add)
{
natsLibuvEvent *newEvent = uvNewEvent(eventType, add);

if (newEvent == NULL)
return NATS_NO_MEMORY;

return uvEnqueueEvent(nle, newEvent);
}

// The socket travels in the event so that `nle` is not written here. Callers
// pass their own socket argument; reading `nle->socket` off the event loop
// thread is not allowed.
static natsStatus
uvScheduleAttachToEventLoop(natsLibuvEvents *nle, natsSock socket)
{
natsLibuvEvent *newEvent = uvNewEvent(NATS_LIBUV_ATTACH, true);

if (newEvent == NULL)
return NATS_NO_MEMORY;

newEvent->socket = socket;

return uvEnqueueEvent(nle, newEvent);
}

static void
natsLibuvPoll(uv_poll_t* handle, int status, int events)
{
Expand Down Expand Up @@ -187,6 +224,14 @@ uvPollUpdate(natsLibuvEvents *nle, int eventType, bool add)
nle->events &= ~UV_WRITABLE;
}

// The poll handle is released as soon as both events have been removed, and further
// events can still be delivered for it: the library asks to stop polling once per
// reconnect attempt and once more when it gives up and closes the connection, and the
// removals of the read and of the write event are two separate events. There is nothing
// left to poll on or to close, so both calls below would dereference a NULL handle.
if (nle->handle == NULL)
return NATS_OK;

if (nle->events)
{
int res = uv_poll_start(nle->handle, nle->events, natsLibuvPoll);
Expand All @@ -204,11 +249,16 @@ uvPollUpdate(natsLibuvEvents *nle, int eventType, bool add)
return NATS_OK;
}

// Runs on the event loop thread only, which is what makes `nle->socket` and
// `nle->events` single-threaded.
static natsStatus
uvAsyncAttach(natsLibuvEvents *nle)
uvAsyncAttach(natsLibuvEvents *nle, natsSock socket)
{
natsStatus s = NATS_OK;

nle->socket = socket;
nle->events = UV_READABLE;

// Even when this is a reconnect, previous nle->handle has already been
// set to NULL (and the memory has or will be freed in uvHandleClosedCb),
// so recreate now.
Expand Down Expand Up @@ -303,7 +353,7 @@ uvAsyncCb(uv_async_t *handle)
{
case NATS_LIBUV_ATTACH:
{
s = uvAsyncAttach(nle);
s = uvAsyncAttach(nle, event->socket);
break;
}
case NATS_LIBUV_READ:
Expand Down Expand Up @@ -404,13 +454,12 @@ natsLibuv_Attach(void **userData, void *loop, natsConnection *nc, natsSock socke

if (s == NATS_OK)
{
nle->socket = socket;
nle->events = UV_READABLE;

if (sched)
s = uvScheduleToEventLoop(nle, NATS_LIBUV_ATTACH, true);
// See comment in natsLibuvRead. Ordering matters here too: a pending
// removal must reach uvPollUpdate before this socket is installed.
if (sched || (nle->head != NULL))
s = uvScheduleAttachToEventLoop(nle, socket);
else
s = uvAsyncAttach(nle);
s = uvAsyncAttach(nle, socket);
}

if (s == NATS_OK)
Expand Down
9 changes: 7 additions & 2 deletions src/conn.c
Original file line number Diff line number Diff line change
Expand Up @@ -2846,8 +2846,13 @@ _close(natsConnection *nc, natsConnStatus status, bool fromPublicClose, bool doC
// one doing it.
if (ttj.readLoop == NULL)
{
// If event loop attached, stop polling...
if (nc->el.attached)
// If event loop attached and still polling, stop polling...
// `_evStopPolling()` has already run for a connection which was reconnecting,
// and it leaves the actual socket close to the event loop adapter, which has
// released its handle by then. Asking a second time would leave the socket of
// the last reconnect attempt open forever, because the adapter only knows the
// socket it was given when it was attached.
if (nc->el.attached && nc->sockCtx.useEventLoop)
{
// This will take care of invalidating the socket and clear SSL,
// but the actual socket close will be done from the event loop
Expand Down
14 changes: 11 additions & 3 deletions src/unix/sock.c
Original file line number Diff line number Diff line change
Expand Up @@ -48,10 +48,18 @@ natsSock_WaitReady(int waitMode, natsSockCtx *ctx)
abort();
}

if (deadline != NULL)
timeout = natsDeadline_GetTimeout(deadline);
do
{
if (deadline != NULL)
timeout = natsDeadline_GetTimeout(deadline);

res = poll(&pfd, 1, timeout);
}
// A signal delivered to the thread (e.g. the sampling profiler of the host
// application) interrupts poll() with EINTR. That is not a socket error:
// wait again for whatever is left of the deadline.
while ((res == NATS_SOCK_ERROR) && (NATS_SOCK_GET_ERROR == EINTR));

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