From ff08afd3da32a34eb0d4864ff270540c355779be Mon Sep 17 00:00:00 2001 From: William Chong <33353798+w1am@users.noreply.github.com> Date: Tue, 4 Aug 2026 18:19:56 +0400 Subject: [PATCH] [DEV-1850] Fix connector accounting when connectors stop on their own (#5694) * [DEV-1850] Fix connector accounting when connectors stop on their own * [DEV-1850] Upgrade Surge and Connectors packages to 1.1.1-alpha.1.72 (cherry picked from commit 1ed8184ebd36fd2698dab41eb34e6003057557b4) --- .../SystemConnectorsFactoryTests.cs | 62 ++++++++ .../Control/ConnectorsActivatorTests.cs | 146 ++++++++++-------- .../Connectors/SystemConnectorsFactory.cs | 63 +++++--- .../Planes/Control/ConnectorsActivator.cs | 50 ++++-- 4 files changed, 220 insertions(+), 101 deletions(-) create mode 100644 src/Connectors/KurrentDB.Connectors.Tests/Infrastructure/SystemConnectorsFactoryTests.cs diff --git a/src/Connectors/KurrentDB.Connectors.Tests/Infrastructure/SystemConnectorsFactoryTests.cs b/src/Connectors/KurrentDB.Connectors.Tests/Infrastructure/SystemConnectorsFactoryTests.cs new file mode 100644 index 00000000000..9952841a189 --- /dev/null +++ b/src/Connectors/KurrentDB.Connectors.Tests/Infrastructure/SystemConnectorsFactoryTests.cs @@ -0,0 +1,62 @@ +// Copyright (c) Kurrent, Inc and/or licensed to Kurrent, Inc under one or more agreements. +// Kurrent, Inc licenses this file to you under the Kurrent License v1 (see LICENSE.md). + +using System.Diagnostics.Metrics; +using Kurrent.Surge.Connectors; +using KurrentDB.Connectors.Infrastructure.Connect.Components.Connectors; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; + +namespace KurrentDB.Connectors.Tests.Infrastructure; + +public class SystemConnectorsFactoryTests(ITestOutputHelper output, ConnectorsAssemblyFixture fixture) : ConnectorsIntegrationTests(output, fixture) { + [Fact] + public Task counts_only_disposed_connector_as_closed() => Fixture.TestWithTimeout(async cts => { + // Arrange + var factory = Fixture.NodeServices.GetRequiredService(); + + var firstId = ConnectorId.From(Fixture.NewConnectorId()); + var secondId = ConnectorId.From(Fixture.NewConnectorId()); + + var closed = new List(); + + using var listener = new MeterListener { + InstrumentPublished = (instrument, meterListener) => { + if (instrument.Meter.Name == "Kurrent.Connectors" && instrument.Name == "kurrent_connector_active_total") + meterListener.EnableMeasurementEvents(instrument); + } + }; + + listener.SetMeasurementEventCallback((_, measurement, tags, _) => { + if (measurement >= 0) + return; + + var connectorId = tags.ToArray().FirstOrDefault(tag => tag.Key == "connector_id").Value?.ToString(); + if (connectorId is not null) + closed.Add(connectorId); + }); + + listener.Start(); + + var first = factory.CreateConnector(firstId, SerilogSinkSettings()); + var second = factory.CreateConnector(secondId, SerilogSinkSettings()); + + await first.Connect(cts.Token); + await second.Connect(cts.Token); + + // Act + await first.DisposeAsync(); + + // Assert + closed.Should().Equal(firstId.ToString()); + + await second.DisposeAsync(); + + closed.Should().Equal(firstId.ToString(), secondId.ToString()); + }); + + static IConfiguration SerilogSinkSettings() => + new ConfigurationBuilder() + .AddInMemoryCollection([new("InstanceTypeName", "serilog-sink")]) + .Build(); +} diff --git a/src/Connectors/KurrentDB.Connectors.Tests/Planes/Control/ConnectorsActivatorTests.cs b/src/Connectors/KurrentDB.Connectors.Tests/Planes/Control/ConnectorsActivatorTests.cs index 1f0983e955d..f735dc64987 100644 --- a/src/Connectors/KurrentDB.Connectors.Tests/Planes/Control/ConnectorsActivatorTests.cs +++ b/src/Connectors/KurrentDB.Connectors.Tests/Planes/Control/ConnectorsActivatorTests.cs @@ -8,115 +8,120 @@ namespace KurrentDB.Connectors.Tests.Planes.Control; [Trait("Category", "ControlPlane")] public class ConnectorsActivatorTests { - [Fact] - public async Task connector_activates() { - // Arrange - var connectorId = ConnectorId.From(Guid.NewGuid()); - var settings = new Dictionary(); - var revision = 1; + const int Revision = 1; + + static (ConnectorsActivator Sut, ConnectorId ConnectorId) CreateSut(TestConnector connector) => + (new ConnectorsActivator((_, _) => connector), ConnectorId.From(Guid.NewGuid())); + + static ValueTask Activate(ConnectorsActivator sut, ConnectorId connectorId) => + sut.Activate(connectorId, NoSettings, Revision); - var testConnector = new TestConnector(failOnConnect: false); + static readonly Dictionary NoSettings = []; - var sut = new ConnectorsActivator(CreateConnector); + [Fact] + public async Task connector_activates() { + var connector = new TestConnector(); + var (sut, connectorId) = CreateSut(connector); - // Act - var result = await sut.Activate(connectorId, settings, revision); + var result = await Activate(sut, connectorId); - // Assert result.Success.Should().BeTrue(); result.Type.Should().Be(ActivateResultType.Activated); - testConnector.IsDisposed.Should().BeFalse(); - testConnector.ConnectionAttempt.Should().Be(1); - return; - - IConnector CreateConnector(ConnectorId connectorId1, IDictionary dictionary) => testConnector; + connector.DisposeCount.Should().Be(0); + connector.ConnectionAttempt.Should().Be(1); } [Fact] public async Task connector_disposed_when_connect_throws_exception() { - // Arrange - var connectorId = ConnectorId.From(Guid.NewGuid()); - var settings = new Dictionary(); - var revision = 1; var exception = new InvalidOperationException("Connection failed"); + var connector = new TestConnector(failOnConnect: true, exception); + var (sut, connectorId) = CreateSut(connector); - var testConnector = new TestConnector(failOnConnect: true, exception); - - var sut = new ConnectorsActivator(CreateConnector); + var result = await Activate(sut, connectorId); - // Act - var result = await sut.Activate(connectorId, settings, revision); - - // Assert result.Failure.Should().BeTrue(); result.Type.Should().Be(ActivateResultType.Unknown); result.Error.Should().Be(exception); - testConnector.IsDisposed.Should().BeTrue(); - testConnector.ConnectionAttempt.Should().Be(1); - return; - - IConnector CreateConnector(ConnectorId connectorId1, IDictionary dictionary) => testConnector; + connector.DisposeCount.Should().Be(1); + connector.ConnectionAttempt.Should().Be(1); + connector.Stopped.Status.Should().Be(TaskStatus.RanToCompletion); } [Fact] public async Task connector_disposed_when_connect_throws_validation_exception() { - // Arrange - var connectorId = ConnectorId.From(Guid.NewGuid()); - var settings = new Dictionary(); - var revision = 1; var validationException = new FluentValidation.ValidationException("Invalid configuration"); + var connector = new TestConnector(failOnConnect: true, validationException); + var (sut, connectorId) = CreateSut(connector); - var testConnector = new TestConnector(failOnConnect: true, validationException); - - var sut = new ConnectorsActivator(CreateConnector); - - // Act - var result = await sut.Activate(connectorId, settings, revision); + var result = await Activate(sut, connectorId); - // Assert result.Failure.Should().BeTrue(); result.Type.Should().Be(ActivateResultType.InvalidConfiguration); result.Error.Should().Be(validationException); - testConnector.IsDisposed.Should().BeTrue(); - testConnector.ConnectionAttempt.Should().Be(1); - return; + connector.DisposeCount.Should().Be(1); + connector.ConnectionAttempt.Should().Be(1); + } + + [Fact] + public async Task deactivates_once() { + var connector = new TestConnector(); + var (sut, connectorId) = CreateSut(connector); - IConnector CreateConnector(ConnectorId connectorId1, IDictionary dictionary) => testConnector; + await Activate(sut, connectorId); + + var result = await sut.Deactivate(connectorId); + + result.Type.Should().Be(DeactivateResultType.Deactivated); + connector.DisposeCount.Should().Be(1); + + var repeated = await sut.Deactivate(connectorId); + + repeated.Type.Should().Be(DeactivateResultType.ConnectorNotFound); + connector.DisposeCount.Should().Be(1); } [Fact] - public async Task connector_stopped_task_completes_on_connect_failure() { - // Arrange - var connectorId = ConnectorId.From(Guid.NewGuid()); - var settings = new Dictionary(); - var revision = 1; - var exception = new InvalidOperationException("Connection failed"); + public async Task deactivates_self_stopped_connector() { + var connector = new TestConnector(); + var (sut, connectorId) = CreateSut(connector); - var testConnector = new TestConnector(failOnConnect: true, exception); + await Activate(sut, connectorId); - var sut = new ConnectorsActivator(CreateConnector); + // a sink failing against an unreachable broker + connector.SimulateSelfTermination(new InvalidOperationException("simulated connector crash")); - // Act - var result = await sut.Activate(connectorId, settings, revision); + var result = await sut.Deactivate(connectorId); - // Assert - result.Failure.Should().BeTrue(); - testConnector.Stopped.IsCompleted.Should().BeTrue(); - testConnector.Stopped.Status.Should().Be(TaskStatus.RanToCompletion); - return; + result.Type.Should().Be(DeactivateResultType.Deactivated); + connector.DisposeCount.Should().Be(1); + } - IConnector CreateConnector(ConnectorId connectorId1, IDictionary dictionary) => testConnector; + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task waits_for_deactivation(bool faulted) { + var connector = new TestConnector(); + var (sut, connectorId) = CreateSut(connector); + + await Activate(sut, connectorId); + + var waiting = sut.WaitForDeactivation(connectorId); + connector.SimulateSelfTermination(faulted ? new InvalidOperationException("simulated connector crash") : null); + var result = await waiting; + + result.Type.Should().Be(DeactivateResultType.Deactivated); + connector.DisposeCount.Should().Be(1); } } internal class TestConnector(bool failOnConnect = false, Exception? exception = null) : IConnector { - readonly TaskCompletionSource _stoppedTcs = new(); + readonly TaskCompletionSource _stoppedTcs = new(TaskCreationOptions.RunContinuationsAsynchronously); public ConnectorId ConnectorId => ConnectorId.From(Guid.NewGuid()); public ConnectorState State { get; private set; } = ConnectorState.Unspecified; public Task Stopped => _stoppedTcs.Task; - public bool IsDisposed { get; private set; } + public int DisposeCount { get; private set; } public int ConnectionAttempt { get; private set; } public Task Connect(CancellationToken stoppingToken) { @@ -131,8 +136,17 @@ public Task Connect(CancellationToken stoppingToken) { return Task.CompletedTask; } + public void SimulateSelfTermination(Exception? error = null) { + State = ConnectorState.Stopped; + + if (error is null) + _stoppedTcs.TrySetResult(); + else + _stoppedTcs.TrySetException(error); + } + public ValueTask DisposeAsync() { - IsDisposed = true; + DisposeCount++; State = ConnectorState.Stopped; _stoppedTcs.TrySetResult(); diff --git a/src/Connectors/KurrentDB.Connectors/Infrastructure/Connect/Components/Connectors/SystemConnectorsFactory.cs b/src/Connectors/KurrentDB.Connectors/Infrastructure/Connect/Components/Connectors/SystemConnectorsFactory.cs index 932e4d3543e..56308e9d170 100644 --- a/src/Connectors/KurrentDB.Connectors/Infrastructure/Connect/Components/Connectors/SystemConnectorsFactory.cs +++ b/src/Connectors/KurrentDB.Connectors/Infrastructure/Connect/Components/Connectors/SystemConnectorsFactory.cs @@ -45,8 +45,6 @@ public class SystemConnectorsFactory(SystemConnectorsFactoryOptions options, ISe SystemConnectorsFactoryOptions Options { get; } = options; IServiceProvider Services { get; } = services; - static DisposeCallback? OnDisposeCallback; - public IConnector CreateConnector(ConnectorId connectorId, IConfiguration configuration) { var options = configuration.GetRequiredOptions(); var validator = Services.GetRequiredService(); @@ -86,29 +84,37 @@ SinkConnector CreateSinkConnector() { connector = new SqlReducerSink(connector, reducer); } - ConnectorMetrics.TrackSinkConnectorCreated(connector.GetType(), connectorId); + Type instanceType = connector.GetType(); - OnDisposeCallback = () => ConnectorMetrics.TrackSinkConnectorClosed(connector.GetType(), connectorId); + ConnectorMetrics.TrackSinkConnectorCreated(instanceType, connectorId); var sinkProxy = new SinkProxy(connectorId, connector, config, Services); var processor = ConfigureSinkProcessor(connectorId, Options.Interceptors, sinkOptions, sinkProxy); - return new SinkConnector(processor, sinkProxy); + return new SinkConnector( + processor, + sinkProxy, + () => ConnectorMetrics.TrackSinkConnectorClosed(instanceType, connectorId) + ); } SourceConnector CreateSourceConnector() { var sourceOptions = configuration.GetRequiredOptions(); - ConnectorMetrics.TrackSourceConnectorCreated(connector.GetType(), connectorId); + Type instanceType = connector.GetType(); - OnDisposeCallback = () => ConnectorMetrics.TrackSourceConnectorClosed(connector.GetType(), connectorId); + ConnectorMetrics.TrackSourceConnectorCreated(instanceType, connectorId); var sourceProxy = new SourceProxy(connectorId, connector, configuration, Services); var processor = ConfigureSourceProcessor(connectorId, Options.Interceptors, sourceOptions, sourceProxy); - return new SourceConnector(connectorId, processor); + return new SourceConnector( + connectorId, + processor, + () => ConnectorMetrics.TrackSourceConnectorClosed(instanceType, connectorId) + ); } dynamic CreateConnectorInstance(string connectorTypeName) { @@ -215,7 +221,7 @@ IProcessor ConfigureSourceProcessor(ConnectorId connectorId, LinkedList (ConnectorState)SourceProcessor.State; + public ConnectorState State => (ConnectorState)sourceProcessor.State; - public Task Stopped => SourceProcessor.Stopped; + public Task Stopped => sourceProcessor.Stopped; protected override async Task ExecuteAsync(CancellationToken stoppingToken) => - await SourceProcessor.Activate(stoppingToken).ConfigureAwait(false); + await sourceProcessor.Activate(stoppingToken).ConfigureAwait(false); public async Task Connect(CancellationToken stoppingToken) => await StartAsync(stoppingToken).ConfigureAwait(false); public async ValueTask DisposeAsync() { - await StopAsync(CancellationToken.None).ConfigureAwait(false); - await SourceProcessor.DisposeAsync().ConfigureAwait(false); - OnDisposeCallback?.Invoke(); + try { + await StopAsync(CancellationToken.None).ConfigureAwait(false); + } + finally { + try { + await sourceProcessor.DisposeAsync().ConfigureAwait(false); + } + finally { + disposeCallback.Invoke(); + } + } } } } diff --git a/src/Connectors/KurrentDB.Connectors/Planes/Control/ConnectorsActivator.cs b/src/Connectors/KurrentDB.Connectors/Planes/Control/ConnectorsActivator.cs index cb42aaec829..28bc1b3ef6d 100644 --- a/src/Connectors/KurrentDB.Connectors/Planes/Control/ConnectorsActivator.cs +++ b/src/Connectors/KurrentDB.Connectors/Planes/Control/ConnectorsActivator.cs @@ -30,7 +30,7 @@ public async ValueTask Activate( if (connector.Revision == revision && connector.Instance.State is ConnectorState.Activating or ConnectorState.Running) return ActivateResult.RevisionAlreadyRunning(); - await Teardown(); + await TryTeardown(connectorId, connector); } try { @@ -45,31 +45,21 @@ public async ValueTask Activate( return ActivateResult.Activated(); } catch (ValidationException ex) { - await Teardown(); + await TryTeardown(connectorId, connector); return ActivateResult.InvalidConfiguration(ex); } catch (Exception ex) { - await Teardown(); + await TryTeardown(connectorId, connector); return ActivateResult.UnknownError(ex); } - - async ValueTask Teardown() { - try { - await connector.Instance.DisposeAsync(); - await connector.Instance.Stopped; - } catch { - // ignore - } - } } public async ValueTask Deactivate(ConnectorId connectorId) { - if (!Connectors.TryRemove(connectorId, out var connector)) + if (!Connectors.TryGetValue(connectorId, out var connector)) return DeactivateResult.ConnectorNotFound(); try { - await connector.Instance.DisposeAsync(); - await connector.Instance.Stopped; + await Teardown(connectorId, connector); return DeactivateResult.Deactivated(); } catch (Exception ex) { @@ -78,17 +68,45 @@ public async ValueTask Deactivate(ConnectorId connectorId) { } public async ValueTask WaitForDeactivation(ConnectorId connectorId) { - if (!Connectors.TryRemove(connectorId, out var connector)) + if (!Connectors.TryGetValue(connectorId, out var connector)) return DeactivateResult.ConnectorNotFound(); try { await connector.Instance.Stopped; + } + catch { + // a faulted stop is still a stop + } + + try { + await Teardown(connectorId, connector); return DeactivateResult.Deactivated(); } catch (Exception ex) { return DeactivateResult.UnknownError(ex); } } + + async ValueTask Teardown(ConnectorId connectorId, (IConnector Instance, int Revision) connector) { + if (Connectors.TryRemove(KeyValuePair.Create(connectorId, connector))) + await connector.Instance.DisposeAsync(); + + try { + await connector.Instance.Stopped; + } + catch { + // a faulted stop is still a stop + } + } + + async ValueTask TryTeardown(ConnectorId connectorId, (IConnector Instance, int Revision) connector) { + try { + await Teardown(connectorId, connector); + } + catch { + // ignore + } + } } public enum ActivateResultType {