numinnex commented on code in PR #3944: URL: https://github.com/apache/iggy/pull/3944#discussion_r3852357895
########## foreign/csharp/Iggy_SDK_Tests/VsrTests/EndpointFailoverTests.cs: ########## @@ -0,0 +1,428 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +using System.Buffers.Binary; +using System.Net; +using System.Net.Sockets; +using System.Text; +using Apache.Iggy.Configuration; +using Apache.Iggy.Contracts.Tcp; +using Apache.Iggy.Enums; +using Apache.Iggy.IggyClient.Implementations; +using Microsoft.Extensions.Logging.Abstractions; + +namespace Apache.Iggy.Tests.VsrTests; + +/// <summary> +/// The node a client signed in on dies; its next request has to complete on a survivor the roster named, +/// under a session established there. Mirrors +/// <c>core/integration/tests/cluster/failover_client_continuity.rs</c>. +/// </summary> +public sealed class EndpointFailoverTests +{ + private const int HeaderSize = 256; + private const int SizeOffset = 48; + private const int CommandOffset = 60; + private const int RequestIdOffset = 168; + private const int RequestOperationOffset = 176; + private const int RequestReservedOffset = 196; + private const int ReplyRequestIdOffset = 200; + private const int ReplyOperationOffset = 208; + private const int ReplyStatusOffset = 216; + + private const byte CommandReply = 8; + private const byte CommandEviction = 13; + private const int EvictionReasonOffset = 255; + private const byte EvictionStaleClient = 13; + private const byte OperationRegister = 1; + private const byte OperationNonReplicated = 2; + private const int GetClusterMetadataCode = 12; + private const int PingCode = 1; + + [Fact] + public async Task ResumesOnASurvivorAfterTheSignedInNodeDies() + { + using var primary = new MockNode(); + using var survivor = new MockNode(); + + // The primary leads, so the sign-in settles there and the roster is only remembered - not acted on - + // until the node dies. + primary.Serve(request => request.Code == GetClusterMetadataCode + ? Reply(OperationNonReplicated, ClusterMetadata(primary.Port, survivor.Port, primary.Port)) + : Answer(request)); + survivor.Serve(request => request.Code == GetClusterMetadataCode + ? Reply(OperationNonReplicated, ClusterMetadata(primary.Port, survivor.Port, survivor.Port)) + : Answer(request)); + + var configuration = new IggyClientConfigurator + { + BaseAddress = $"127.0.0.1:{primary.Port}", + Protocol = Protocol.Tcp, + ReconnectionSettings = new ReconnectionSettings + { + Enabled = true, + MaxRetries = 4, + InitialDelay = TimeSpan.FromMilliseconds(20) + } + }; + using var client = new TcpMessageStream(configuration, NullLoggerFactory.Instance); + + await client.ConnectAsync(TestContext.Current.CancellationToken); + // No auto login: the credentials come from the caller's own sign-in, which is the shape that could not + // reconnect at all before. + await client.LoginUserAsync("iggy", "iggy", TestContext.Current.CancellationToken); + await client.PingAsync(TestContext.Current.CancellationToken); + Assert.Equal(1, primary.Pings); + + primary.Kill(); + + // The request in flight when the node 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. + var (resumed, lastError) = await ResumedWithin(client, TimeSpan.FromSeconds(10)); + Assert.True(resumed, + $"the client has to resume on the survivor the roster named ({lastError}, survivor saw " + + $"{survivor.Registrations} registrations and {survivor.Pings} pings)"); + Assert.True(survivor.Registrations >= 1, "the remembered credentials signed in again on the survivor"); + Assert.True(survivor.Pings >= 1, "the request landed on the survivor"); + } + + /// <summary> + /// Mirrors the integration contract (HeartbeatTests + /// EvictedClient_WithoutAutoLogin_Should_FailFast_And_NotReconnect): a server-side eviction ends the + /// session authoritatively, so the credentials a manual sign-in remembered must not resurrect it - the + /// evicted request surfaces the loss with no reconnect attempt. + /// </summary> + [Fact] + public async Task ServerEvictionForgetsTheRememberedSignIn() + { + using var node = new MockNode(); + var evict = false; + node.Serve(request => + { + if (request.Operation == OperationRegister) + { + return Reply(OperationRegister, RegisterBody(session: 128)); + } + + return evict + ? EvictionFrame(EvictionStaleClient) + : Reply(OperationNonReplicated, request.Code == GetClusterMetadataCode + ? ClusterMetadata(node.Port, node.Port, node.Port) + : []); + }); + + var configuration = new IggyClientConfigurator + { + BaseAddress = $"127.0.0.1:{node.Port}", + Protocol = Protocol.Tcp, + ReconnectionSettings = new ReconnectionSettings + { + Enabled = true, + MaxRetries = 2, + InitialDelay = TimeSpan.FromMilliseconds(20) + } + }; + using var client = new TcpMessageStream(configuration, NullLoggerFactory.Instance); + + await client.ConnectAsync(TestContext.Current.CancellationToken); + await client.LoginUserAsync("iggy", "iggy", TestContext.Current.CancellationToken); + await client.PingAsync(TestContext.Current.CancellationToken); + var connectionsBeforeEviction = node.Connections; + + evict = true; + await Assert.ThrowsAnyAsync<Exception>(() => client.PingAsync(TestContext.Current.CancellationToken)); + Assert.Equal(connectionsBeforeEviction, node.Connections); + + // The dropped connection leaves the next call transport-shaped, but the eviction forgot the remembered + // sign-in, so it must fail fast instead of reconnecting into a resurrected session. + await Assert.ThrowsAnyAsync<Exception>(() => client.PingAsync(TestContext.Current.CancellationToken)); + Assert.Equal(connectionsBeforeEviction, node.Connections); + } + + private static byte[] EvictionFrame(byte reason) + { + var frame = new byte[HeaderSize]; + BinaryPrimitives.WriteUInt32LittleEndian(frame.AsSpan(SizeOffset, 4), HeaderSize); + frame[CommandOffset] = CommandEviction; + frame[EvictionReasonOffset] = reason; + return frame; + } + + [Fact] + public async Task FailsFastWhenNothingEverSignedIn() + { + using var node = new MockNode(); + node.Serve(request => request.Code == GetClusterMetadataCode + ? Reply(OperationNonReplicated, ClusterMetadata(node.Port, node.Port, node.Port)) + : Answer(request)); + + var configuration = new IggyClientConfigurator + { + BaseAddress = $"127.0.0.1:{node.Port}", + Protocol = Protocol.Tcp, + ReconnectionSettings = new ReconnectionSettings + { + Enabled = true, + MaxRetries = 2, + InitialDelay = TimeSpan.FromMilliseconds(20) + } + }; + using var client = new TcpMessageStream(configuration, NullLoggerFactory.Instance); + + await client.ConnectAsync(TestContext.Current.CancellationToken); + await client.PingAsync(TestContext.Current.CancellationToken); + + node.Kill(); + + var (resumed, _) = await ResumedWithin(client, TimeSpan.FromSeconds(2)); Review Comment: Replaced with the event assertion: the test subscribes through `SubscribeConnectionEvents`, kills the node, expects the request to throw, and asserts no `Connecting` transition was ever announced. No polling window. ########## foreign/node/src/client/client.connection.ts: ########## @@ -312,24 +319,28 @@ export class IggyConnection extends EventEmitter { if (this.connected || this.socket !== expectedSocket) return this.connect(); - const options = this._reconnectTarget(attempt); - attempt += 1; - const socket = this._installSocket( - getTransport({ ...this.config, options }) - ); - this.socket = socket; - expectedSocket = socket; - try { - await this._waitForConnection(socket); - if (this.socket !== socket) - return this.connect(); - this.config.options = options; - return this; - } catch (error) { - lastError = error instanceof Error - ? error - : new Error(String(error)); - debug('reconnect attempt failed', lastError); + // Every endpoint gets its turn inside one attempt, so a full pass over + // the cluster costs one retry rather than one per endpoint: a pass that + // stopped at the first refusal would never reach the survivors of a + // client configured for a single retry. + for (const options of this._redialCandidates()) { Review Comment: Fixed. `this.ending` and the supersession check are re-tested at the top of every iteration of the pass, and a dial that wins after an `ending` flip destroys its socket and throws instead of publishing it. `stops a redial pass that is destroyed part-way through` pins it: the current endpoint is dead and a live endpoint sits behind it in the roster, and the destroy fires on the pass's first failure — so it lands between two candidates. The test asserts the live listener never accepts and no `'connect'` fires after the destroy. Without the per-iteration checks the live endpoint is accepted and it fails. ########## foreign/node/src/client/client.connection.ts: ########## @@ -312,24 +319,28 @@ export class IggyConnection extends EventEmitter { if (this.connected || this.socket !== expectedSocket) return this.connect(); - const options = this._reconnectTarget(attempt); - attempt += 1; - const socket = this._installSocket( - getTransport({ ...this.config, options }) - ); - this.socket = socket; - expectedSocket = socket; - try { - await this._waitForConnection(socket); - if (this.socket !== socket) - return this.connect(); - this.config.options = options; - return this; - } catch (error) { - lastError = error instanceof Error - ? error - : new Error(String(error)); - debug('reconnect attempt failed', lastError); + // Every endpoint gets its turn inside one attempt, so a full pass over + // the cluster costs one retry rather than one per endpoint: a pass that + // stopped at the first refusal would never reach the survivors of a + // client configured for a single retry. + for (const options of this._redialCandidates()) { + const socket = this._installSocket( + getTransport({ ...this.config, options }) + ); + this.socket = socket; + expectedSocket = socket; + try { + await this._waitForConnection(socket); Review Comment: Fixed. `_dialWithin` races the dial against a `FAILOVER_DIAL_TIMEOUT_MS` (2s) timer when there is more than one candidate, and destroys the socket on expiry so the pending dial settles and the handle is released. The timer is unref'd and always cleared. ########## foreign/node/src/client/client.connection.ts: ########## @@ -300,7 +308,6 @@ export class IggyConnection extends EventEmitter { ): Promise<this> { let lastError = initialError; let expectedSocket = this.socket; - let attempt = 0; while (enabled && this.reconnectCount < maxRetries) { this.connecting = true; this.reconnectCount += 1; Review Comment: Fixed: the first pass no longer sleeps when more than one candidate is known. Later passes still back off, and a single-endpoint client keeps the full interval, so the 'does not reconnect after destruction' test is unaffected. ########## foreign/node/src/client/client.connection.ts: ########## @@ -341,16 +352,43 @@ export class IggyConnection extends EventEmitter { } /** - * Alternates reconnect dials between the current endpoint and the - * configured seed. After a leader redirect the current endpoint may die - * with the leader, and the seed is the way back to the rest of the cluster. + * Records the cluster roster as redial candidates. + * + * Replaced wholesale rather than merged: the roster is the cluster's own + * answer about where its nodes are, so a node it dropped stops being + * dialed. The configured seed is kept separately and outlives it. + */ + rememberRoster(endpoints: { host: string, port: number }[]): void { + if (endpoints.length === 0) + return; + this.rosterEndpoints = endpoints; + } + + /** + * Endpoints a redial rotates through, likeliest first: where the client + * currently is, the endpoint it was configured with, then the roster it + * learned while connected. After a leader redirect the current endpoint may + * die with the leader, and the rest of the list is the way back to the + * cluster. Duplicates are dropped, so an endpoint the roster merely spells + * differently does not earn a second attempt. */ - private _reconnectTarget(attempt: number): ClientConfig['options'] { - const current = this.config.options; - if (this.seedOptions.host === current.host && - this.seedOptions.port === current.port) - return current; - return attempt % 2 === 0 ? current : this.seedOptions; + _redialCandidates(): ClientConfig['options'][] { + const candidates = [this.config.options]; + const known = [ + this.seedOptions, + ...this.rosterEndpoints.map( + ({ host, port }) => ({ ...this.config.options, host, port }) + ) + ]; + for (const candidate of known) { + const duplicate = candidates.some( + (known) => known.port === candidate.port && Review Comment: Both done: the inner binding is `existing`, and `Endpoint` is now an exported type (`export type Endpoint = { host: string, port: number }`) used by `rememberRoster` and `rosterEndpoints`. ########## foreign/node/src/client/client.connection.ts: ########## @@ -341,16 +352,43 @@ export class IggyConnection extends EventEmitter { } /** - * Alternates reconnect dials between the current endpoint and the - * configured seed. After a leader redirect the current endpoint may die - * with the leader, and the seed is the way back to the rest of the cluster. + * Records the cluster roster as redial candidates. + * + * Replaced wholesale rather than merged: the roster is the cluster's own + * answer about where its nodes are, so a node it dropped stops being + * dialed. The configured seed is kept separately and outlives it. + */ + rememberRoster(endpoints: { host: string, port: number }[]): void { + if (endpoints.length === 0) + return; + this.rosterEndpoints = endpoints; + } + + /** + * Endpoints a redial rotates through, likeliest first: where the client + * currently is, the endpoint it was configured with, then the roster it + * learned while connected. After a leader redirect the current endpoint may + * die with the leader, and the rest of the list is the way back to the + * cluster. Duplicates are dropped, so an endpoint the roster merely spells Review Comment: Softened, and it now states the consequence rather than implying dedup is complete: the loopback aliases and IPv4-mapped addresses collapse, names are not resolved, so a DNS seed and the roster's IP for one node stay two candidates — one wasted dial per pass, not a correctness problem. -- 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]
