From 19edbc411ed49bcfeab27b8517d745e5b3ac92a9 Mon Sep 17 00:00:00 2001 From: Groene AI <270696204+groeneai@users.noreply.github.com> Date: Tue, 15 Sep 2026 06:34:31 +0000 Subject: [PATCH] Destroy the protocol parser on disconnect when an event loop is attached Backport of nats-io/nats.c#1008 by @ arnaudhe, which fixes nats-io/nats.c#1007. Both are still open upstream, where the only discussion is how to test the change, so this carries the patch onto the branch ClickHouse pins and adds the regression test upstream does not have. The parser holds half-decoded bytes that belong to exactly one socket, so it must be discarded before any byte from the next socket is parsed. The library honours that on the blocking-IO path and violates it on the external event-loop path. `natsParser_Destroy` has two call sites in `src/conn.c`: `_freeConn`, and the exit of `_readLoop`, which also NULLs `nc->ps`. `_readLoop` is the blocking-IO reader and every reconnect creates a new one, so there the parser is per socket. When `nc->opts->evLoop` is set, `_processConnInit` attaches the adapter instead, reads arrive through `natsConnection_ProcessReadEvent`, and that function creates the parser lazily and never destroys it. The disconnect notification point for that path is `_processOpError`, which stops polling and leaves `nc->ps` alone. A drop while the parser is inside `MSG_PAYLOAD` therefore survives into the reconnect, and the next bytes read on the new socket are appended to the stale pending payload. Measured on the unpatched tree with the new test: a header announcing 12 payload bytes followed by 4 of them, then a close, then `MSG foo 1 5\r\nhello\r\n` on the reconnected socket, delivers `natsConn_processMsg(buf="AAAAMSG foo ", bufLen=12)` to the subscriber. The stale parser takes the first 8 bytes of the new protocol line as the payload it was still owed, `MSG_END` then skips to the next `\n`, and the leftover `hello` is not a valid op, so the connection is torn down with `NATS_PROTOCOL_ERROR` right after handing out the wrong message. The handshake itself does not pass through the parser and cannot clear it: `_processExpectedInfo` reads the greeting through `_readProto`, which loops `natsSock_Read` one byte at a time, and `_sendConnect` reads the PONG the same way. The flag is deferred rather than a destroy at the disconnect because both direct alternatives are unavailable. `natsConnection_ProcessCloseEvent` takes `(natsSock *socket)` and has no connection pointer, and the adapters call it as `natsConnection_ProcessCloseEvent(&(nle->socket))`, so widening it is an ABI break for every adapter in tree and out. Destroying `nc->ps` inside `_processOpError` would be a use-after-free: `natsConnection_ProcessReadEvent` unlocks before calling `natsParser_Parse`, and `_processOpError` is entered from the ping timer and the write path on other threads. Setting the flag under `natsConn_Lock` and consuming it at the top of the next `ProcessReadEvent` keeps the destroy on the event-loop thread, where the previous socket's `natsParser_Parse` has already returned. Two comments are reworded from upstream's, and nothing else in the code differs: the `natsp.h` field's trailing comment and the one in `_processOpError` both say "external event loop" instead of naming libuv, because the guard is `nc->el.attached`, which is any adapter, and the test below drives the suite's generic event loop. The new `EventLoopParserResetOnDisconnect` needs no broker binary and no libuv. It drives the suite's generic external event loop, whose thread calls the patched `natsConnection_ProcessReadEvent`, against a two-accept `_startMockupServer` thread that serves the initial connection and the reconnect on one listener. Reading the resent `SUB` is what orders the post-reconnect send after the handshake, so the test does not sleep. It is a new registered case rather than an extension of `test_EventLoop` or `test_EventLoopRetryOnFailedConnect` because both of those drive a real `nats-server` through `_startServer`, and a real server never announces 12 payload bytes and then sends 4, so the defect is unreachable from either. The three checks do not return on failure, so the event loop thread is always stopped by the teardown before the library is closed. Validated with clang-22, Debug, `BUILD_TESTING=ON`, `NATS_BUILD_STREAMING=OFF`, `NATS_BUILD_LIB_STATIC=ON`, out of tree. With the two `src/` hunks reverted and the test present the case fails on 5 of 5 runs of this tree, always at the content check, with the payload above; with them applied it passes on 7 of 7. Two earlier forms of the test, before its teardown and comments were cleaned up, measured the same direction on their own binaries: 13 base failures, 0 base passes, 0 fix failures. The 114 tests of this suite that need no broker binary go 113 of 114 on the reverted tree, the single failure being this new case, and 114 of 114 with the fix, so nothing else moves. The two source files also build clean in ClickHouse's own configuration, which is where this branch is consumed. Reported on: https://github.com/ClickHouse/ClickHouse/pull/76867#issuecomment-5674537202 CI report: https://s3.amazonaws.com/clickhouse-test-reports/praktika.html?PR=76867&sha=dfc5ef24014997273075680988a2e6a1712b21b3&name_0=PR&name_1=Integration%20tests%20(amd_asan_ubsan,%20db%20disk,%20old%20analyzer,%204%2F8) Related: https://github.com/ClickHouse/ClickHouse/issues/116417#issuecomment-5627881074 Related: https://github.com/nats-io/nats.c/issues/1007 Related: https://github.com/nats-io/nats.c/pull/1008 Co-Authored-By: Claude Opus 5 --- src/conn.c | 13 ++++ src/natsp.h | 1 + test/list_test.txt | 1 + test/test.c | 181 +++++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 196 insertions(+) diff --git a/src/conn.c b/src/conn.c index e03fed0cb..13218817a 100644 --- a/src/conn.c +++ b/src/conn.c @@ -2214,6 +2214,10 @@ _processOpError(natsConnection *nc, natsStatus s, bool initialConnect) // but the actual socket close will be done from the event loop // adapter by calling natsConnection_ProcessCloseEvent(). ls = _evStopPolling(nc); + + // The blocking-IO reader destroys the parser once per socket at _readLoop's exit; + // an attached event loop has no such per-socket boundary, so mark it here. + nc->psReset = true; } // Fail pending flush requests. @@ -4110,6 +4114,15 @@ natsConnection_ProcessReadEvent(natsConnection *nc) return; } + // Same event-loop thread runs both this destroy and natsParser_Parse below, + // so the in-flight Parse on the previous socket has already returned. + if (nc->psReset) + { + natsParser_Destroy(nc->ps); + nc->ps = NULL; + nc->psReset = false; + } + if (nc->ps == NULL) { s = natsParser_Create(&(nc->ps)); diff --git a/src/natsp.h b/src/natsp.h index 1a6cc1410..41ebf1570 100644 --- a/src/natsp.h +++ b/src/natsp.h @@ -708,6 +708,7 @@ struct __natsConnection char errStr[256]; natsParser *ps; + bool psReset; // external event loop: destroy nc->ps on next ProcessReadEvent natsTimer *ptmr; int pout; diff --git a/test/list_test.txt b/test/list_test.txt index 8a2045581..58ab69766 100644 --- a/test/list_test.txt +++ b/test/list_test.txt @@ -51,6 +51,7 @@ _test(DrainSubStops) _test(ErrOnConnectAndDeadlock) _test(ErrOnMaxPayloadLimit) _test(EventLoop) +_test(EventLoopParserResetOnDisconnect) _test(EventLoopRetryOnFailedConnect) _test(EventLoopTLS) _test(ExtendedReconnectFunctionality) diff --git a/test/test.c b/test/test.c index d55d33ec6..f0f6d7eea 100644 --- a/test/test.c +++ b/test/test.c @@ -20443,6 +20443,187 @@ void test_EventLoop(void) _stopServer(pid); } +static void +_parserResetMockupServerThread(void *closure) +{ + natsStatus s = NATS_OK; + natsSock sock = NATS_SOCK_INVALID; + struct threadArg *arg = (struct threadArg*) closure; + natsSockCtx ctx[2]; + char buffer[1024]; + int i; + // A header announcing 12 payload bytes followed by only 4 of them leaves the + // client's parser inside MSG_PAYLOAD when the socket goes away. + const char *truncatedMsg = "MSG foo 1 12\r\nAAAA"; + const char *completeMsg = "MSG foo 1 5\r\nhello\r\n"; + + memset(&ctx, 0, sizeof(ctx)); + ctx[0].fd = NATS_SOCK_INVALID; + ctx[1].fd = NATS_SOCK_INVALID; + + s = _startMockupServer(&sock, "localhost", "4222"); + natsMutex_Lock(arg->m); + arg->status = s; + natsCondition_Signal(arg->c); + natsMutex_Unlock(arg->m); + + // Serve the initial connection, then the reconnect, on the same listener. + for (i = 0; (s == NATS_OK) && (i < 2); i++) + { + if (((ctx[i].fd = accept(sock, NULL, NULL)) == NATS_SOCK_INVALID) + || (natsSock_SetCommonTcpOptions(ctx[i].fd) != NATS_OK)) + { + s = NATS_SYS_ERROR; + break; + } + + s = natsSock_WriteFully(&(ctx[i]), arg->string, (int) strlen(arg->string)); + // natsSock_ReadLine keeps the bytes after the line it returns, so the + // buffer is cleared once per socket, before its first read. + buffer[0] = '\0'; + // CONNECT, then PING. + IFOK(s, natsSock_ReadLine(&(ctx[i]), buffer, sizeof(buffer))); + IFOK(s, natsSock_ReadLine(&(ctx[i]), buffer, sizeof(buffer))); + IFOK(s, natsSock_WriteFully(&(ctx[i]), _PONG_PROTO_, _PONG_PROTO_LEN_)); + // Reading the SUB orders the send below after the client's handshake. + IFOK(s, natsSock_ReadLine(&(ctx[i]), buffer, sizeof(buffer))); + + if (i == 0) + { + IFOK(s, natsSock_WriteFully(&(ctx[i]), truncatedMsg, (int) strlen(truncatedMsg))); + natsSock_Close(ctx[i].fd); + ctx[i].fd = NATS_SOCK_INVALID; + } + else + { + IFOK(s, natsSock_WriteFully(&(ctx[i]), completeMsg, (int) strlen(completeMsg))); + } + } + + if (s == NATS_OK) + { + 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[1].fd); + natsSock_Close(ctx[0].fd); + natsSock_Close(sock); +} + +void test_EventLoopParserResetOnDisconnect(void) +{ + natsStatus s; + natsConnection *nc = NULL; + natsOptions *opts = NULL; + natsSubscription *sub = NULL; + natsMsg *msg = NULL; + natsMsg *msg2 = 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_SetMaxReconnect(opts, 100)); + IFOK(s, natsOptions_SetReconnectWait(opts, 50)); + IFOK(s, natsOptions_SetEventLoop(opts, (void*) &arg, + _evLoopAttach, + _evLoopRead, + _evLoopWrite, + _evLoopDetach)); + IFOK(s, natsOptions_SetDisconnectedCB(opts, _disconnectedCb, (void*) &arg)); + IFOK(s, natsOptions_SetReconnectedCB(opts, _reconnectedCb, (void*) &arg)); + IFOK(s, natsOptions_SetClosedCB(opts, _closedCb, (void*) &arg)); + testCond(s == NATS_OK); + + 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, _parserResetMockupServerThread, (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: "); + s = natsConnection_Connect(&nc, opts); + testCond(s == NATS_OK); + + test("Create sub: "); + s = natsConnection_SubscribeSync(&sub, nc, "foo"); + testCond(s == NATS_OK); + + test("Wait for reconnect after the truncated message: "); + natsMutex_Lock(arg.m); + while ((s != NATS_TIMEOUT) && !arg.reconnected) + s = natsCondition_TimedWait(arg.c, arg.m, 5000); + natsMutex_Unlock(arg.m); + testCond(s == NATS_OK); + + // The three checks below do not return on failure: the event loop thread has + // to be stopped by the teardown before the library is closed. + test("A message is delivered after the reconnect: "); + s = natsSubscription_NextMsg(&msg, sub, 5000); + testCondNoReturn((s == NATS_OK) && (msg != NULL)); + + test("It is not corrupted by the payload pending when the socket went away: "); + testCondNoReturn((msg != NULL) && (natsMsg_GetDataLength(msg) == 5) + && (strncmp(natsMsg_GetData(msg), "hello", 5) == 0)); + + test("No stale message follows and the connection survived: "); + s = natsSubscription_NextMsg(&msg2, sub, 250); + testCondNoReturn((s == NATS_TIMEOUT) + && (natsConnection_Status(nc) == NATS_CONN_STATUS_CONNECTED)); + + natsMutex_Lock(sarg.m); + sarg.done = true; + natsCondition_Broadcast(sarg.c); + natsMutex_Unlock(sarg.m); + natsThread_Join(t); + natsThread_Destroy(t); + + natsConnection_Close(nc); + + natsMutex_Lock(arg.m); + arg.done = true; + natsCondition_Broadcast(arg.c); + natsMutex_Unlock(arg.m); + + natsThread_Join(arg.t); + natsThread_Destroy(arg.t); + + natsMsg_Destroy(msg); + natsMsg_Destroy(msg2); + natsSubscription_Destroy(sub); + natsConnection_Destroy(nc); + natsOptions_Destroy(opts); + + _destroyDefaultThreadArgs(&sarg); + _destroyDefaultThreadArgs(&arg); +} + void test_EventLoopRetryOnFailedConnect(void) { natsStatus s;