numinnex commented on code in PR #3944:
URL: https://github.com/apache/iggy/pull/3944#discussion_r3852359216
##########
foreign/node/src/client/client.socket.ts:
##########
@@ -478,6 +478,13 @@ export class CommandResponseStream extends EventEmitter {
{ last: false }
);
const metadata = GET_CLUSTER_METADATA.deserialize(response);
+ // Every read feeds the redial candidates, leaderless ones included: a
+ // roster with no leader still names where the nodes are.
+ this.connection.rememberRoster(
Review Comment:
Both fixed.
The roster now comes from every `GetClusterMetadata` reply: `_processVsr`
feeds `_rememberRoster(parsed)` whenever that command succeeds, whoever asked
for it, so it no longer goes stale between logins. `_readLeaderEndpoint` lost
its own feed since it is covered.
For the leader move: a `TRANSIENT_NOT_ACCEPTED` that outlives a 2s window is
thrown out of the queue as an internal `LeaderMovedError`, and `sendCommand`
re-reads the roster, follows the leader if it moved, and re-issues the command,
bounded by `MAX_LEADER_REDIRECTS`. That is why the recheck lives in
`sendCommand` and not in `_processVsr`: the roster read is itself a queued
command and `_processQueue` is single-flighted, so rechecking from inside a job
would deadlock on its own queue. 57 stays where it was, on the same connection,
as you said — its outcome is unknown anywhere else. Logins are excluded, since
`_settleOnLeader` already owns that path.
##########
foreign/node/src/client/client.connection.test.ts:
##########
@@ -393,6 +393,31 @@ describe('IggyConnection', () => {
}
);
+ it('rotates a redial through the roster it learned while connected',
Review Comment:
Added three cases: dial order (current endpoint, then the seed, then the
roster), spelling dedup (localhost / `::1` / `::ffff:` collapsing onto the
current endpoint), and the destroy-mid-pass case — which is indeed the one that
catches the leaked socket, and it is red against the code as it stood.
##########
foreign/node/src/client/client.socket.test.ts:
##########
@@ -504,6 +504,111 @@ describe('VSR client socket', () => {
}
});
+ // The node a client authenticated on dies; its next command has to complete
+ // on a survivor the roster named, under a session established there.
+ // Mirrors `core/integration/tests/cluster/failover_client_continuity.rs`.
+ it('resumes on a survivor after the node it authenticated on dies',
+ async () => {
+ const primarySockets = new Set<Socket>();
+ let primaryDead = false;
+
+ const survivor = await startVsrServer((frame, socket) => {
+ const operation = frame.readUInt8(REQUEST_OFFSET.operation);
+ if (operation === Operation.Register) {
+ socket.write(replyFrame(Operation.Register, registerReplyBody()));
+ return;
+ }
+ const code = frame.readUInt32LE(REQUEST_OFFSET.reserved);
+ if (code === COMMAND_CODE.GetClusterMetadata) {
+ // The survivor leads once the primary is gone.
+ socket.write(replyFrame(
+ Operation.NonReplicated,
+ twoNodeMetadataBody(primary.port, survivor.port)
+ ));
+ return;
+ }
+ socket.write(replyFrame(operation));
+ });
+
+ const primary = await startVsrServer((frame, socket) => {
+ primarySockets.add(socket);
+ if (primaryDead) {
+ socket.destroy();
+ return;
+ }
+ const operation = frame.readUInt8(REQUEST_OFFSET.operation);
+ if (operation === Operation.Register) {
+ socket.write(replyFrame(Operation.Register, registerReplyBody()));
+ return;
+ }
+ const code = frame.readUInt32LE(REQUEST_OFFSET.reserved);
+ if (code === COMMAND_CODE.GetClusterMetadata) {
+ // The primary leads, so the login settles here and the roster is
+ // only remembered, not acted on, until the node dies.
+ socket.write(replyFrame(
+ Operation.NonReplicated,
+ twoNodeMetadataBody(survivor.port, primary.port)
+ ));
+ return;
+ }
+ socket.write(replyFrame(operation));
+ });
+
+ const config: ClientConfig = {
+ ...vsrConfig(primary.port),
+ reconnect: { enabled: true, interval: 1, maxRetries: 3 }
+ };
+ const client = new CommandResponseStream(config);
+ try {
+ await client.authenticate(config.credentials);
+ await client.sendCommand(60_021, Buffer.alloc(0));
+ assert.ok(
+ primary.frames.some(
+ (frame) => frame.readUInt32LE(REQUEST_OFFSET.reserved) === 60_021
+ ),
+ 'the live primary answered the first command'
+ );
+
+ primaryDead = true;
+ for (const socket of primarySockets)
+ socket.destroy();
+ await primary.close();
+
+ // The attempt in flight when the socket died is allowed to fail; what
+ // is not allowed is never completing one, which is what a client that
+ // only knows the dead endpoint does.
+ let resumed = false;
+ let lastError: unknown;
+ for (let attempt = 0; attempt < 20 && !resumed; attempt += 1) {
Review Comment:
Capped at 2 and dropped the sleep between them, so the comment's 'at most
one failed submission' is what the test actually asserts.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]