Skip to content
Open
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
21 changes: 20 additions & 1 deletion src/adapters/libevent.h
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ typedef struct
struct event *read;
struct event *write;
struct event *keepActive;
// This object owns the reference the library takes on `nc` at the first attach.
bool releaseConnOnFree;

} natsLibeventEvents;

Expand Down Expand Up @@ -100,6 +102,7 @@ natsLibevent_Attach(void **userData, void *loop, natsConnection *nc, natsSock so
struct event_base *libeventLoop = (struct event_base*) loop;
natsLibeventEvents *nle = (natsLibeventEvents*) (*userData);
natsStatus s = NATS_OK;
bool created = false;

// This is the first attach (when reconnecting, nle will be non-NULL).
if (nle == NULL)
Expand All @@ -108,6 +111,8 @@ natsLibevent_Attach(void **userData, void *loop, natsConnection *nc, natsSock so
if (nle == NULL)
return NATS_NO_MEMORY;

// Indicate that we have created the object here (in case we get a failure).
created = true;
nle->nc = nc;
nle->loop = libeventLoop;

Expand Down Expand Up @@ -151,9 +156,19 @@ natsLibevent_Attach(void **userData, void *loop, natsConnection *nc, natsSock so
}

if (s == NATS_OK)
{
// The library retains only at a first attach, so only a first attach owns it.
if (created)
nle->releaseConnOnFree = true;

*userData = (void*) nle;
else
}
else if (created)
{
// A failure on a successive attach must leave `nle` untouched: the library
// keeps it in nc->el.data and reuses it on the next reconnect attempt.
natsLibevent_Detach((void*) nle);
}

return s;
}
Expand Down Expand Up @@ -245,6 +260,10 @@ natsLibevent_Detach(void *userData)
event_free(nle->keepActive);
}

// Reads nle->nc, so it has to happen before the free below.
if (nle->releaseConnOnFree)
natsConnection_ProcessDetachedEvent(nle->nc);

free(nle);

return NATS_OK;
Expand Down
97 changes: 84 additions & 13 deletions src/adapters/libuv.h
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,8 @@ typedef struct
uv_mutex_t *lock;
natsLibuvEvent *head;
natsLibuvEvent *tail;
// This object owns the reference the library takes on `nc` at the first attach.
bool releaseConnOnFree;

} natsLibuvEvents;

Expand Down Expand Up @@ -262,7 +264,9 @@ uvAsyncAttach(natsLibuvEvents *nle, natsSock socket)
// 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.
nle->handle = (uv_poll_t*) malloc(sizeof(uv_poll_t));
// Zero-initialized so that `type` stays UV_UNKNOWN_HANDLE until uv_poll_init
// links the handle into the loop, which is what the teardown below asks.
nle->handle = (uv_poll_t*) calloc(1, sizeof(uv_poll_t));
if (nle->handle == NULL)
s = NATS_NO_MEMORY;

Expand All @@ -283,13 +287,25 @@ uvAsyncAttach(natsLibuvEvents *nle, natsSock socket)
s = NATS_ERR;
}

if ((s != NATS_OK) && (nle->handle != NULL))
{
// uv_poll_init sets `type` and links the handle in one step but can fail on either
// side of it, so `type` is the exact witness: what the loop never saw is freed, what
// it linked is closed, which is safe mid-init while no poll request is submitted.
if (nle->handle->type == UV_UNKNOWN_HANDLE)
free(nle->handle);
else
uv_close((uv_handle_t*) nle->handle, uvHandleClosedCb);

nle->handle = NULL;
}

return s;
}

static void
uvFinalCloseCb(uv_handle_t* handle)
natsLibuvEvents_free(natsLibuvEvents *nle, bool releaseConn)
{
natsLibuvEvents *nle = (natsLibuvEvents*) handle->data;
natsLibuvEvent *event;

while ((event = nle->head) != NULL)
Expand All @@ -298,14 +314,38 @@ uvFinalCloseCb(uv_handle_t* handle)
free(event);
}
free(nle->scheduler);
uv_mutex_destroy(nle->lock);
free(nle->lock);
if (nle->lock != NULL)
{
uv_mutex_destroy(nle->lock);
free(nle->lock);
}
// Reads nle->nc, so it has to happen before the free below.
if (releaseConn)
natsConnection_ProcessDetachedEvent(nle->nc);
free(nle);
}

static void
uvFinalCloseCb(uv_handle_t* handle)
{
natsLibuvEvents *nle = (natsLibuvEvents*) handle->data;

natsLibuvEvents_free(nle, nle->releaseConnOnFree);
}

static void
uvAsyncDetach(natsLibuvEvents *nle)
{
// The library asks to stop polling before it detaches, but that request can fail
// to reach this loop and its failure is not propagated, so a handle still armed
// here is this adapter's to close: the free below invalidates the `nle` that
// handle's `data` points at.
if (nle->handle != NULL)
{
uv_close((uv_handle_t*) nle->handle, uvHandleClosedCb);
nle->handle = NULL;
}

uv_close((uv_handle_t*) nle->scheduler, uvFinalCloseCb);
}

Expand Down Expand Up @@ -396,6 +436,8 @@ natsLibuv_Attach(void **userData, void *loop, natsConnection *nc, natsSock socke
bool sched = false;
natsLibuvEvents *nle = (natsLibuvEvents*) (*userData);
natsStatus s = NATS_OK;
bool created = false;
bool schedulerInitialized = false;

sched = ((uv_key_get(&uvLoopThreadKey) != loop) ? true : false);

Expand All @@ -410,30 +452,45 @@ natsLibuv_Attach(void **userData, void *loop, natsConnection *nc, natsSock socke
if (nle == NULL)
return NATS_NO_MEMORY;

// Indicate that we have created the object here (in case we get a failure).
created = true;

nle->lock = (uv_mutex_t*) malloc(sizeof(uv_mutex_t));
if (nle->lock == NULL)
s = NATS_NO_MEMORY;

if ((s == NATS_OK) && (uv_mutex_init(nle->lock) != 0))
{
// A non-NULL nle->lock has to mean an initialized mutex:
// uv_mutex_destroy aborts on one that was never initialized.
free(nle->lock);
nle->lock = NULL;
s = NATS_ERR;
}

if ((s == NATS_OK)
&& ((nle->scheduler = (uv_async_t*) malloc(sizeof(uv_async_t))) == NULL))
{
s = NATS_NO_MEMORY;
}

if ((s == NATS_OK)
&& (uv_async_init(uvLoop, nle->scheduler, uvAsyncCb) != 0))
if (s == NATS_OK)
{
s = NATS_ERR;
if (uv_async_init(uvLoop, nle->scheduler, uvAsyncCb) != 0)
s = NATS_ERR;
else
{
// uv_async_init links the handle into the loop's handle and async
// queues, so from here on it has to be closed, not freed.
schedulerInitialized = true;
nle->scheduler->data = (void*) nle;
}
}

if (s == NATS_OK)
{
nle->nc = nc;
nle->loop = uvLoop;
nle->scheduler->data = (void*) nle;
nle->nc = nc;
nle->loop = uvLoop;
}
}

Expand All @@ -448,9 +505,23 @@ natsLibuv_Attach(void **userData, void *loop, natsConnection *nc, natsSock socke
}

if (s == NATS_OK)
{
// Gated on a first attach, which cannot run off the loop thread, so this
// field is written and read (uvFinalCloseCb) on that thread only.
if (created)
nle->releaseConnOnFree = true;

*userData = (void*) nle;
else
natsLibuv_Detach((void*) nle);
}
else if (created)
{
// A failure on a successive attach must leave `nle` untouched: the library
// keeps it in nc->el.data and reuses it on the next reconnect attempt.
if (schedulerInitialized)
uv_close((uv_handle_t*) nle->scheduler, uvFinalCloseCb);
else
natsLibuvEvents_free(nle, false);
}

return s;
}
Expand Down
20 changes: 16 additions & 4 deletions src/conn.c
Original file line number Diff line number Diff line change
Expand Up @@ -2004,18 +2004,28 @@ _processConnInit(natsConnection *nc)
// event just after this call returns.
nc->sockCtx.useEventLoop = true;

// For the very first attach, we will retain the connection.
if (!nc->el.retained)
_retain(nc);

s = nc->opts->evCbs.attach(&(nc->el.data),
nc->opts->evLoop,
nc,
(int) nc->sockCtx.fd);
if (s == NATS_OK)
{
nc->el.attached = true;
nc->el.retained = true;
}
else
{
nc->sockCtx.useEventLoop = false;

// If this was the very first attach and we failed, release the connection
// to compensate for the retain above.
if (!nc->el.retained)
_release(nc);

nats_setError(s,
"Error attaching to the event loop: %d - %s",
s, natsStatus_GetText(s));
Expand Down Expand Up @@ -4134,8 +4144,6 @@ natsConnection_ProcessReadEvent(natsConnection *nc)
}
}

_retain(nc);

buffer = nc->el.buffer;
size = nc->opts->ioBufSize;

Expand All @@ -4154,8 +4162,6 @@ natsConnection_ProcessReadEvent(natsConnection *nc)

if (s != NATS_OK)
_processOpError(nc, s, false);

natsConn_release(nc);
}

void
Expand Down Expand Up @@ -4215,6 +4221,12 @@ natsConnection_ProcessCloseEvent(natsSock *socket)
*socket = NATS_SOCK_INVALID;
}

void
natsConnection_ProcessDetachedEvent(natsConnection *nc)
{
natsConn_release(nc);
}

natsStatus
natsConnection_GetClientID(natsConnection *nc, uint64_t *cid)
{
Expand Down
13 changes: 13 additions & 0 deletions src/nats.h
Original file line number Diff line number Diff line change
Expand Up @@ -4221,6 +4221,19 @@ natsConnection_ProcessCloseEvent(natsSock *socket);
NATS_EXTERN void
natsConnection_ProcessWriteEvent(natsConnection *nc);

/** \brief Process a detach event when using external event loop.
*
* When a connection is closed, the library will invoke the adapter's
* #natsEvLoop_Detach callback. But the code in the adapter may run
* asynchronously. The adapter will invoke this function when the
* adapter has fully detached the NATS connection from the event loop,
* so that resources held by the library for this connection can be released.
*
* @param nc the pointer to the #natsConnection object.
*/
NATS_EXTERN void
natsConnection_ProcessDetachedEvent(natsConnection *nc);

/** \brief Connects to a `NATS Server` using any of the URL from the given list.
*
* Attempts to connect to a `NATS Server`.
Expand Down
1 change: 1 addition & 0 deletions src/natsp.h
Original file line number Diff line number Diff line change
Expand Up @@ -756,6 +756,7 @@ struct __natsConnection
{
bool attached;
bool writeAdded;
bool retained; // Will be set to true at the very first successful attach.
void *buffer;
void *data;
} el;
Expand Down
1 change: 1 addition & 0 deletions test/list_test.txt
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@ _test(DrainSubStops)
_test(ErrOnConnectAndDeadlock)
_test(ErrOnMaxPayloadLimit)
_test(EventLoop)
_test(EventLoopDestroyWhileAttached)
_test(EventLoopParserResetOnDisconnect)
_test(EventLoopRetryOnFailedConnect)
_test(EventLoopTLS)
Expand Down
Loading