diff --git a/src/adapters/libevent.h b/src/adapters/libevent.h index 6952bb44c..ad41ae6de 100644 --- a/src/adapters/libevent.h +++ b/src/adapters/libevent.h @@ -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; @@ -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) @@ -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; @@ -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; } @@ -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; diff --git a/src/adapters/libuv.h b/src/adapters/libuv.h index 3de674c84..9583ef9bb 100644 --- a/src/adapters/libuv.h +++ b/src/adapters/libuv.h @@ -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; @@ -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; @@ -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) @@ -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); } @@ -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); @@ -410,12 +452,21 @@ 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)) @@ -423,17 +474,23 @@ natsLibuv_Attach(void **userData, void *loop, natsConnection *nc, natsSock socke 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; } } @@ -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; } diff --git a/src/conn.c b/src/conn.c index 13218817a..dae68df11 100644 --- a/src/conn.c +++ b/src/conn.c @@ -2004,6 +2004,10 @@ _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, @@ -2011,11 +2015,17 @@ _processConnInit(natsConnection *nc) 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)); @@ -4134,8 +4144,6 @@ natsConnection_ProcessReadEvent(natsConnection *nc) } } - _retain(nc); - buffer = nc->el.buffer; size = nc->opts->ioBufSize; @@ -4154,8 +4162,6 @@ natsConnection_ProcessReadEvent(natsConnection *nc) if (s != NATS_OK) _processOpError(nc, s, false); - - natsConn_release(nc); } void @@ -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) { diff --git a/src/nats.h b/src/nats.h index c3bfbb230..7e21ec6e2 100644 --- a/src/nats.h +++ b/src/nats.h @@ -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`. diff --git a/src/natsp.h b/src/natsp.h index 41ebf1570..10ba8c9c2 100644 --- a/src/natsp.h +++ b/src/natsp.h @@ -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; diff --git a/test/list_test.txt b/test/list_test.txt index 58ab69766..49fc98d19 100644 --- a/test/list_test.txt +++ b/test/list_test.txt @@ -51,6 +51,7 @@ _test(DrainSubStops) _test(ErrOnConnectAndDeadlock) _test(ErrOnMaxPayloadLimit) _test(EventLoop) +_test(EventLoopDestroyWhileAttached) _test(EventLoopParserResetOnDisconnect) _test(EventLoopRetryOnFailedConnect) _test(EventLoopTLS) diff --git a/test/test.c b/test/test.c index aa0c006a3..53b3edde5 100644 --- a/test/test.c +++ b/test/test.c @@ -20382,6 +20382,27 @@ _evLoopWrite(void *userData, bool add) static natsStatus _evLoopDetach(void *userData) +{ + struct threadArg *arg = (struct threadArg *) userData; + natsConnection *nc = NULL; + + natsMutex_Lock(arg->m); + nc = arg->nc; + natsMutex_Unlock(arg->m); + natsConnection_Destroy(nc); + + natsMutex_Lock(arg->m); + arg->detached++; + natsCondition_Broadcast(arg->c); + natsMutex_Unlock(arg->m); + + return NATS_OK; +} + +// Same as _evLoopDetach, but leaves the release to the test, which is what lets +// it observe the connection after the user has destroyed it. +static natsStatus +_evLoopDetachNoRelease(void *userData) { struct threadArg *arg = (struct threadArg *) userData; @@ -20468,6 +20489,8 @@ void test_EventLoop(void) IFOK(s, natsOptions_SetClosedCB(opts, _closedCb, (void*) &arg)); testCond(s == NATS_OK); + arg.nc = nc; + pid = _startServer("nats://127.0.0.1:4222", NULL, true); CHECK_SERVER_STARTED(pid); @@ -20668,6 +20691,8 @@ void test_EventLoopParserResetOnDisconnect(void) IFOK(s, natsOptions_SetClosedCB(opts, _closedCb, (void*) &arg)); testCond(s == NATS_OK); + arg.nc = nc; + test("Start mockup server: "); if (s == NATS_OK) { @@ -20750,6 +20775,177 @@ void test_EventLoopParserResetOnDisconnect(void) _destroyDefaultThreadArgs(&arg); } +static void +_evLoopKeepOpenMockupServerThread(void *closure) +{ + natsStatus s = NATS_OK; + natsSock sock = NATS_SOCK_INVALID; + struct threadArg *arg = (struct threadArg*) closure; + natsSockCtx ctx; + char buffer[1024]; + + memset(&ctx, 0, sizeof(ctx)); + ctx.fd = NATS_SOCK_INVALID; + + s = _startMockupServer(&sock, "localhost", "4222"); + natsMutex_Lock(arg->m); + arg->status = s; + natsCondition_Signal(arg->c); + natsMutex_Unlock(arg->m); + + if ((s == NATS_OK) + && (((ctx.fd = accept(sock, NULL, NULL)) == NATS_SOCK_INVALID) + || (natsSock_SetCommonTcpOptions(ctx.fd) != NATS_OK))) + { + s = NATS_SYS_ERROR; + } + + if (s == NATS_OK) + { + s = natsSock_WriteFully(&ctx, arg->string, (int) strlen(arg->string)); + // natsSock_ReadLine keeps the bytes after the line it returns, so the + // buffer is cleared once, before the first read. + buffer[0] = '\0'; + // CONNECT, then PING. + IFOK(s, natsSock_ReadLine(&ctx, buffer, sizeof(buffer))); + IFOK(s, natsSock_ReadLine(&ctx, buffer, sizeof(buffer))); + IFOK(s, natsSock_WriteFully(&ctx, _PONG_PROTO_, _PONG_PROTO_LEN_)); + } + + // The handshake above is served by blocking reads inside the connect call, so + // a PING is what forces at least one read event through the event loop. + if (s == NATS_OK) + { + s = natsSock_WriteFully(&ctx, _PING_PROTO_, _PING_PROTO_LEN_); + IFOK(s, natsSock_ReadLine(&ctx, buffer, sizeof(buffer))); + if ((s == NATS_OK) && (strncmp(buffer, "PONG", 4) == 0)) + { + natsMutex_Lock(arg->m); + arg->msgReceived = true; + natsCondition_Broadcast(arg->c); + natsMutex_Unlock(arg->m); + } + } + + natsMutex_Lock(arg->m); + while ((s != NATS_TIMEOUT) && !(arg->done)) + s = natsCondition_TimedWait(arg->c, arg->m, 10000); + natsMutex_Unlock(arg->m); + + natsSock_Close(ctx.fd); + natsSock_Close(sock); +} + +void test_EventLoopDestroyWhileAttached(void) +{ + natsStatus s; + natsConnection *nc = NULL; + natsOptions *opts = NULL; + natsThread *t = NULL; + struct threadArg arg; + struct threadArg sarg; + + test("Set options: "); + s = _createDefaultThreadArgsForCbTests(&arg); + IFOK(s, _createDefaultThreadArgsForCbTests(&sarg)); + IFOK(s, natsOptions_Create(&opts)); + IFOK(s, natsOptions_SetURL(opts, "nats://localhost:4222")); + IFOK(s, natsOptions_SetEventLoop(opts, (void*) &arg, + _evLoopAttach, + _evLoopRead, + _evLoopWrite, + _evLoopDetachNoRelease)); + testCond(s == NATS_OK); + + arg.nc = nc; + + test("Start mockup server: "); + if (s == NATS_OK) + { + // Set to error, the mockup server thread sets it to OK once listening. + sarg.status = NATS_ERR; + sarg.string = "INFO {\"server_id\":\"22\",\"version\":\"latest\",\"go\":\"latest\",\"port\":4222,\"max_payload\":1048576}\r\n"; + s = natsThread_Create(&t, _evLoopKeepOpenMockupServerThread, (void*) &sarg); + } + if (s == NATS_OK) + { + natsMutex_Lock(sarg.m); + while ((s != NATS_TIMEOUT) && (sarg.status != NATS_OK)) + s = natsCondition_TimedWait(sarg.c, sarg.m, 2000); + IFOK(s, sarg.status); + natsMutex_Unlock(sarg.m); + } + testCond(s == NATS_OK); + + test("Start event loop: "); + natsMutex_Lock(arg.m); + arg.sock = NATS_SOCK_INVALID; + natsMutex_Unlock(arg.m); + s = natsThread_Create(&arg.t, _eventLoop, (void*) &arg); + testCond(s == NATS_OK); + + test("Connect and attach: "); + s = natsConnection_Connect(&nc, opts); + if (s == NATS_OK) + { + natsMutex_Lock(arg.m); + while ((s != NATS_TIMEOUT) && (arg.attached == 0)) + s = natsCondition_TimedWait(arg.c, arg.m, 2000); + natsMutex_Unlock(arg.m); + } + testCond((s == NATS_OK) && (nc != NULL)); + + // Positive control: without a read event having gone through the event loop, + // the checks below would pass for the wrong reason. + test("A read event went through the event loop: "); + natsMutex_Lock(sarg.m); + while ((s != NATS_TIMEOUT) && !sarg.msgReceived) + s = natsCondition_TimedWait(sarg.c, sarg.m, 5000); + natsMutex_Unlock(sarg.m); + testCond(s == NATS_OK); + + test("Destroy the connection while still attached: "); + natsConnection_Destroy(nc); + natsMutex_Lock(arg.m); + while ((s != NATS_TIMEOUT) && (arg.detached == 0)) + s = natsCondition_TimedWait(arg.c, arg.m, 5000); + natsMutex_Unlock(arg.m); + testCondNoReturn(s == NATS_OK); + + natsMutex_Lock(arg.m); + arg.done = true; + natsCondition_Broadcast(arg.c); + natsMutex_Unlock(arg.m); + natsThread_Join(arg.t); + natsThread_Destroy(arg.t); + arg.t = NULL; + + // The event loop can still deliver an event for a socket whose removal it has + // not processed yet. The library's own reference, taken at the first attach + // and released by the adapter below, is what keeps that from being a + // use-after-free on `nc`. + test("A late event after the destroy is a no-op: "); + natsConnection_ProcessReadEvent(nc); + natsConnection_ProcessWriteEvent(nc); + testCondNoReturn(natsConnection_Status(nc) == NATS_CONN_STATUS_CLOSED); + + test("The adapter releases on detach: "); + natsConnection_ProcessDetachedEvent(nc); + testCondNoReturn(true); + + natsMutex_Lock(sarg.m); + sarg.done = true; + natsCondition_Broadcast(sarg.c); + natsMutex_Unlock(sarg.m); + natsThread_Join(t); + natsThread_Destroy(t); + + natsOptions_Destroy(opts); + + _destroyDefaultThreadArgs(&sarg); + _destroyDefaultThreadArgs(&arg); +} + void test_EventLoopRetryOnFailedConnect(void) { natsStatus s; @@ -20777,6 +20973,8 @@ void test_EventLoopRetryOnFailedConnect(void) _evLoopDetach)); testCond(s == NATS_OK); + arg.nc = nc; + test("Start event loop: "); natsMutex_Lock(arg.m); arg.sock = NATS_SOCK_INVALID; @@ -20864,6 +21062,8 @@ void test_EventLoopTLS(void) _evLoopDetach)); testCond(s == NATS_OK); + arg.nc = nc; + test("Start server: "); pid = _startServer("nats://127.0.0.1:4443", "-config tls.conf", true); CHECK_SERVER_STARTED(pid);