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
13 changes: 13 additions & 0 deletions src/conn.c
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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));
Expand Down
1 change: 1 addition & 0 deletions src/natsp.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;

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(EventLoopParserResetOnDisconnect)
_test(EventLoopRetryOnFailedConnect)
_test(EventLoopTLS)
_test(ExtendedReconnectFunctionality)
Expand Down
181 changes: 181 additions & 0 deletions test/test.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down