From 64f016ee48095cca97115f6de5d7ed2f0bbaaa1a Mon Sep 17 00:00:00 2001 From: John Simons Date: Wed, 19 Aug 2026 16:29:40 +1000 Subject: [PATCH 1/2] Add support for error ingestion scale-out MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Introduces the `--error-ingestion-only` command to allow multiple ServiceControl processes to share the load of draining error queues into a single database. This mode disables singleton background tasks—such as retention sweeps, the retry pipeline, and heartbeat monitoring—to ensure they only run on the primary instance. Support is currently limited to SQL Server and PostgreSQL storage. --- ...viceControl.AcceptanceTests.RavenDB.csproj | 3 + .../When_hosting_error_ingestion_only.cs | 267 ++++++++++++++++++ .../Abstractions/BasePersistence.cs | 5 +- .../PersistenceSettings.cs | 6 + .../CustomChecks/CustomChecksComponent.cs | 6 +- .../ExternalIntegrationsComponent.cs | 6 +- .../HostApplicationBuilderExtensions.cs | 40 ++- .../Commands/ErrorIngestionOnlyCommand.cs | 72 +++++ src/ServiceControl/Hosting/Help.txt | 12 + src/ServiceControl/Hosting/HostArguments.cs | 16 ++ .../Infrastructure/Settings/Settings.cs | 3 + .../HeartbeatMonitoringComponent.cs | 11 +- .../HeartbeatMonitoringHostedService.cs | 13 +- .../Persistence/PersistenceFactory.cs | 1 + .../Recoverability/RecoverabilityComponent.cs | 6 +- 15 files changed, 454 insertions(+), 13 deletions(-) create mode 100644 src/ServiceControl.AcceptanceTests/Recoverability/When_hosting_error_ingestion_only.cs create mode 100644 src/ServiceControl/Hosting/Commands/ErrorIngestionOnlyCommand.cs diff --git a/src/ServiceControl.AcceptanceTests.RavenDB/ServiceControl.AcceptanceTests.RavenDB.csproj b/src/ServiceControl.AcceptanceTests.RavenDB/ServiceControl.AcceptanceTests.RavenDB.csproj index bc922ad126..7cc11bc96c 100644 --- a/src/ServiceControl.AcceptanceTests.RavenDB/ServiceControl.AcceptanceTests.RavenDB.csproj +++ b/src/ServiceControl.AcceptanceTests.RavenDB/ServiceControl.AcceptanceTests.RavenDB.csproj @@ -33,6 +33,9 @@ + + + diff --git a/src/ServiceControl.AcceptanceTests/Recoverability/When_hosting_error_ingestion_only.cs b/src/ServiceControl.AcceptanceTests/Recoverability/When_hosting_error_ingestion_only.cs new file mode 100644 index 0000000000..a5dff19984 --- /dev/null +++ b/src/ServiceControl.AcceptanceTests/Recoverability/When_hosting_error_ingestion_only.cs @@ -0,0 +1,267 @@ +namespace ServiceControl.AcceptanceTests.Recoverability +{ + using System; + using System.Collections.Generic; + using System.Diagnostics; + using System.IO; + using System.Linq; + using System.Runtime.Loader; + using System.Threading.Tasks; + using Microsoft.AspNetCore.TestHost; + using Microsoft.EntityFrameworkCore; + using Microsoft.Extensions.DependencyInjection; + using Microsoft.Extensions.Hosting; + using Microsoft.Extensions.Logging; + using NServiceBus; + using NServiceBus.Routing; + using NServiceBus.Transport; + using NUnit.Framework; + using Particular.ServiceControl.Hosting; + using ServiceBus.Management.Infrastructure.Settings; + using ServiceControl.ExternalIntegrations; + using ServiceControl.Hosting.Commands; + using ServiceControl.Infrastructure; + using ServiceControl.MessageFailures; + using ServiceControl.Operations; + using ServiceControl.Persistence.EFCore.DbContexts; + using ServiceControl.Persistence.EFCore.Entities; + using ServiceControl.Persistence.EFCore.Infrastructure; + using ServiceControl.Recoverability; + using ServiceControl.Transports; + + class When_hosting_error_ingestion_only : AcceptanceTest + { + [Test] + public async Task Should_ingest_without_an_endpoint_and_without_the_single_owner_services() + { + var settings = await CreateSettings(); + + var host = ErrorIngestionOnlyCommand.BuildHost(settings); + + try + { + var hostedServices = host.Services.GetServices() + .Select(hostedService => hostedService.GetType().Name) + .ToArray(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(host.Services.GetService(), Is.Null, + "the host must not run an NServiceBus endpoint"); + Assert.That(host.Services.GetService(), Is.Not.Null); + + Assert.That(hostedServices, Is.EquivalentTo(ExpectedHostedServices), + "the set of hosted services in the error ingestion only host changed. Every one of " + + "these runs on every ingestion node, so decide whether that is safe before updating " + + "this list. ReturnToSenderDequeuer would steal messages from a retry batch, " + + "RetentionSweeper would duplicate the sweep, and EventDispatcherHostedService would " + + "publish every integration event once per node."); + } + } + finally + { + await host.DisposeAsync(); + } + } + + static readonly string[] ExpectedHostedServices = + [ + "GenericWebHostService", // health endpoint only, no ServiceControl API + nameof(ErrorIngestion), // the reason this host exists + "HeartbeatMonitoringHostedService", // warms the endpoint monitor, does not check heartbeats + "InternalCustomChecksHostedService", // reports this node's ingestion health to the database + "MetricsReporterHostedService", + "ExternalIntegrationRequestsDataStore" // its drain is inert here, nothing calls Subscribe + ]; + + [Test] + public void Should_refuse_to_start_against_unsupported_storage() + { + var settings = new Settings(TransportIntegration.TypeName, "RavenDB", CreateLoggingSettings(), + forwardErrorMessages: false, errorRetentionPeriod: TimeSpan.FromDays(10)); + + var exception = Assert.ThrowsAsync(() => + new ErrorIngestionOnlyCommand().Execute(new HostArguments([]), settings)); + + Assert.That(exception.Message, Does.Contain("SQL Server or PostgreSQL")); + } + + [Test] + public async Task Should_ingest_a_failed_message_into_the_shared_database() + { + var settings = await CreateSettings(); + + // The schema and the queues are provisioned by a normal instance, never by an ingest only host. + await new SetupCommand().Execute(new HostArguments([]), settings); + + var messageId = Guid.NewGuid().ToString(); + var host = ErrorIngestionOnlyCommand.BuildHost(settings, builder => builder.WebHost.UseTestServer()); + + try + { + await host.StartAsync(); + + await DispatchFailedMessage(settings, messageId); + + var failedMessage = await WaitForFailedMessage(host, messageId); + + using (Assert.EnterMultipleScope()) + { + Assert.That(failedMessage.Status, Is.EqualTo(FailedMessageStatus.Unresolved)); + Assert.That(failedMessage.ExceptionType, Is.EqualTo("System.InvalidOperationException")); + Assert.That(failedMessage.ExceptionMessage, Is.EqualTo("Simulated failure")); + Assert.That(failedMessage.FailingEndpointAddress, Is.EqualTo("IngestOnly.Receiver@IngestOnlyHost")); + + Assert.That(failedMessage.SendingEndpointName, Is.EqualTo("IngestOnly.Sender")); + Assert.That(failedMessage.SendingEndpointHost, Is.EqualTo("IngestOnlyHost")); + Assert.That(failedMessage.SendingEndpointHostId, Is.Not.Null); + Assert.That(failedMessage.ReceivingEndpointName, Is.EqualTo("IngestOnly.Receiver")); + Assert.That(failedMessage.ReceivingEndpointHost, Is.EqualTo("IngestOnlyHost")); + Assert.That(failedMessage.ReceivingEndpointHostId, Is.Not.Null); + } + + await using var scope = host.Services.GetRequiredService().CreateAsyncScope(); + var dbContext = scope.ServiceProvider.GetRequiredService(); + + var groups = await dbContext.FailedMessageGroups.AsNoTracking() + .Where(group => group.FailedMessageUniqueId == failedMessage.UniqueMessageId) + .ToListAsync(); + var knownEndpoints = await dbContext.KnownEndpoints.AsNoTracking().ToListAsync(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(groups.Select(group => group.Type), Does.Contain("Endpoint Name")); + Assert.That(groups.Select(group => group.Type), Does.Contain("Exception Type and Stack Trace")); + Assert.That(knownEndpoints.Select(endpoint => endpoint.Name), Does.Contain("IngestOnly.Receiver")); + } + + await WaitFor(host, async dbContext => await dbContext.EventLogItems.AsNoTracking() + .AnyAsync(item => item.EventType == "MessageFailed" && item.Description == "Simulated failure"), + "an event log entry for the failure"); + } + finally + { + await host.StopAsync(); + await host.DisposeAsync(); + } + } + + static async Task DispatchFailedMessage(Settings settings, string messageId) + { + var dispatchSettings = new TransportSettings + { + EndpointName = "IngestOnly.Dispatcher", + TransportType = settings.TransportType, + ConnectionString = settings.TransportConnectionString, + ErrorQueue = settings.ErrorQueue, + MaxConcurrency = 1, + AssemblyLoadContextResolver = settings.AssemblyLoadContextResolver + }; + + var customization = TransportFactory.Create(dispatchSettings); + var infrastructure = await customization.CreateTransportInfrastructure("IngestOnly.Dispatcher", dispatchSettings); + + try + { + var headers = new Dictionary + { + [Headers.MessageId] = messageId, + [Headers.EnclosedMessageTypes] = "IngestOnly.SomeCommand, IngestOnly, Version=1.0.0.0, Culture=neutral, PublicKeyToken=null", + [Headers.ConversationId] = Guid.NewGuid().ToString(), + [Headers.OriginatingEndpoint] = "IngestOnly.Sender", + [Headers.OriginatingMachine] = "IngestOnlyHost", + [Headers.OriginatingHostId] = Guid.NewGuid().ToString("N"), + [Headers.ProcessingEndpoint] = "IngestOnly.Receiver", + [Headers.HostDisplayName] = "IngestOnlyHost", + [Headers.HostId] = Guid.NewGuid().ToString("N"), + ["NServiceBus.FailedQ"] = "IngestOnly.Receiver@IngestOnlyHost", + ["NServiceBus.TimeOfFailure"] = DateTimeOffsetHelper.ToWireFormattedString(DateTimeOffset.UtcNow), + ["NServiceBus.TimeSent"] = DateTimeOffsetHelper.ToWireFormattedString(DateTimeOffset.UtcNow), + ["NServiceBus.ExceptionInfo.ExceptionType"] = "System.InvalidOperationException", + ["NServiceBus.ExceptionInfo.Message"] = "Simulated failure", + ["NServiceBus.ExceptionInfo.Source"] = "IngestOnly", + ["NServiceBus.ExceptionInfo.StackTrace"] = " at IngestOnly.Receiver.Handle()" + }; + + var outgoing = new OutgoingMessage(messageId, headers, "{}"u8.ToArray()); + + await infrastructure.Dispatcher.Dispatch( + new TransportOperations(new TransportOperation(outgoing, new UnicastAddressTag(settings.ErrorQueue))), + new TransportTransaction()); + } + finally + { + await infrastructure.Shutdown(); + } + } + + static async Task WaitFor(IHost host, Func> condition, string description) + { + var scopeFactory = host.Services.GetRequiredService(); + var timeout = Stopwatch.StartNew(); + + while (timeout.Elapsed < TimeSpan.FromSeconds(60)) + { + await using var scope = scopeFactory.CreateAsyncScope(); + var dbContext = scope.ServiceProvider.GetRequiredService(); + + if (await condition(dbContext)) + { + return; + } + + await Task.Delay(TimeSpan.FromMilliseconds(200)); + } + + Assert.Fail($"Timed out waiting for {description}."); + } + + static async Task WaitForFailedMessage(IHost host, string messageId) + { + var scopeFactory = host.Services.GetRequiredService(); + var timeout = Stopwatch.StartNew(); + + while (timeout.Elapsed < TimeSpan.FromSeconds(60)) + { + await using var scope = scopeFactory.CreateAsyncScope(); + var dbContext = scope.ServiceProvider.GetRequiredService(); + + var failedMessage = await dbContext.FailedMessages.AsNoTracking() + .SingleOrDefaultAsync(message => message.MessageId == messageId); + + if (failedMessage != null) + { + return failedMessage; + } + + await Task.Delay(TimeSpan.FromMilliseconds(200)); + } + + Assert.Fail($"The failed message {messageId} was not ingested within the timeout."); + return null; + } + + async Task CreateSettings() + { + var settings = new Settings(TransportIntegration.TypeName, StorageConfiguration.PersistenceType, + CreateLoggingSettings(), forwardErrorMessages: false, errorRetentionPeriod: TimeSpan.FromDays(10)) + { + InstanceName = $"IngestOnly.{Guid.NewGuid():n}", + TransportConnectionString = TransportIntegration.ConnectionString, + MaximumConcurrencyLevel = 2, + AssemblyLoadContextResolver = static _ => AssemblyLoadContext.Default + }; + + await StorageConfiguration.CustomizeSettings(settings); + + return settings; + } + + static LoggingSettings CreateLoggingSettings() + { + var logPath = Path.Combine(Path.GetTempPath(), Path.GetRandomFileName()); + Directory.CreateDirectory(logPath); + return new LoggingSettings(Settings.SettingsRootNamespace, defaultLevel: LogLevel.Debug, logPath: logPath); + } + } +} diff --git a/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs b/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs index 4a94742649..632bb0cbbb 100644 --- a/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs +++ b/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs @@ -33,7 +33,10 @@ protected static void RegisterDataStores(IServiceCollection services, EFPersiste services.AddSingleton(p => p.GetRequiredService()); services.AddHostedService(p => p.GetRequiredService()); - services.AddHostedService(); + if (settings.RunRetentionSweep) + { + services.AddHostedService(); + } services.AddSingleton(); services.AddSingleton(); diff --git a/src/ServiceControl.Persistence/PersistenceSettings.cs b/src/ServiceControl.Persistence/PersistenceSettings.cs index c544578a9e..a11b359c22 100644 --- a/src/ServiceControl.Persistence/PersistenceSettings.cs +++ b/src/ServiceControl.Persistence/PersistenceSettings.cs @@ -11,6 +11,12 @@ public abstract class PersistenceSettings //HINT: This needs to be here so that ServerControl instance can add an instance specific metadata to tweak the DatabasePath value public string? DatabasePath { get; set; } + /// + /// Whether this host owns the background deletion of data past its retention period. Only one + /// host in a deployment should, so error ingestion only hosts turn it off. + /// + public bool RunRetentionSweep { get; set; } = true; + public bool EnableFullTextSearchOnBodies { get; set; } = true; public TimeSpan? OverrideCustomCheckRepeatTime { get; set; } diff --git a/src/ServiceControl/CustomChecks/CustomChecksComponent.cs b/src/ServiceControl/CustomChecks/CustomChecksComponent.cs index 0bbda85dee..5e4b20e731 100644 --- a/src/ServiceControl/CustomChecks/CustomChecksComponent.cs +++ b/src/ServiceControl/CustomChecks/CustomChecksComponent.cs @@ -28,7 +28,11 @@ public override void Configure(Settings settings, ITransportCustomization transp hostBuilder.Services.AddEventLogMapping(); hostBuilder.Services.AddEventLogMapping(); hostBuilder.Services.AddEventLogMapping(); - hostBuilder.Services.AddPlatformConnectionProvider(); + + if (!settings.ErrorIngestionOnly) + { + hostBuilder.Services.AddPlatformConnectionProvider(); + } hostBuilder.Services.AddSingleton(); } } diff --git a/src/ServiceControl/ExternalIntegrations/ExternalIntegrationsComponent.cs b/src/ServiceControl/ExternalIntegrations/ExternalIntegrationsComponent.cs index 882e54b35b..6d86407288 100644 --- a/src/ServiceControl/ExternalIntegrations/ExternalIntegrationsComponent.cs +++ b/src/ServiceControl/ExternalIntegrations/ExternalIntegrationsComponent.cs @@ -16,8 +16,12 @@ public override void Configure(Settings settings, ITransportCustomization transp if (!settings.DisableExternalIntegrationsPublishing) { - services.AddHostedService(); services.AddDomainEventHandler(); + + if (!settings.ErrorIngestionOnly) + { + services.AddHostedService(); + } } } } diff --git a/src/ServiceControl/HostApplicationBuilderExtensions.cs b/src/ServiceControl/HostApplicationBuilderExtensions.cs index 56d0d6ff3c..732a60d181 100644 --- a/src/ServiceControl/HostApplicationBuilderExtensions.cs +++ b/src/ServiceControl/HostApplicationBuilderExtensions.cs @@ -3,6 +3,7 @@ namespace Particular.ServiceControl using System; using System.Diagnostics; using System.Runtime.InteropServices; + using System.Threading.Tasks; using global::ServiceControl.CustomChecks; using global::ServiceControl.Hosting; using global::ServiceControl.Infrastructure; @@ -22,6 +23,7 @@ namespace Particular.ServiceControl using Microsoft.Extensions.Logging; using NServiceBus; using NServiceBus.Configuration.AdvancedExtensibility; + using NServiceBus.Hosting; using NServiceBus.Transport; using Particular.LicensingComponent; using ServiceBus.Management.Infrastructure; @@ -32,7 +34,10 @@ static class HostApplicationBuilderExtensions { public static void AddServiceControl(this IHostApplicationBuilder hostBuilder, Settings settings, EndpointConfiguration configuration, params ReadOnlySpan components) { - ArgumentNullException.ThrowIfNull(configuration); + if (!settings.ErrorIngestionOnly) + { + ArgumentNullException.ThrowIfNull(configuration); + } RecordStartup(settings, configuration); @@ -85,14 +90,34 @@ public static void AddServiceControl(this IHostApplicationBuilder hostBuilder, S // directly and to make things more complex of course the order of registration still matters ;) services.AddSingleton(provider => new Lazy(provider.GetRequiredService)); - services.AddLicenseCheck(); services.AddPersistence(settings); services.AddMetrics(settings.PrintMetrics); - NServiceBusFactory.Configure(settings, transportCustomization, transportSettings, configuration); - hostBuilder.Services.AddNServiceBusEndpoint(configuration); + if (settings.ErrorIngestionOnly) + { + // Ingestion receives through its own transport infrastructure and forwards through + // that same infrastructure's dispatcher, so the endpoint is not hosted at all. + var machineName = NServiceBus.Support.RuntimeEnvironment.MachineName; + services.AddSingleton(new HostInformation( + DeterministicGuid.MakeId(machineName, settings.InstanceName), + machineName)); + services.AddSingleton(provider => new CriticalError((context, _) => + { + provider.GetRequiredService>().LogCritical(context.Exception, "{CriticalError}", context.Error); + provider.GetRequiredService().StopApplication(); + return Task.CompletedTask; + })); + } + else + { + services.AddLicenseCheck(); + + NServiceBusFactory.Configure(settings, transportCustomization, transportSettings, configuration); + hostBuilder.Services.AddNServiceBusEndpoint(configuration); + + hostBuilder.AddEmailNotifications(); + } - hostBuilder.AddEmailNotifications(); hostBuilder.AddAsyncTimer(); if (!settings.DisableHealthChecks) @@ -125,6 +150,7 @@ static void RecordStartup(Settings settings, EndpointConfiguration endpointConfi Audit Retention Period (optional): {settings.AuditRetentionPeriod} Error Retention Period: {settings.ErrorRetentionPeriod} Ingest Error Messages: {settings.IngestErrorMessages} +Error Ingestion Only: {settings.ErrorIngestionOnly} Forwarding Error Messages: {settings.ForwardErrorMessages} ServiceControl Logging Level: {settings.LoggingSettings.LogLevel} Selected Transport Customization: {settings.TransportType} @@ -133,7 +159,9 @@ Audit Retention Period (optional): {settings.AuditRetentionPeriod} var logger = LoggerUtil.CreateStaticLogger(typeof(HostApplicationBuilderExtensions), settings.LoggingSettings.LogLevel); logger.LogInformation(startupMessage); - endpointConfiguration.GetSettings().AddStartupDiagnosticsSection("Startup", new + + // There is no endpoint to hang diagnostics off in error ingestion only mode. + endpointConfiguration?.GetSettings().AddStartupDiagnosticsSection("Startup", new { Settings = settings, }); diff --git a/src/ServiceControl/Hosting/Commands/ErrorIngestionOnlyCommand.cs b/src/ServiceControl/Hosting/Commands/ErrorIngestionOnlyCommand.cs new file mode 100644 index 0000000000..42599249f8 --- /dev/null +++ b/src/ServiceControl/Hosting/Commands/ErrorIngestionOnlyCommand.cs @@ -0,0 +1,72 @@ +namespace ServiceControl.Hosting.Commands +{ + using System; + using System.Linq; + using System.Threading; + using System.Threading.Tasks; + using Microsoft.AspNetCore.Builder; + using NServiceBus; + using Particular.ServiceControl; + using Particular.ServiceControl.Hosting; + using ServiceBus.Management.Infrastructure.Settings; + using ServiceControl.EventLog; + using ServiceControl.ExternalIntegrations; + using ServiceControl.Monitoring; + using ServiceControl.Persistence; + using ServiceControl.Recoverability; + + /// + /// Runs a host that does nothing but drain the error queue into the shared database, so several + /// processes can ingest against one database. Everything a deployment may only run once, the + /// retry pipeline, the retention sweep, integration event dispatch and heartbeat monitoring, + /// stays with the primary instance. + /// + class ErrorIngestionOnlyCommand : AbstractCommand + { + static readonly string[] SupportedStorageNames = ["SQLServer", "PostgreSQL"]; + + public override async Task Execute(HostArguments args, Settings settings, CancellationToken cancellationToken = default) + { + EnsureStorageCanScaleOut(settings); + + var app = BuildHost(settings); + + await app.RunAsync(settings.RootUrl); + } + + internal static WebApplication BuildHost(Settings settings, Action customize = null) + { + settings.ErrorIngestionOnly = true; + settings.IngestErrorMessages = true; + settings.RunRetryProcessor = false; + + var hostBuilder = WebApplication.CreateBuilder(); + + hostBuilder.AddServiceControl(settings, configuration: null, Components); + + customize?.Invoke(hostBuilder); + + return hostBuilder.Build(); + } + + static void EnsureStorageCanScaleOut(Settings settings) + { + var manifest = PersistenceManifestLibrary.Find(settings.PersistenceType); + + if (manifest == null || !SupportedStorageNames.Contains(manifest.Name, StringComparer.OrdinalIgnoreCase)) + { + throw new Exception( + $"--error-ingestion-only requires SQL Server or PostgreSQL storage, but this instance is configured to use '{settings.PersistenceType}'. Scaling out error ingestion is not supported for this storage type."); + } + } + + static ServiceControlComponent[] Components => + [ + new EventLogComponent(), + new ExternalIntegrationsComponent(), + new RecoverabilityComponent(), + new HeartbeatMonitoringComponent(), + new CustomChecks.CustomChecksComponent() + ]; + } +} diff --git a/src/ServiceControl/Hosting/Help.txt b/src/ServiceControl/Hosting/Help.txt index 30fad00314..4925b8c494 100644 --- a/src/ServiceControl/Hosting/Help.txt +++ b/src/ServiceControl/Hosting/Help.txt @@ -10,6 +10,18 @@ The REST API, message importers and the background document expiry are all disab This mode is only supported when run interactively. +ERROR INGESTION ONLY + + ServiceControl.exe --error-ingestion-only + +Runs a host that only drains the error queue into the configured database, so several processes can +share the ingestion load. Requires SQL Server or PostgreSQL storage, and requires that the database +has already been provisioned by a normal instance. Exactly one normal instance must still be running: +it owns the retry pipeline, the retention sweep, integration event dispatch and heartbeat monitoring. + +Message bodies must be stored somewhere every host can read, so this mode should not be combined with +file system body storage unless the path is a shared mount. + SERVICE INSTALL AND UNINSTALL AND CONFIGURATION OPTIONS As of Service Control 1.7 the command line uninstall and install switches have been removed. diff --git a/src/ServiceControl/Hosting/HostArguments.cs b/src/ServiceControl/Hosting/HostArguments.cs index f2151d8968..b260543662 100644 --- a/src/ServiceControl/Hosting/HostArguments.cs +++ b/src/ServiceControl/Hosting/HostArguments.cs @@ -53,6 +53,15 @@ public HostArguments(string[] args) } }; + var errorIngestionOnlyOptions = new OptionSet + { + { + "error-ingestion-only", + "Run only error ingestion, for scaling out ingestion across several processes", + s => Command = typeof(ErrorIngestionOnlyCommand) + } + }; + try { externalInstallerOptions.Parse(args); @@ -76,6 +85,13 @@ public HostArguments(string[] args) return; } + errorIngestionOnlyOptions.Parse(args); + + if (Command == typeof(ErrorIngestionOnlyCommand)) + { + return; + } + defaultOptions.Parse(args); } catch (Exception e) diff --git a/src/ServiceControl/Infrastructure/Settings/Settings.cs b/src/ServiceControl/Infrastructure/Settings/Settings.cs index ef46272162..3850d65b52 100644 --- a/src/ServiceControl/Infrastructure/Settings/Settings.cs +++ b/src/ServiceControl/Infrastructure/Settings/Settings.cs @@ -201,6 +201,9 @@ public TimeSpan HeartbeatGracePeriod public bool IngestErrorMessages { get; set; } = true; public bool RunRetryProcessor { get; set; } = true; + // Set by the --error-ingestion-only command, never read from configuration. + public bool ErrorIngestionOnly { get; set; } + public TimeSpan? AuditRetentionPeriod { get; set; } public TimeSpan ErrorRetentionPeriod { get; } diff --git a/src/ServiceControl/Monitoring/HeartbeatMonitoringComponent.cs b/src/ServiceControl/Monitoring/HeartbeatMonitoringComponent.cs index 48a5159dfe..542901eed0 100644 --- a/src/ServiceControl/Monitoring/HeartbeatMonitoringComponent.cs +++ b/src/ServiceControl/Monitoring/HeartbeatMonitoringComponent.cs @@ -30,7 +30,11 @@ public override void Setup(Settings settings, IComponentInstallationContext cont public override void Configure(Settings settings, ITransportCustomization transportCustomization, IHostApplicationBuilder hostBuilder) { hostBuilder.Services.AddHostedService(); - hostBuilder.Services.AddHostedService(); + + if (!settings.ErrorIngestionOnly) + { + hostBuilder.Services.AddHostedService(); + } hostBuilder.Services.AddSingleton(); hostBuilder.Services.AddSingleton(); @@ -48,7 +52,10 @@ public override void Configure(Settings settings, ITransportCustomization transp hostBuilder.Services.AddErrorMessageEnricher(); - hostBuilder.Services.AddPlatformConnectionProvider(); + if (!settings.ErrorIngestionOnly) + { + hostBuilder.Services.AddPlatformConnectionProvider(); + } } } } \ No newline at end of file diff --git a/src/ServiceControl/Monitoring/HeartbeatMonitoringHostedService.cs b/src/ServiceControl/Monitoring/HeartbeatMonitoringHostedService.cs index 00e57af2ce..779a2fb46f 100644 --- a/src/ServiceControl/Monitoring/HeartbeatMonitoringHostedService.cs +++ b/src/ServiceControl/Monitoring/HeartbeatMonitoringHostedService.cs @@ -17,15 +17,25 @@ public HeartbeatMonitoringHostedService(IEndpointInstanceMonitoring monitor, IMo this.persistence = persistence; this.scheduler = scheduler; this.logger = logger; + this.settings = settings; gracePeriod = settings.HeartbeatGracePeriod; } public async Task StartAsync(CancellationToken cancellationToken = default) { await persistence.WarmupMonitoringFromPersistence(monitor, cancellationToken); + + // An ingestion only host receives no heartbeats, so it has nothing to check and would + // only report every endpoint as dead. It still warms the monitor, because the error + // enricher asks it whether an endpoint is new before recording it. + if (settings.ErrorIngestionOnly) + { + return; + } + timer = scheduler.Schedule(CheckEndpoints, TimeSpan.Zero, TimeSpan.FromSeconds(5), e => logger.LogError(e, "Exception occurred when monitoring endpoint instances")); } - public Task StopAsync(CancellationToken cancellationToken = default) => timer.Stop(cancellationToken); + public Task StopAsync(CancellationToken cancellationToken = default) => timer?.Stop(cancellationToken) ?? Task.CompletedTask; async Task CheckEndpoints(CancellationToken cancellationToken) { @@ -42,6 +52,7 @@ async Task CheckEndpoints(CancellationToken cancellatio IAsyncTimer scheduler; TimerJob timer; TimeSpan gracePeriod; + readonly Settings settings; readonly ILogger logger; } diff --git a/src/ServiceControl/Persistence/PersistenceFactory.cs b/src/ServiceControl/Persistence/PersistenceFactory.cs index 745faa5ef3..bc1c655cd8 100644 --- a/src/ServiceControl/Persistence/PersistenceFactory.cs +++ b/src/ServiceControl/Persistence/PersistenceFactory.cs @@ -13,6 +13,7 @@ public static IPersistence Create(Settings settings, bool maintenanceMode = fals //HINT: This is false when executed from acceptance tests settings.PersisterSpecificSettings ??= persistenceConfiguration.CreateSettings(Settings.SettingsRootNamespace); settings.PersisterSpecificSettings.MaintenanceMode = maintenanceMode; + settings.PersisterSpecificSettings.RunRetentionSweep = !settings.ErrorIngestionOnly; var persistence = persistenceConfiguration.Create(settings.PersisterSpecificSettings); return persistence; diff --git a/src/ServiceControl/Recoverability/RecoverabilityComponent.cs b/src/ServiceControl/Recoverability/RecoverabilityComponent.cs index 3e59f938ec..ab17793427 100644 --- a/src/ServiceControl/Recoverability/RecoverabilityComponent.cs +++ b/src/ServiceControl/Recoverability/RecoverabilityComponent.cs @@ -79,7 +79,11 @@ public override void Configure(Settings settings, ITransportCustomization transp services.AddSingleton(); services.AddSingleton(); services.AddSingleton(); - services.AddHostedService(provider => provider.GetRequiredService()); + + if (!settings.ErrorIngestionOnly) + { + services.AddHostedService(provider => provider.GetRequiredService()); + } //Error importer services.AddSingleton(); From 5853be3e2527413c94001852de942cb1c4470bcc Mon Sep 17 00:00:00 2001 From: John Simons Date: Thu, 20 Aug 2026 12:56:49 +1000 Subject: [PATCH 2/2] Patch approval file --- .../APIApprovals.PlatformSampleSettings.approved.txt | 1 + 1 file changed, 1 insertion(+) diff --git a/src/ServiceControl.UnitTests/ApprovalFiles/APIApprovals.PlatformSampleSettings.approved.txt b/src/ServiceControl.UnitTests/ApprovalFiles/APIApprovals.PlatformSampleSettings.approved.txt index 13a979c0e0..203d0facb5 100644 --- a/src/ServiceControl.UnitTests/ApprovalFiles/APIApprovals.PlatformSampleSettings.approved.txt +++ b/src/ServiceControl.UnitTests/ApprovalFiles/APIApprovals.PlatformSampleSettings.approved.txt @@ -68,6 +68,7 @@ "ForwardErrorMessages": false, "IngestErrorMessages": true, "RunRetryProcessor": true, + "ErrorIngestionOnly": false, "AuditRetentionPeriod": null, "ErrorRetentionPeriod": "10.00:00:00", "EventsRetentionPeriod": "14.00:00:00",