diff --git a/src/ServiceControl.AcceptanceTests.RavenDB/ServiceControl.AcceptanceTests.RavenDB.csproj b/src/ServiceControl.AcceptanceTests.RavenDB/ServiceControl.AcceptanceTests.RavenDB.csproj index 7cc11bc96c..5a749a1181 100644 --- a/src/ServiceControl.AcceptanceTests.RavenDB/ServiceControl.AcceptanceTests.RavenDB.csproj +++ b/src/ServiceControl.AcceptanceTests.RavenDB/ServiceControl.AcceptanceTests.RavenDB.csproj @@ -36,6 +36,9 @@ + + + diff --git a/src/ServiceControl.AcceptanceTests.RavenDB/StartupModeTests.cs b/src/ServiceControl.AcceptanceTests.RavenDB/StartupModeTests.cs index 0260330da6..1cd5328ce9 100644 --- a/src/ServiceControl.AcceptanceTests.RavenDB/StartupModeTests.cs +++ b/src/ServiceControl.AcceptanceTests.RavenDB/StartupModeTests.cs @@ -51,6 +51,10 @@ public async Task CanRunMaintenanceMode() using var host = hostBuilder.Build(); await host.StartAsync(); + + // Maintenance mode never calls AddServiceControl, so the persister has to supply the clock itself. + Assert.That(host.Services.GetRequiredService(), Is.Not.Null); + await host.StopAsync(); } diff --git a/src/ServiceControl.AcceptanceTests/ClockRegistrationTests.cs b/src/ServiceControl.AcceptanceTests/ClockRegistrationTests.cs new file mode 100644 index 0000000000..09977fdebe --- /dev/null +++ b/src/ServiceControl.AcceptanceTests/ClockRegistrationTests.cs @@ -0,0 +1,55 @@ +namespace ServiceControl.AcceptanceTests +{ + using System; + using System.Runtime.Loader; + using System.Threading.Tasks; + using Microsoft.Extensions.DependencyInjection; + using Microsoft.Extensions.Hosting; + using NServiceBus; + using NUnit.Framework; + using Particular.ServiceControl; + using ServiceBus.Management.Infrastructure.Settings; + + class ClockRegistrationTests : AcceptanceTest + { + Settings settings; + + [SetUp] + public async Task InitializeSettings() + { + settings = new Settings( + transportType: TransportIntegration.TypeName, + persisterType: StorageConfiguration.PersistenceType, + forwardErrorMessages: false, + errorRetentionPeriod: TimeSpan.FromDays(1)) + { + TransportConnectionString = TransportIntegration.ConnectionString, + IngestErrorMessages = false, + RunRetryProcessor = false, + DisableHealthChecks = true, + AssemblyLoadContextResolver = static _ => AssemblyLoadContext.Default + }; + + await StorageConfiguration.CustomizeSettings(settings); + } + + [Test] + public void Should_keep_the_clock_the_host_registered() + { + TimeProvider hostClock = new StubTimeProvider(); + + var endpointConfiguration = new EndpointConfiguration(settings.InstanceName); + endpointConfiguration.AssemblyScanner().Disable = true; + + var hostBuilder = Host.CreateApplicationBuilder(); + hostBuilder.Services.AddSingleton(hostClock); + hostBuilder.AddServiceControl(settings, endpointConfiguration); + + using var host = hostBuilder.Build(); + + Assert.That(host.Services.GetRequiredService(), Is.SameAs(hostClock)); + } + + sealed class StubTimeProvider : TimeProvider; + } +} diff --git a/src/ServiceControl.AcceptanceTests/MaintenanceModeTests.cs b/src/ServiceControl.AcceptanceTests/MaintenanceModeTests.cs new file mode 100644 index 0000000000..549cd87221 --- /dev/null +++ b/src/ServiceControl.AcceptanceTests/MaintenanceModeTests.cs @@ -0,0 +1,31 @@ +namespace ServiceControl.AcceptanceTests +{ + using System; + using System.Runtime.Loader; + using Microsoft.Extensions.DependencyInjection; + using NUnit.Framework; + using Persistence; + using ServiceBus.Management.Infrastructure.Settings; + + class MaintenanceModeTests : AcceptanceTest + { + // The refusal happens before any database is touched, so this needs no storage set up. + [Test] + public void Should_refuse_maintenance_mode() + { + var settings = new Settings( + transportType: TransportIntegration.TypeName, + persisterType: StorageConfiguration.PersistenceType, + forwardErrorMessages: false, + errorRetentionPeriod: TimeSpan.FromDays(1)) + { + TransportConnectionString = TransportIntegration.ConnectionString, + AssemblyLoadContextResolver = static _ => AssemblyLoadContext.Default + }; + + var exception = Assert.Throws(() => new ServiceCollection().AddPersistence(settings, maintenanceMode: true)); + + Assert.That(exception.Message, Does.Contain("Maintenance mode is not supported").And.Contain(StorageConfiguration.PersistenceType)); + } + } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs b/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs index e88ea72bfc..f96d37b3ca 100644 --- a/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs +++ b/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs @@ -23,7 +23,8 @@ public abstract class BasePersistence { protected static void RegisterDataStores(IServiceCollection services, EFPersisterSettings settings) { - services.AddSingleton(TimeProvider.System); + services.TryAddSingleton(TimeProvider.System); + services.AddSingleton(); services.AddSingleton(); diff --git a/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersistenceConfigurationBase.cs b/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersistenceConfigurationBase.cs index 62529efbf2..1957faad3a 100644 --- a/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersistenceConfigurationBase.cs +++ b/src/ServiceControl.Persistence.EFCore/Abstractions/EFPersistenceConfigurationBase.cs @@ -29,6 +29,8 @@ public abstract class EFPersistenceConfigurationBase : PersistenceConfiguration, const string SubscriptionCacheDurationKey = "SubscriptionCacheDuration"; const string ExternalIntegrationsDispatchingBatchSizeKey = "ExternalIntegrationsDispatchingBatchSize"; + public bool SupportsMaintenanceMode => false; + public PersistenceSettings CreateSettings(SettingsRootNamespace settingsRootNamespace) { var settings = CreateSettings( diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageLifecycleDataStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageLifecycleDataStore.cs index e17695d8a6..bdb0201b1b 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageLifecycleDataStore.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageLifecycleDataStore.cs @@ -7,7 +7,8 @@ namespace ServiceControl.Persistence.EFCore.Implementation; /// /// Every operation here has to update both StatusChangedAt and LastModified. /// -public class FailedMessageLifecycleDataStore(IServiceScopeFactory scopeFactory) : DataStoreBase(scopeFactory), IFailedMessageLifecycleDataStore +public class FailedMessageLifecycleDataStore(IServiceScopeFactory scopeFactory, TimeProvider timeProvider) + : DataStoreBase(scopeFactory), IFailedMessageLifecycleDataStore { public async Task MarkAsArchived(string failedMessageId, CancellationToken cancellationToken = default) { @@ -18,7 +19,7 @@ await ExecuteWithDbContext(async (dbContext, token) => return; } - var now = DateTime.UtcNow; + var now = timeProvider.GetUtcNow().UtcDateTime; await dbContext.FailedMessages .Where(fm => fm.UniqueMessageId == uniqueMessageId) @@ -38,7 +39,7 @@ public async Task MarkAsResolved(string failedMessageId, CancellationToken return false; } - var now = DateTime.UtcNow; + var now = timeProvider.GetUtcNow().UtcDateTime; var affected = await dbContext.FailedMessages .Where(fm => fm.UniqueMessageId == uniqueMessageId && fm.Status != FailedMessageStatus.Resolved) @@ -65,7 +66,7 @@ public async Task UnArchiveMessages(IEnumerable failedMessageI return await ExecuteWithDbContext(async (dbContext, token) => { - var now = DateTime.UtcNow; + var now = timeProvider.GetUtcNow().UtcDateTime; // Query which messages will actually be unarchived (must be Archived status) var unarchivableIds = await dbContext.FailedMessages @@ -93,7 +94,7 @@ public async Task UnArchiveMessagesByRange(DateTime from, DateTime to, { return await ExecuteWithDbContext(async (dbContext, token) => { - var now = DateTime.UtcNow; + var now = timeProvider.GetUtcNow().UtcDateTime; // Query which messages will be unarchived (must be Archived and within the date range) var unarchivableIds = await dbContext.FailedMessages @@ -128,7 +129,7 @@ await ExecuteWithDbContext(async (dbContext, token) => return; } - var now = DateTime.UtcNow; + var now = timeProvider.GetUtcNow().UtcDateTime; await dbContext.FailedMessages .Where(fm => fm.UniqueMessageId == uniqueMessageId && fm.Status == FailedMessageStatus.RetryIssued) diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreArchivingManager.cs b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreArchivingManager.cs index 0602bfe2e0..188825d935 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreArchivingManager.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreArchivingManager.cs @@ -9,14 +9,14 @@ namespace ServiceControl.Persistence.EFCore.Implementation; /// EFCore equivalent of the RavenDB . Wraps the shared /// singleton to manage in-memory archive progress state. /// -class EFCoreArchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager, ArchiveMetrics metrics) +class EFCoreArchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager, ArchiveMetrics metrics, TimeProvider timeProvider) { InMemoryArchive GetOrCreate(ArchiveType archiveType, string requestId) { var id = InMemoryArchive.MakeId(requestId, archiveType); if (!operationsManager.ArchiveOperations.TryGetValue(id, out var summary)) { - summary = new InMemoryArchive(requestId, archiveType, domainEvents, metrics); + summary = new InMemoryArchive(requestId, archiveType, domainEvents, timeProvider, metrics); operationsManager.ArchiveOperations[id] = summary; } @@ -43,7 +43,7 @@ public Task StartArchiving(string requestId, ArchiveType archiveType, Cancellati summary.TotalNumberOfMessages = 0; summary.NumberOfMessagesArchived = 0; - summary.Started = DateTime.UtcNow; + summary.Started = timeProvider.GetUtcNow().UtcDateTime; summary.GroupName = "Undefined"; summary.NumberOfBatches = 0; summary.CurrentBatch = 0; diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreUnarchivingManager.cs b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreUnarchivingManager.cs index 276ae07788..167b9dd432 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreUnarchivingManager.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreUnarchivingManager.cs @@ -9,14 +9,14 @@ namespace ServiceControl.Persistence.EFCore.Implementation; /// EFCore equivalent of the RavenDB . Wraps the shared /// singleton to manage in-memory unarchive progress state. /// -class EFCoreUnarchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager, ArchiveMetrics metrics) +class EFCoreUnarchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager, ArchiveMetrics metrics, TimeProvider timeProvider) { InMemoryUnarchive GetOrCreate(ArchiveType archiveType, string requestId) { var id = InMemoryUnarchive.MakeId(requestId, archiveType); if (!operationsManager.UnarchiveOperations.TryGetValue(id, out var summary)) { - summary = new InMemoryUnarchive(requestId, archiveType, domainEvents, metrics); + summary = new InMemoryUnarchive(requestId, archiveType, domainEvents, timeProvider, metrics); operationsManager.UnarchiveOperations[id] = summary; } @@ -43,7 +43,7 @@ public Task StartUnarchiving(string requestId, ArchiveType archiveType, Cancella summary.TotalNumberOfMessages = 0; summary.NumberOfMessagesUnarchived = 0; - summary.Started = DateTime.UtcNow; + summary.Started = timeProvider.GetUtcNow().UtcDateTime; summary.GroupName = "Undefined"; summary.NumberOfBatches = 0; summary.CurrentBatch = 0; diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs index 579abd6211..90a73be8ca 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs @@ -31,8 +31,8 @@ ILogger logger this.logger = logger; this.domainEvents = domainEvents; - archivingManager = new EFCoreArchivingManager(domainEvents, operationsManager, metrics); - unarchivingManager = new EFCoreUnarchivingManager(domainEvents, operationsManager, metrics); + archivingManager = new EFCoreArchivingManager(domainEvents, operationsManager, metrics, timeProvider); + unarchivingManager = new EFCoreUnarchivingManager(domainEvents, operationsManager, metrics, timeProvider); } public async Task ArchiveAllInGroup(string groupId, AuditUser? initiatedBy = null, string? operationId = null, CancellationToken cancellationToken = default) diff --git a/src/ServiceControl.Persistence.RavenDB/CustomChecks/CheckRavenDBIndexLag.cs b/src/ServiceControl.Persistence.RavenDB/CustomChecks/CheckRavenDBIndexLag.cs index 57c191756b..6cac66e4ef 100644 --- a/src/ServiceControl.Persistence.RavenDB/CustomChecks/CheckRavenDBIndexLag.cs +++ b/src/ServiceControl.Persistence.RavenDB/CustomChecks/CheckRavenDBIndexLag.cs @@ -39,6 +39,7 @@ int CheckAndReportIndexesWithTooMuchIndexLag(IndexInformation[] indexes) { if (indexStats.IsStale && indexStats.LastIndexingTime.HasValue) { + // Machine clock on purpose: LastIndexingTime comes from the server, so an injected clock would give a meaningless lag. var indexLag = DateTime.UtcNow - indexStats.LastIndexingTime.Value; if (indexLag > IndexLagThresholdError) diff --git a/src/ServiceControl.Persistence.RavenDB/ExpirationManager.cs b/src/ServiceControl.Persistence.RavenDB/ExpirationManager.cs index bca8b10025..7a8040de71 100644 --- a/src/ServiceControl.Persistence.RavenDB/ExpirationManager.cs +++ b/src/ServiceControl.Persistence.RavenDB/ExpirationManager.cs @@ -12,12 +12,14 @@ class ExpirationManager public const string DeleteExpirationFieldExpression = "delete msg['@metadata']['@expires']"; readonly TimeSpan errorRetentionPeriod; + readonly TimeProvider timeProvider; readonly TimeSpan eventsRetentionPeriod; - public ExpirationManager(RavenPersisterSettings settings) + public ExpirationManager(RavenPersisterSettings settings, TimeProvider timeProvider) { errorRetentionPeriod = settings.ErrorRetentionPeriod; eventsRetentionPeriod = settings.EventsRetentionPeriod; + this.timeProvider = timeProvider; } public void CancelExpiration(IAsyncDocumentSession session, FailedMessage failedMessage) @@ -27,14 +29,14 @@ public void CancelExpiration(IAsyncDocumentSession session, FailedMessage failed public void EnableExpiration(IAsyncDocumentSession session, FailedMessage failedMessage) { - var expiresAt = DateTime.UtcNow + errorRetentionPeriod; + var expiresAt = timeProvider.GetUtcNow().UtcDateTime + errorRetentionPeriod; session.Advanced.GetMetadataFor(failedMessage)[Constants.Documents.Metadata.Expires] = expiresAt; } public void EnableExpiration(IAsyncDocumentSession session, EventLogItem eventLogItem) { - var expiresAt = DateTime.UtcNow + eventsRetentionPeriod; + var expiresAt = timeProvider.GetUtcNow().UtcDateTime + eventsRetentionPeriod; session.Advanced.GetMetadataFor(eventLogItem)[Constants.Documents.Metadata.Expires] = expiresAt; } @@ -45,7 +47,7 @@ public void EnableExpiration(IAsyncDocumentSession session, EventLogItem eventLo // document down one branch and so cannot have it appended to the end. public string EnableExpirationScript(PatchRequest request) { - var expiredAt = DateTime.UtcNow + errorRetentionPeriod; + var expiredAt = timeProvider.GetUtcNow().UtcDateTime + errorRetentionPeriod; request.Values.Add("Expires", expiredAt); diff --git a/src/ServiceControl.Persistence.RavenDB/RavenPersistence.cs b/src/ServiceControl.Persistence.RavenDB/RavenPersistence.cs index 09863a002a..aaf21b227f 100644 --- a/src/ServiceControl.Persistence.RavenDB/RavenPersistence.cs +++ b/src/ServiceControl.Persistence.RavenDB/RavenPersistence.cs @@ -1,9 +1,11 @@ namespace ServiceControl.Persistence.RavenDB; +using System; using CustomChecks; using Editing; using MessageRedirects; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.DependencyInjection.Extensions; using NServiceBus.Unicast.Subscriptions.MessageDrivenSubscriptions; using Operations.BodyStorage; using Operations.BodyStorage.RavenAttachments; @@ -81,6 +83,8 @@ public void AddInstaller(IServiceCollection services) void ConfigureLifecycle(IServiceCollection services) { + services.TryAddSingleton(TimeProvider.System); + services.AddSingleton(settings); services.AddSingleton(settings); diff --git a/src/ServiceControl.Persistence.RavenDB/RavenPersistenceConfiguration.cs b/src/ServiceControl.Persistence.RavenDB/RavenPersistenceConfiguration.cs index 10a68fec56..5f2cdc27c6 100644 --- a/src/ServiceControl.Persistence.RavenDB/RavenPersistenceConfiguration.cs +++ b/src/ServiceControl.Persistence.RavenDB/RavenPersistenceConfiguration.cs @@ -18,6 +18,8 @@ class RavenPersistenceConfiguration : PersistenceConfiguration, IPersistenceConf const string ExternalIntegrationsDispatchingBatchSizeKey = "ExternalIntegrationsDispatchingBatchSize"; const string MaintenanceModeKey = "MaintenanceMode"; + public bool SupportsMaintenanceMode => true; + public PersistenceSettings CreateSettings(SettingsRootNamespace settingsRootNamespace) { var ravenDbLogLevel = SettingsReader.Read(settingsRootNamespace, RavenBootstrapper.RavenDbLogLevelKey, "Warn"); diff --git a/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/ArchiveDocumentManager.cs b/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/ArchiveDocumentManager.cs index e725e5215d..e955d37087 100644 --- a/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/ArchiveDocumentManager.cs +++ b/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/ArchiveDocumentManager.cs @@ -13,7 +13,7 @@ using Raven.Client.Documents.Operations; using Raven.Client.Documents.Session; - class ArchiveDocumentManager(ExpirationManager expirationManager, ILogger logger) + class ArchiveDocumentManager(ExpirationManager expirationManager, ILogger logger, TimeProvider timeProvider) { public Task LoadArchiveOperation(IAsyncDocumentSession session, string groupId, ArchiveType archiveType, CancellationToken cancellationToken = default) => session.LoadAsync(ArchiveOperation.MakeId(groupId, archiveType), cancellationToken); @@ -26,7 +26,7 @@ public async Task CreateArchiveOperation(IAsyncDocumentSession ArchiveType = archiveType, TotalNumberOfMessages = numberOfMessages, NumberOfMessagesArchived = 0, - Started = DateTime.UtcNow, + Started = timeProvider.GetUtcNow().UtcDateTime, GroupName = groupName, NumberOfBatches = (int)Math.Ceiling(numberOfMessages / (float)batchSize), CurrentBatch = 0, diff --git a/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/ArchivingManager.cs b/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/ArchivingManager.cs index 873df977b2..92ac880cab 100644 --- a/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/ArchivingManager.cs +++ b/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/ArchivingManager.cs @@ -9,10 +9,11 @@ class ArchivingManager { - public ArchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager) + public ArchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager, TimeProvider timeProvider) { this.domainEvents = domainEvents; this.operationsManager = operationsManager; + this.timeProvider = timeProvider; } public bool IsArchiveInProgressFor(string requestId) @@ -34,7 +35,7 @@ InMemoryArchive GetOrCreate(ArchiveType archiveType, string requestId) { if (!operationsManager.ArchiveOperations.TryGetValue(InMemoryArchive.MakeId(requestId, archiveType), out var summary)) { - summary = new InMemoryArchive(requestId, archiveType, domainEvents); + summary = new InMemoryArchive(requestId, archiveType, domainEvents, timeProvider); operationsManager.ArchiveOperations[InMemoryArchive.MakeId(requestId, archiveType)] = summary; } @@ -61,7 +62,7 @@ public Task StartArchiving(string requestId, ArchiveType archiveType, Cancellati summary.TotalNumberOfMessages = 0; summary.NumberOfMessagesArchived = 0; - summary.Started = DateTime.UtcNow; + summary.Started = timeProvider.GetUtcNow().UtcDateTime; summary.GroupName = "Undefined"; summary.NumberOfBatches = 0; summary.CurrentBatch = 0; @@ -106,6 +107,7 @@ void RemoveArchiveOperation(string requestId, ArchiveType archiveType) } IDomainEvents domainEvents; + readonly TimeProvider timeProvider; OperationsManager operationsManager; diff --git a/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/MessageArchiver.cs b/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/MessageArchiver.cs index 022fb17d3a..9171a7df6f 100644 --- a/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/MessageArchiver.cs +++ b/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/MessageArchiver.cs @@ -20,6 +20,7 @@ public MessageArchiver( IDomainEvents domainEvents, ExpirationManager expirationManager, IMessageActionAuditLog auditLog, + TimeProvider timeProvider, ILogger logger ) { @@ -30,11 +31,11 @@ ILogger logger this.logger = logger; this.operationsManager = operationsManager; - archiveDocumentManager = new ArchiveDocumentManager(expirationManager, logger); - archivingManager = new ArchivingManager(domainEvents, operationsManager); + archiveDocumentManager = new ArchiveDocumentManager(expirationManager, logger, timeProvider); + archivingManager = new ArchivingManager(domainEvents, operationsManager, timeProvider); - unarchiveDocumentManager = new UnarchiveDocumentManager(); - unarchivingManager = new UnarchivingManager(domainEvents, operationsManager); + unarchiveDocumentManager = new UnarchiveDocumentManager(timeProvider); + unarchivingManager = new UnarchivingManager(domainEvents, operationsManager, timeProvider); } public async Task ArchiveAllInGroup(string groupId, AuditUser? initiatedBy = null, string operationId = null, CancellationToken cancellationToken = default) diff --git a/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/UnarchiveDocumentManager.cs b/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/UnarchiveDocumentManager.cs index 3759aa3a6a..8b1c06c00b 100644 --- a/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/UnarchiveDocumentManager.cs +++ b/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/UnarchiveDocumentManager.cs @@ -12,7 +12,7 @@ using Raven.Client.Documents.Operations; using Raven.Client.Documents.Session; - class UnarchiveDocumentManager + class UnarchiveDocumentManager(TimeProvider timeProvider) { public Task LoadUnarchiveOperation(IAsyncDocumentSession session, string groupId, ArchiveType archiveType, CancellationToken cancellationToken = default) => session.LoadAsync(UnarchiveOperation.MakeId(groupId, archiveType), cancellationToken); @@ -25,7 +25,7 @@ public async Task CreateUnarchiveOperation(IAsyncDocumentSes ArchiveType = archiveType, TotalNumberOfMessages = numberOfMessages, NumberOfMessagesUnarchived = 0, - Started = DateTime.UtcNow, + Started = timeProvider.GetUtcNow().UtcDateTime, GroupName = groupName, NumberOfBatches = (int)Math.Ceiling(numberOfMessages / (float)batchSize), CurrentBatch = 0, diff --git a/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/UnarchivingManager.cs b/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/UnarchivingManager.cs index 105a4a6ed1..1e7572807c 100644 --- a/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/UnarchivingManager.cs +++ b/src/ServiceControl.Persistence.RavenDB/Recoverability/Archiving/UnarchivingManager.cs @@ -9,10 +9,11 @@ class UnarchivingManager { - public UnarchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager) + public UnarchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager, TimeProvider timeProvider) { this.domainEvents = domainEvents; this.operationsManager = operationsManager; + this.timeProvider = timeProvider; } public bool IsUnarchiveInProgressFor(string requestId) @@ -34,7 +35,7 @@ InMemoryUnarchive GetOrCreate(ArchiveType archiveType, string requestId) { if (!operationsManager.UnarchiveOperations.TryGetValue(InMemoryUnarchive.MakeId(requestId, archiveType), out var summary)) { - summary = new InMemoryUnarchive(requestId, archiveType, domainEvents); + summary = new InMemoryUnarchive(requestId, archiveType, domainEvents, timeProvider); operationsManager.UnarchiveOperations[InMemoryUnarchive.MakeId(requestId, archiveType)] = summary; } @@ -61,7 +62,7 @@ public Task StartUnarchiving(string requestId, ArchiveType archiveType, Cancella summary.TotalNumberOfMessages = 0; summary.NumberOfMessagesUnarchived = 0; - summary.Started = DateTime.UtcNow; + summary.Started = timeProvider.GetUtcNow().UtcDateTime; summary.GroupName = "Undefined"; summary.NumberOfBatches = 0; summary.CurrentBatch = 0; @@ -106,6 +107,7 @@ void RemoveUnarchiveOperation(string requestId, ArchiveType archiveType) } IDomainEvents domainEvents; + readonly TimeProvider timeProvider; OperationsManager operationsManager; } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.RavenDB/Throughput/LicensingDataStore.cs b/src/ServiceControl.Persistence.RavenDB/Throughput/LicensingDataStore.cs index ba91a7af15..6e14a20845 100644 --- a/src/ServiceControl.Persistence.RavenDB/Throughput/LicensingDataStore.cs +++ b/src/ServiceControl.Persistence.RavenDB/Throughput/LicensingDataStore.cs @@ -17,7 +17,8 @@ namespace ServiceControl.Persistence.RavenDB.Throughput; class LicensingDataStore( IRavenDocumentStoreProvider storeProvider, - ThroughputDatabaseConfiguration databaseConfiguration) : ILicensingDataStore + ThroughputDatabaseConfiguration databaseConfiguration, + TimeProvider timeProvider) : ILicensingDataStore { internal const string ThroughputTimeSeriesName = "INC: throughput data"; const string AuditServiceMetadataDocumentId = "AuditServiceMetadata"; @@ -142,7 +143,7 @@ public async Task>> GetEndpointT var store = await storeProvider.GetDocumentStore(cancellationToken); using IAsyncDocumentSession session = store.OpenAsyncSession(databaseConfiguration.Name); - var from = DateTime.UtcNow.AddMonths(-ThroughputReporting.ReportedMonths); + var from = timeProvider.GetUtcNow().UtcDateTime.AddMonths(-ThroughputReporting.ReportedMonths); var query = session.Query() .Where(document => document.SanitizedName.In(queueNames)) .Include(builder => builder.IncludeTimeSeries(ThroughputTimeSeriesName, from)); @@ -270,10 +271,11 @@ public async Task IsThereThroughputForLastXDaysForSource(int days, Through return result; } - static async Task IsThereThroughputForLastXDaysInternal(IRavenQueryable baseQuery, int days, bool includeToday, CancellationToken cancellationToken) + async Task IsThereThroughputForLastXDaysInternal(IRavenQueryable baseQuery, int days, bool includeToday, CancellationToken cancellationToken) { - DateTime fromDate = DateTime.UtcNow.AddDays(-days).Date; - DateTime toDate = includeToday ? DateTime.UtcNow.Date : DateTime.UtcNow.AddDays(-1).Date; + DateTime today = timeProvider.GetUtcNow().UtcDateTime.Date; + DateTime fromDate = today.AddDays(-days); + DateTime toDate = includeToday ? today : today.AddDays(-1); var documents = await baseQuery .Select(e => RavenQuery.TimeSeries(e, ThroughputTimeSeriesName, fromDate, toDate).ToList()) diff --git a/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs index d05de5de27..48628f6ad0 100644 --- a/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs +++ b/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs @@ -39,7 +39,9 @@ public async Task Setup(IHostApplicationBuilder hostBuilder) PersistenceSettings = new PostgreSqlPersisterSettings { ConnectionString = connectionStringBuilder.ConnectionString, - BodyStorage = new FileSystemBodyStorageSettings { StoragePath = bodyStoragePath } + BodyStorage = new FileSystemBodyStorageSettings { StoragePath = bodyStoragePath }, + ErrorRetentionPeriod = DefaultRetentionPeriod, + EventsRetentionPeriod = DefaultRetentionPeriod }; var persistence = new PostgreSqlPersistenceConfiguration().Create(PersistenceSettings); diff --git a/src/ServiceControl.Persistence.Tests.RavenDB/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.RavenDB/PersistenceTestsContext.cs index 2910d3c45a..5bcae71f20 100644 --- a/src/ServiceControl.Persistence.Tests.RavenDB/PersistenceTestsContext.cs +++ b/src/ServiceControl.Persistence.Tests.RavenDB/PersistenceTestsContext.cs @@ -7,6 +7,7 @@ namespace ServiceControl.Persistence.Tests; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Time.Testing; using NUnit.Framework; using Raven.Client.Documents; using Raven.Client.Documents.Session; @@ -43,6 +44,8 @@ public async Task Setup(IHostApplicationBuilder hostBuilder) var persistence = new RavenPersistenceConfiguration().Create(PersistenceSettings); + hostBuilder.Services.AddSingleton(FakeTime); + persistence.AddPersistence(hostBuilder.Services); persistence.AddInstaller(hostBuilder.Services); } @@ -81,13 +84,13 @@ public async Task InsertFailedMessages(params FailedMessage[] messages) public IRavenSessionProvider SessionProvider { get; private set; } - // Nothing to do: Raven versions come from document and index etags, which move on every write, so - // no test needs to push its clock. There is no hook for the server clock anyway. - public void AdvanceClock(TimeSpan by) - { - } + public FakeTimeProvider FakeTime { get; } = new(DateTimeOffset.UtcNow); + + // Moves only what the persister stamps itself. The Raven server expires documents on its own + // clock, so advancing this changes the @expires value written, never when the server acts on it. + public void AdvanceClock(TimeSpan by) => FakeTime.Advance(by); - public DateTime UtcNow => DateTime.UtcNow; + public DateTime UtcNow => FakeTime.GetUtcNow().UtcDateTime; public Task CompleteDatabaseOperation() { diff --git a/src/ServiceControl.Persistence.Tests.RavenDB/ServiceControl.Persistence.Tests.RavenDB.csproj b/src/ServiceControl.Persistence.Tests.RavenDB/ServiceControl.Persistence.Tests.RavenDB.csproj index 83f29d47a4..af6adc14d5 100644 --- a/src/ServiceControl.Persistence.Tests.RavenDB/ServiceControl.Persistence.Tests.RavenDB.csproj +++ b/src/ServiceControl.Persistence.Tests.RavenDB/ServiceControl.Persistence.Tests.RavenDB.csproj @@ -15,6 +15,7 @@ + diff --git a/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs index 9077dd38dd..cc9afebabe 100644 --- a/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs +++ b/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs @@ -38,7 +38,9 @@ public async Task Setup(IHostApplicationBuilder hostBuilder) PersistenceSettings = new SqlServerPersisterSettings { ConnectionString = connectionStringBuilder.ConnectionString, - BodyStorage = new FileSystemBodyStorageSettings { StoragePath = bodyStoragePath } + BodyStorage = new FileSystemBodyStorageSettings { StoragePath = bodyStoragePath }, + ErrorRetentionPeriod = DefaultRetentionPeriod, + EventsRetentionPeriod = DefaultRetentionPeriod }; var persistence = new SqlServerPersistenceConfiguration().Create(PersistenceSettings); diff --git a/src/ServiceControl.Persistence.Tests/EFCore/ArchiveMetricsTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/ArchiveMetricsTests.cs index 37317fc884..7a386c5da0 100644 --- a/src/ServiceControl.Persistence.Tests/EFCore/ArchiveMetricsTests.cs +++ b/src/ServiceControl.Persistence.Tests/EFCore/ArchiveMetricsTests.cs @@ -62,7 +62,7 @@ public async Task An_archive_operation_records_batch_gaps_messages_and_total_dur { var metrics = new ArchiveMetrics(MeterFactory, fakeTime); using var recorded = new RecordedArchiveMetrics(MeterFactory); - var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), metrics) { TotalNumberOfMessages = 1500 }; + var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), TimeProvider.System, metrics) { TotalNumberOfMessages = 1500 }; await archive.Start(); Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "started"), Is.EqualTo(1)); @@ -96,7 +96,7 @@ public async Task An_unarchive_operation_is_tagged_unarchive() { var metrics = new ArchiveMetrics(MeterFactory, fakeTime); using var recorded = new RecordedArchiveMetrics(MeterFactory); - var unarchive = new InMemoryUnarchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), metrics) { TotalNumberOfMessages = 200 }; + var unarchive = new InMemoryUnarchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), TimeProvider.System, metrics) { TotalNumberOfMessages = 200 }; await unarchive.Start(); fakeTime.Advance(TimeSpan.FromSeconds(5)); @@ -118,7 +118,7 @@ public async Task An_operation_that_never_completes_stays_visible_on_the_gauge() { var metrics = new ArchiveMetrics(MeterFactory, fakeTime); using var recorded = new RecordedArchiveMetrics(MeterFactory); - var archive = new InMemoryArchive("group-stuck", ArchiveType.FailureGroup, new FakeDomainEvents(), metrics) { TotalNumberOfMessages = 2000 }; + var archive = new InMemoryArchive("group-stuck", ArchiveType.FailureGroup, new FakeDomainEvents(), TimeProvider.System, metrics) { TotalNumberOfMessages = 2000 }; await archive.Start(); await archive.BatchArchived(1000); @@ -135,7 +135,7 @@ public async Task A_restarted_operation_counts_once_on_the_gauge() { var metrics = new ArchiveMetrics(MeterFactory, fakeTime); using var recorded = new RecordedArchiveMetrics(MeterFactory); - var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), metrics) { TotalNumberOfMessages = 10 }; + var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), TimeProvider.System, metrics) { TotalNumberOfMessages = 10 }; await archive.Start(); await archive.BatchArchived(10); @@ -153,7 +153,7 @@ public async Task A_restarted_operation_counts_once_on_the_gauge() [Test] public async Task Without_metrics_the_state_machine_runs_unchanged() { - var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents()) { TotalNumberOfMessages = 10 }; + var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), TimeProvider.System) { TotalNumberOfMessages = 10 }; await archive.Start(); await archive.BatchArchived(10); diff --git a/src/ServiceControl.Persistence.Tests/EFCore/FailedMessageLifecycleDataStoreTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/FailedMessageLifecycleDataStoreTests.cs new file mode 100644 index 0000000000..cc0c0398a1 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EFCore/FailedMessageLifecycleDataStoreTests.cs @@ -0,0 +1,110 @@ +namespace ServiceControl.Persistence.Tests; + +using System; +using System.Threading.Tasks; +using NUnit.Framework; +using ServiceControl.MessageFailures; +using ServiceControl.Persistence.EFCore.Entities; + +class FailedMessageLifecycleDataStoreTests : ErrorIngestionTestBase +{ + [Test] + public async Task Marking_as_archived_stamps_the_injected_clock() + { + var messageId = await Seed(FailedMessageStatus.Unresolved); + + AdvanceClock(TimeSpan.FromDays(7)); + + await FailedMessageLifecycleStore.MarkAsArchived(messageId.ToString()); + + await AssertStampedWithCurrentTime(messageId, FailedMessageStatus.Archived); + } + + [Test] + public async Task Marking_as_resolved_stamps_the_injected_clock() + { + var messageId = await Seed(FailedMessageStatus.Unresolved); + + AdvanceClock(TimeSpan.FromDays(7)); + + Assert.That(await FailedMessageLifecycleStore.MarkAsResolved(messageId.ToString()), Is.True); + + await AssertStampedWithCurrentTime(messageId, FailedMessageStatus.Resolved); + } + + [Test] + public async Task Unarchiving_stamps_the_injected_clock() + { + var messageId = await Seed(FailedMessageStatus.Archived); + + AdvanceClock(TimeSpan.FromDays(7)); + + var unarchived = await FailedMessageLifecycleStore.UnArchiveMessages([messageId.ToString()]); + Assert.That(unarchived, Is.EquivalentTo(new[] { messageId.ToString() })); + + await AssertStampedWithCurrentTime(messageId, FailedMessageStatus.Unresolved); + } + + [Test] + public async Task Unarchiving_by_range_stamps_the_injected_clock() + { + var archivedAt = Now; + var messageId = await Seed(FailedMessageStatus.Archived, archivedAt); + + AdvanceClock(TimeSpan.FromDays(7)); + + var unarchived = await FailedMessageLifecycleStore.UnArchiveMessagesByRange(archivedAt.AddDays(-1), archivedAt.AddDays(1)); + Assert.That(unarchived, Is.EquivalentTo(new[] { messageId.ToString() })); + + await AssertStampedWithCurrentTime(messageId, FailedMessageStatus.Unresolved); + } + + [Test] + public async Task Reverting_a_retry_stamps_the_injected_clock() + { + var messageId = await Seed(FailedMessageStatus.RetryIssued); + + AdvanceClock(TimeSpan.FromDays(7)); + + await FailedMessageLifecycleStore.RevertRetry(messageId.ToString()); + + await AssertStampedWithCurrentTime(messageId, FailedMessageStatus.Unresolved); + } + + async Task AssertStampedWithCurrentTime(Guid messageId, FailedMessageStatus expectedStatus) + { + var row = await GetFailedMessage(messageId); + + using (Assert.EnterMultipleScope()) + { + Assert.That(row.Status, Is.EqualTo(expectedStatus)); + Assert.That(row.StatusChangedAt, Is.EqualTo(Now)); + Assert.That(row.LastModified, Is.EqualTo(Now)); + } + } + + async Task Seed(FailedMessageStatus status, DateTime? statusChangedAt = null) + { + var id = Guid.NewGuid(); + var timestamp = statusChangedAt ?? Now; + + await Store(new FailedMessageEntity + { + UniqueMessageId = id, + Status = status, + StatusChangedAt = timestamp, + LastModified = timestamp, + NumberOfProcessingAttempts = 1, + FirstTimeOfFailure = timestamp, + LastTimeOfFailure = timestamp, + LastAttemptedAt = timestamp, + IsSystemMessage = false, + HeadersJson = "{}", + BodyStoredExternally = false, + BodySize = 0, + FailingEndpointAddress = "Shipping" + }); + + return id; + } +} diff --git a/src/ServiceControl.Persistence.Tests/EFCore/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests/EFCore/PersistenceTestsContext.cs index 17de14b2dd..c28e78a997 100644 --- a/src/ServiceControl.Persistence.Tests/EFCore/PersistenceTestsContext.cs +++ b/src/ServiceControl.Persistence.Tests/EFCore/PersistenceTestsContext.cs @@ -18,6 +18,9 @@ namespace ServiceControl.Persistence.Tests; public partial class PersistenceTestsContext { + // Advancing the clock wakes the live retention sweeper, so retention has to outrun every advance a test makes. + public static readonly TimeSpan DefaultRetentionPeriod = TimeSpan.FromDays(365); + public FakeTimeProvider FakeTime { get; } = new(StorableUtcNow()); // PostgreSQL timestamps only keep microseconds, so a clock seeded straight from diff --git a/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs index 761e624e12..5503a34d00 100644 --- a/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs +++ b/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs @@ -190,6 +190,44 @@ public async Task Archived_messages_are_swept_after_the_archiver_updates_the_tim Assert.That(await FindFailedMessage(messageId), Is.Null); } + [Test] + public async Task Archived_messages_are_swept_after_the_lifecycle_store_updates_the_timestamp() + { + var messageId = await SeedFailedMessage(FailedMessageStatus.Unresolved, Now.AddDays(-40)); + + await FailedMessageLifecycleStore.MarkAsArchived(messageId.ToString()); + + var archived = await FindFailedMessage(messageId); + Assert.That(archived, Is.Not.Null); + Assert.That(archived!.Status, Is.EqualTo(FailedMessageStatus.Archived)); + Assert.That(archived.StatusChangedAt, Is.EqualTo(Now), "the lifecycle store should stamp the current fake time"); + + AdvanceClock(TimeSpan.FromDays(31)); + + await RunRetentionSweep(); + + Assert.That(await FindFailedMessage(messageId), Is.Null); + } + + [Test] + public async Task Resolved_messages_are_swept_after_the_lifecycle_store_updates_the_timestamp() + { + var messageId = await SeedFailedMessage(FailedMessageStatus.Unresolved, Now.AddDays(-40)); + + Assert.That(await FailedMessageLifecycleStore.MarkAsResolved(messageId.ToString()), Is.True); + + var resolved = await FindFailedMessage(messageId); + Assert.That(resolved, Is.Not.Null); + Assert.That(resolved!.Status, Is.EqualTo(FailedMessageStatus.Resolved)); + Assert.That(resolved.StatusChangedAt, Is.EqualTo(Now), "the lifecycle store should stamp the current fake time"); + + AdvanceClock(TimeSpan.FromDays(31)); + + await RunRetentionSweep(); + + Assert.That(await FindFailedMessage(messageId), Is.Null); + } + [Test] public async Task Counts_the_rows_it_deletes() { diff --git a/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs b/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs index ae814fcece..3db037dc57 100644 --- a/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs +++ b/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs @@ -1,7 +1,6 @@ namespace ServiceControl.Persistence.Tests; using System; -using System.Threading; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/ArchivingProgressClockTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/ArchivingProgressClockTests.cs new file mode 100644 index 0000000000..39a76d6ef0 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/Recoverability/ArchivingProgressClockTests.cs @@ -0,0 +1,58 @@ +namespace ServiceControl.Persistence.Tests.Recoverability; + +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using NUnit.Framework; +using ServiceControl.Infrastructure.DomainEvents; +using ServiceControl.Recoverability; + +class ArchivingProgressClockTests : PersistenceTestBase +{ + readonly CapturingDomainEvents events = new(); + + public ArchivingProgressClockTests() => + RegisterServices = services => services.AddSingleton(events); + + [Test] + public async Task Starting_an_archive_stamps_the_injected_clock() + { + AdvanceClock(TimeSpan.FromDays(7)); + + await ArchiveMessages.StartArchiving("group-1", ArchiveType.FailureGroup); + + var starting = events.Raised.OfType().Single(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(starting.StartTime, Is.EqualTo(Now)); + Assert.That(ArchiveMessages.GetArchivalOperations().Single().Started, Is.EqualTo(Now)); + } + } + + [Test] + public async Task Starting_an_unarchive_stamps_the_injected_clock() + { + AdvanceClock(TimeSpan.FromDays(7)); + + await ArchiveMessages.StartUnarchiving("group-1", ArchiveType.FailureGroup); + + var starting = events.Raised.OfType().Single(); + + Assert.That(starting.StartTime, Is.EqualTo(Now)); + } + + sealed class CapturingDomainEvents : IDomainEvents + { + public List Raised { get; } = []; + + public Task Raise(T domainEvent, CancellationToken cancellationToken = default) where T : IDomainEvent + { + Raised.Add(domainEvent); + return Task.CompletedTask; + } + } +} diff --git a/src/ServiceControl.Persistence.Tests/RetryStateTests.cs b/src/ServiceControl.Persistence.Tests/RetryStateTests.cs index ea748e65d5..68df8d8d83 100644 --- a/src/ServiceControl.Persistence.Tests/RetryStateTests.cs +++ b/src/ServiceControl.Persistence.Tests/RetryStateTests.cs @@ -30,7 +30,7 @@ class RetryStateTests : PersistenceTestBase public async Task When_a_group_is_processed_it_is_set_to_the_Preparing_state() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, 1); var status = retryManager.GetStatusForRetryOperation("Test-group", RetryType.FailureGroup); @@ -42,7 +42,7 @@ public async Task When_a_group_is_processed_it_is_set_to_the_Preparing_state() public async Task When_a_group_is_prepared_and_SC_is_started_the_group_is_marked_as_failed() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", false, 1); @@ -83,7 +83,7 @@ public async Task When_the_dequeuer_is_created_then_the_error_address_is_cached( public async Task When_a_group_is_prepared_with_three_batches_and_SC_is_restarted_while_the_first_group_is_being_forwarded_then_the_count_still_matches() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, 2001); @@ -110,7 +110,7 @@ public async Task When_a_group_is_prepared_with_three_batches_and_SC_is_restarte await processor.ProcessBatches(); // mark ready // Simulate SC restart - retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); + retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); var documentManager = new CustomRetryDocumentManager(false, RetryBatchStore, retryManager); @@ -142,7 +142,7 @@ public async Task When_a_group_is_prepared_with_three_batches_and_SC_is_restarte public async Task When_a_group_is_forwarded_the_status_is_Completed() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, 1); @@ -197,7 +197,7 @@ public async Task When_the_batch_being_forwarded_is_gone_the_forwarding_pointer_ public async Task When_there_is_one_poison_message_it_is_removed_from_batch_and_the_status_is_Complete() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, "A", "B", "C"); @@ -246,7 +246,7 @@ public async Task When_there_is_one_poison_message_it_is_removed_from_batch_and_ public async Task When_a_group_has_one_batch_out_of_two_forwarded_the_status_is_Forwarding() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, 1001); @@ -269,7 +269,7 @@ public async Task When_a_group_has_one_batch_out_of_two_forwarded_the_status_is_ public async Task When_a_selection_is_staged_each_message_is_audited_as_a_batch() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); var user = new AuditUser("alice-sub", "Alice"); const string operationId = "op-sel"; var ids = new[] { "A", "B" }; @@ -319,7 +319,7 @@ public async Task When_a_selection_is_staged_each_message_is_audited_as_a_batch( public async Task When_a_group_is_staged_each_message_is_audited_with_the_initiating_user() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); var user = new AuditUser("alice-sub", "Alice"); const string operationId = "op-abc"; @@ -363,7 +363,7 @@ RetryProcessor CreateProcessor(IDomainEvents domainEvents, TestSender sender) => MessageRedirectsDataStore, domainEvents, new TestReturnToSenderDequeuer(new ReturnToSender(BodyStorage, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization()), - new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance), + new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System), TestRetryMetrics.Create(), new Lazy(() => sender), new RecordingMessageActionAuditLog(), NullLogger.Instance); @@ -426,7 +426,7 @@ async Task CreateAFailedMessageAndMarkAsPartOfRetryBatch(RetryingManager retryMa class CustomRetriesGateway : RetriesGateway { public CustomRetriesGateway(bool progressToStaged, IRetryBatchStore store, RetryingManager retryManager) - : base(store, retryManager, TestRetryMetrics.Create(), NullLogger.Instance) + : base(store, retryManager, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System) { this.progressToStaged = progressToStaged; } diff --git a/src/ServiceControl.Persistence/IPersistenceConfiguration.cs b/src/ServiceControl.Persistence/IPersistenceConfiguration.cs index 5b2e794562..3c85cc8558 100644 --- a/src/ServiceControl.Persistence/IPersistenceConfiguration.cs +++ b/src/ServiceControl.Persistence/IPersistenceConfiguration.cs @@ -4,6 +4,9 @@ public interface IPersistenceConfiguration { + /// Whether maintenance mode has anything to expose: only RavenDB, which hosts the database in-process. + bool SupportsMaintenanceMode { get; } + PersistenceSettings CreateSettings(SettingsRootNamespace settingsRootNamespace); IPersistence Create(PersistenceSettings settings); } diff --git a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs index 90f933d506..bef83dd441 100644 --- a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs +++ b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs @@ -8,12 +8,13 @@ public class InMemoryArchive // in memory { - public InMemoryArchive(string requestId, ArchiveType archiveType, IDomainEvents domainEvents, ArchiveMetrics? metrics = null) + public InMemoryArchive(string requestId, ArchiveType archiveType, IDomainEvents domainEvents, TimeProvider timeProvider, ArchiveMetrics? metrics = null) { RequestId = requestId; ArchiveType = archiveType; this.domainEvents = domainEvents; operationMetrics = metrics?.CreateOperation(ArchiveOperationKind.Archive); + this.timeProvider = timeProvider; } public int TotalNumberOfMessages { get; set; } @@ -63,7 +64,7 @@ public Task BatchArchived(int numberOfMessagesArchivedInBatch, CancellationToken ArchiveState = ArchiveState.ArchiveProgressing; NumberOfMessagesArchived += numberOfMessagesArchivedInBatch; CurrentBatch++; - Last = DateTime.UtcNow; + Last = timeProvider.GetUtcNow().UtcDateTime; operationMetrics?.BatchCompleted(numberOfMessagesArchivedInBatch); return domainEvents.Raise(new ArchiveOperationBatchCompleted @@ -80,7 +81,7 @@ public Task FinalizeArchive(CancellationToken cancellationToken = default) { ArchiveState = ArchiveState.ArchiveFinalizing; NumberOfMessagesArchived = TotalNumberOfMessages; - Last = DateTime.UtcNow; + Last = timeProvider.GetUtcNow().UtcDateTime; operationMetrics?.Finalizing(); return domainEvents.Raise(new ArchiveOperationFinalizing @@ -97,8 +98,9 @@ public Task Complete(CancellationToken cancellationToken = default) { ArchiveState = ArchiveState.ArchiveCompleted; NumberOfMessagesArchived = TotalNumberOfMessages; - CompletionTime = DateTime.UtcNow; - Last = DateTime.UtcNow; + var completedAt = timeProvider.GetUtcNow().UtcDateTime; + CompletionTime = completedAt; + Last = completedAt; operationMetrics?.Completed(); return domainEvents.Raise(new ArchiveOperationCompleted @@ -120,5 +122,6 @@ public bool NeedsAcknowledgement() IDomainEvents domainEvents; readonly ArchiveOperationMetrics? operationMetrics; + readonly TimeProvider timeProvider; } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs index bbba1d009b..375f98f38d 100644 --- a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs +++ b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs @@ -8,12 +8,13 @@ public class InMemoryUnarchive // in memory { - public InMemoryUnarchive(string requestId, ArchiveType archiveType, IDomainEvents domainEvents, ArchiveMetrics? metrics = null) + public InMemoryUnarchive(string requestId, ArchiveType archiveType, IDomainEvents domainEvents, TimeProvider timeProvider, ArchiveMetrics? metrics = null) { RequestId = requestId; ArchiveType = archiveType; this.domainEvents = domainEvents; operationMetrics = metrics?.CreateOperation(ArchiveOperationKind.Unarchive); + this.timeProvider = timeProvider; } public int TotalNumberOfMessages { get; set; } @@ -63,7 +64,7 @@ public Task BatchUnarchived(int numberOfMessagesUnarchivedInBatch, CancellationT ArchiveState = ArchiveState.ArchiveProgressing; NumberOfMessagesUnarchived += numberOfMessagesUnarchivedInBatch; CurrentBatch++; - Last = DateTime.UtcNow; + Last = timeProvider.GetUtcNow().UtcDateTime; operationMetrics?.BatchCompleted(numberOfMessagesUnarchivedInBatch); return domainEvents.Raise(new UnarchiveOperationBatchCompleted @@ -80,7 +81,7 @@ public Task FinalizeUnarchive(CancellationToken cancellationToken = default) { ArchiveState = ArchiveState.ArchiveFinalizing; NumberOfMessagesUnarchived = TotalNumberOfMessages; - Last = DateTime.UtcNow; + Last = timeProvider.GetUtcNow().UtcDateTime; operationMetrics?.Finalizing(); return domainEvents.Raise(new UnarchiveOperationFinalizing @@ -97,8 +98,9 @@ public Task Complete(CancellationToken cancellationToken = default) { ArchiveState = ArchiveState.ArchiveCompleted; NumberOfMessagesUnarchived = TotalNumberOfMessages; - CompletionTime = DateTime.UtcNow; - Last = DateTime.UtcNow; + var completedAt = timeProvider.GetUtcNow().UtcDateTime; + CompletionTime = completedAt; + Last = completedAt; operationMetrics?.Completed(); return domainEvents.Raise(new UnarchiveOperationCompleted @@ -120,5 +122,6 @@ internal bool NeedsAcknowledgement() IDomainEvents domainEvents; readonly ArchiveOperationMetrics? operationMetrics; + readonly TimeProvider timeProvider; } } \ No newline at end of file diff --git a/src/ServiceControl.UnitTests/MessageRedirects/MessageRedirectsControllerClockTests.cs b/src/ServiceControl.UnitTests/MessageRedirects/MessageRedirectsControllerClockTests.cs new file mode 100644 index 0000000000..9ccd1c6ffe --- /dev/null +++ b/src/ServiceControl.UnitTests/MessageRedirects/MessageRedirectsControllerClockTests.cs @@ -0,0 +1,78 @@ +namespace ServiceControl.UnitTests.MessageRedirects; + +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using NServiceBus.Testing; +using NUnit.Framework; +using ServiceControl.MessageRedirects.Api; +using ServiceControl.Persistence.MessageRedirects; +using ServiceControl.UnitTests.Operations; + +[TestFixture] +public class MessageRedirectsControllerClockTests +{ + static readonly DateTime FixedNow = new(2026, 3, 4, 5, 6, 7, DateTimeKind.Utc); + + [Test] + public async Task Creating_a_redirect_stamps_the_injected_clock() + { + var store = new RecordingRedirectsStore(); + var controller = NewController(store); + + await controller.NewRedirects(new MessageRedirectsController.MessageRedirectRequest { FromPhysicalAddress = "A", ToPhysicalAddress = "B" }); + + Assert.That(store.Redirects.Single().LastModified, Is.EqualTo(FixedNow)); + } + + [Test] + public async Task Updating_a_redirect_stamps_the_injected_clock() + { + var existing = new MessageRedirect + { + FromPhysicalAddress = "A", + ToPhysicalAddress = "B", + LastModified = FixedNow.AddDays(-10) + }; + + var store = new RecordingRedirectsStore(); + store.Redirects.Add(existing); + var controller = NewController(store); + + await controller.UpdateRedirect(existing.MessageRedirectId, new MessageRedirectsController.MessageRedirectRequest { ToPhysicalAddress = "C" }); + + Assert.That(store.Redirects.Single().LastModified, Is.EqualTo(FixedNow)); + } + + static MessageRedirectsController NewController(IMessageRedirectsDataStore store) => + new(new TestableMessageSession(), store, new FakeDomainEvents(), new FixedClock(FixedNow)); + + sealed class RecordingRedirectsStore : IMessageRedirectsDataStore + { + public List Redirects { get; } = []; + + public Task> GetRedirects(CancellationToken cancellationToken = default) => + Task.FromResult>(Redirects); + + public Task AddRedirect(MessageRedirect redirect, CancellationToken cancellationToken = default) + { + Redirects.Add(redirect); + return Task.CompletedTask; + } + + public Task UpdateRedirect(MessageRedirect redirect, CancellationToken cancellationToken = default) => Task.CompletedTask; + + public Task RemoveRedirect(MessageRedirect redirect, CancellationToken cancellationToken = default) + { + Redirects.Remove(redirect); + return Task.CompletedTask; + } + } + + sealed class FixedClock(DateTime now) : TimeProvider + { + public override DateTimeOffset GetUtcNow() => new(now, TimeSpan.Zero); + } +} diff --git a/src/ServiceControl.UnitTests/Recoverability/FailureGroupsRetryControllerAuditTests.cs b/src/ServiceControl.UnitTests/Recoverability/FailureGroupsRetryControllerAuditTests.cs index af91ad48ab..9d55ffba33 100644 --- a/src/ServiceControl.UnitTests/Recoverability/FailureGroupsRetryControllerAuditTests.cs +++ b/src/ServiceControl.UnitTests/Recoverability/FailureGroupsRetryControllerAuditTests.cs @@ -22,8 +22,8 @@ public async Task Emits_group_retry_operation_entry() var session = new TestableMessageSession(); var audit = new RecordingMessageActionAuditLog(); var user = new AuditUser("alice-sub", "Alice"); - var retryingManager = new RetryingManager(new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); - var controller = new FailureGroupsRetryController(session, retryingManager, new StubCurrentUserAccessor(user), audit); + var retryingManager = new RetryingManager(new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); + var controller = new FailureGroupsRetryController(session, retryingManager, new StubCurrentUserAccessor(user), audit, TimeProvider.System); await controller.ArchiveGroupErrors("group-42"); @@ -41,9 +41,9 @@ public async Task Group_retry_skipped_as_already_in_progress_is_not_audited() { var session = new TestableMessageSession(); var audit = new RecordingMessageActionAuditLog(); - var retryingManager = new RetryingManager(new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var retryingManager = new RetryingManager(new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await retryingManager.Preparing("group-42", RetryType.FailureGroup, totalNumberOfMessages: 10); - var controller = new FailureGroupsRetryController(session, retryingManager, new StubCurrentUserAccessor(new AuditUser("alice-sub", "Alice")), audit); + var controller = new FailureGroupsRetryController(session, retryingManager, new StubCurrentUserAccessor(new AuditUser("alice-sub", "Alice")), audit, TimeProvider.System); await controller.ArchiveGroupErrors("group-42"); diff --git a/src/ServiceControl.UnitTests/Recoverability/OperationProgressClockTests.cs b/src/ServiceControl.UnitTests/Recoverability/OperationProgressClockTests.cs new file mode 100644 index 0000000000..533a5aa283 --- /dev/null +++ b/src/ServiceControl.UnitTests/Recoverability/OperationProgressClockTests.cs @@ -0,0 +1,116 @@ +namespace ServiceControl.UnitTests.Recoverability; + +using System; +using System.Threading.Tasks; +using Microsoft.Extensions.Logging.Abstractions; +using NUnit.Framework; +using ServiceControl.Persistence; +using ServiceControl.Recoverability; +using ServiceControl.UnitTests.Operations; + +[TestFixture] +public class OperationProgressClockTests +{ + static readonly DateTime FixedNow = new(2026, 3, 4, 5, 6, 7, DateTimeKind.Utc); + + [Test] + public async Task Archive_progress_and_completion_come_from_the_injected_clock() + { + var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), new FixedClock(FixedNow)); + + await archive.BatchArchived(1); + Assert.That(archive.Last, Is.EqualTo(FixedNow)); + + await archive.FinalizeArchive(); + Assert.That(archive.Last, Is.EqualTo(FixedNow)); + + await archive.Complete(); + using (Assert.EnterMultipleScope()) + { + Assert.That(archive.CompletionTime, Is.EqualTo(FixedNow)); + Assert.That(archive.Last, Is.EqualTo(FixedNow)); + } + } + + [Test] + public async Task Unarchive_progress_and_completion_come_from_the_injected_clock() + { + var unarchive = new InMemoryUnarchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), new FixedClock(FixedNow)); + + await unarchive.BatchUnarchived(1); + Assert.That(unarchive.Last, Is.EqualTo(FixedNow)); + + await unarchive.FinalizeUnarchive(); + Assert.That(unarchive.Last, Is.EqualTo(FixedNow)); + + await unarchive.Complete(); + using (Assert.EnterMultipleScope()) + { + Assert.That(unarchive.CompletionTime, Is.EqualTo(FixedNow)); + Assert.That(unarchive.Last, Is.EqualTo(FixedNow)); + } + } + + [Test] + public async Task Completing_an_archive_reads_the_clock_once() + { + var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), new TickingClock(FixedNow)); + + await archive.Complete(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(archive.CompletionTime, Is.EqualTo(FixedNow)); + Assert.That(archive.Last, Is.EqualTo(FixedNow)); + } + } + + [Test] + public async Task Completing_an_unarchive_reads_the_clock_once() + { + var unarchive = new InMemoryUnarchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), new TickingClock(FixedNow)); + + await unarchive.Complete(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(unarchive.CompletionTime, Is.EqualTo(FixedNow)); + Assert.That(unarchive.Last, Is.EqualTo(FixedNow)); + } + } + + [Test] + public async Task Retry_completion_comes_from_the_injected_clock() + { + var retry = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, new FixedClock(FixedNow)); + + await retry.Prepare(1000); + await retry.PrepareBatch(1000); + await retry.Forwarding(); + await retry.BatchForwarded(1000); + + using (Assert.EnterMultipleScope()) + { + Assert.That(retry.RetryState, Is.EqualTo(RetryState.Completed)); + Assert.That(retry.CompletionTime, Is.EqualTo(FixedNow)); + } + } + + sealed class FixedClock(DateTime now) : TimeProvider + { + public override DateTimeOffset GetUtcNow() => new(now, TimeSpan.Zero); + } + + // Moves on every read, so a member that reads twice cannot get the same value twice. + sealed class TickingClock(DateTime start) : TimeProvider + { + DateTime next = start; + + public override DateTimeOffset GetUtcNow() + { + var value = next; + next = next.AddSeconds(1); + return new DateTimeOffset(value, TimeSpan.Zero); + } + } +} diff --git a/src/ServiceControl.UnitTests/Recoverability/RetryOperationTests.cs b/src/ServiceControl.UnitTests/Recoverability/RetryOperationTests.cs index 95c99d9428..2e7141842f 100644 --- a/src/ServiceControl.UnitTests/Recoverability/RetryOperationTests.cs +++ b/src/ServiceControl.UnitTests/Recoverability/RetryOperationTests.cs @@ -14,7 +14,7 @@ public class RetryOperationTests [Test] public async Task Wait_should_set_wait_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await summary.Wait(DateTime.UtcNow, "FailureGroup1"); using (Assert.EnterMultipleScope()) { @@ -30,7 +30,7 @@ public async Task Wait_should_set_wait_state() [Test] public void Fail_should_set_failed() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); summary.Fail(); Assert.That(summary.Failed, Is.True); } @@ -38,7 +38,7 @@ public void Fail_should_set_failed() [Test] public async Task Prepare_should_set_prepare_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await summary.Prepare(1000); using (Assert.EnterMultipleScope()) { @@ -51,7 +51,7 @@ public async Task Prepare_should_set_prepare_state() [Test] public async Task Prepared_batch_should_set_prepare_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await summary.Prepare(1000); await summary.PrepareBatch(1000); using (Assert.EnterMultipleScope()) @@ -65,7 +65,7 @@ public async Task Prepared_batch_should_set_prepare_state() [Test] public async Task Forwarding_should_set_forwarding_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await summary.Prepare(1000); await summary.PrepareBatch(1000); await summary.Forwarding(); @@ -81,7 +81,7 @@ public async Task Forwarding_should_set_forwarding_state() [Test] public async Task Batch_forwarded_should_set_forwarding_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await summary.Prepare(1000); await summary.PrepareBatch(1000); await summary.Forwarding(); @@ -99,7 +99,7 @@ public async Task Batch_forwarded_should_set_forwarding_state() public async Task Should_raise_domain_events() { var domainEvents = new FakeDomainEvents(); - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, domainEvents, TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await summary.Prepare(1000); await summary.PrepareBatch(1000); await summary.Forwarding(); @@ -118,7 +118,7 @@ public async Task Should_raise_domain_events() [Test] public async Task Batch_forwarded_all_forwarded_should_set_completed_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await summary.Prepare(1000); await summary.PrepareBatch(1000); await summary.Forwarding(); @@ -135,7 +135,7 @@ public async Task Batch_forwarded_all_forwarded_should_set_completed_state() [Test] public async Task Skip_should_set_update_skipped_messages() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await summary.Wait(DateTime.UtcNow); await summary.Prepare(2000); await summary.PrepareBatch(1000); @@ -151,7 +151,7 @@ public async Task Skip_should_set_update_skipped_messages() [Test] public async Task Skip_should_complete_when_all_skipped() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await summary.Wait(DateTime.UtcNow); await summary.Prepare(1000); await summary.PrepareBatch(1000); @@ -167,7 +167,7 @@ public async Task Skip_should_complete_when_all_skipped() [Test] public async Task Skip_and_forward_combination_should_complete_when_done() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, TimeProvider.System); await summary.Wait(DateTime.UtcNow); await summary.Prepare(2000); await summary.PrepareBatch(1000); diff --git a/src/ServiceControl.UnitTests/Recoverability/RetryStartTimeTests.cs b/src/ServiceControl.UnitTests/Recoverability/RetryStartTimeTests.cs new file mode 100644 index 0000000000..4d2a3a3179 --- /dev/null +++ b/src/ServiceControl.UnitTests/Recoverability/RetryStartTimeTests.cs @@ -0,0 +1,63 @@ +#nullable enable +namespace ServiceControl.UnitTests.Recoverability; + +using System; +using System.Linq; +using System.Threading.Tasks; +using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Time.Testing; +using NServiceBus.Testing; +using NUnit.Framework; +using ServiceControl.Infrastructure.Auth; +using ServiceControl.Persistence; +using ServiceControl.Recoverability; +using ServiceControl.Recoverability.API; +using ServiceControl.UnitTests.Operations; + +[TestFixture] +public class RetryStartTimeTests +{ + static readonly DateTimeOffset ClockStart = new(2020, 1, 1, 0, 0, 0, TimeSpan.Zero); + + [Test] + public async Task Group_retry_takes_its_start_time_from_the_injected_clock() + { + var clock = new FakeTimeProvider(ClockStart); + var session = new TestableMessageSession(); + var retryingManager = new RetryingManager(new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, clock); + + await NewController(session, retryingManager, clock).ArchiveGroupErrors("group-42"); + + var sent = (RetryAllInGroup)session.SentMessages.Single().Message; + using (Assert.EnterMultipleScope()) + { + Assert.That(sent.Started, Is.EqualTo(ClockStart.UtcDateTime)); + Assert.That(retryingManager.GetStatusForRetryOperation("group-42", RetryType.FailureGroup).Started, Is.EqualTo(ClockStart.UtcDateTime)); + } + } + + [Test] + public async Task A_completed_group_retry_never_finishes_before_it_started() + { + var clock = new FakeTimeProvider(ClockStart); + var retryingManager = new RetryingManager(new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance, clock); + + await NewController(new TestableMessageSession(), retryingManager, clock).ArchiveGroupErrors("group-42"); + + clock.Advance(TimeSpan.FromMinutes(5)); + await retryingManager.Preparing("group-42", RetryType.FailureGroup, totalNumberOfMessages: 1); + await retryingManager.PreparedBatch("group-42", RetryType.FailureGroup, numberOfMessagesPrepared: 1); + await retryingManager.Forwarding("group-42", RetryType.FailureGroup); + await retryingManager.ForwardedBatch("group-42", RetryType.FailureGroup, numberOfMessagesForwarded: 1); + + var operation = retryingManager.GetStatusForRetryOperation("group-42", RetryType.FailureGroup); + using (Assert.EnterMultipleScope()) + { + Assert.That(operation.Started, Is.EqualTo(ClockStart.UtcDateTime)); + Assert.That(operation.CompletionTime, Is.EqualTo(ClockStart.UtcDateTime.AddMinutes(5))); + } + } + + static FailureGroupsRetryController NewController(TestableMessageSession session, RetryingManager retryingManager, TimeProvider clock) => + new(session, retryingManager, new StubCurrentUserAccessor(new AuditUser("alice-sub", "Alice")), new RecordingMessageActionAuditLog(), clock); +} diff --git a/src/ServiceControl/HostApplicationBuilderExtensions.cs b/src/ServiceControl/HostApplicationBuilderExtensions.cs index b511f16ac2..9b682fe1a6 100644 --- a/src/ServiceControl/HostApplicationBuilderExtensions.cs +++ b/src/ServiceControl/HostApplicationBuilderExtensions.cs @@ -70,6 +70,9 @@ public static void AddServiceControl(this IHostApplicationBuilder hostBuilder, S transportCustomization.AddTransportForPrimary(services, transportSettings); services.Configure(options => options.ShutdownTimeout = settings.ShutdownTimeout); + + services.TryAddSingleton(TimeProvider.System); + services.AddSingleton(); // Message-action audit trail. Registered here rather than in AddServiceControlAuthorization diff --git a/src/ServiceControl/MessageRedirects/Api/MessageRedirectsController.cs b/src/ServiceControl/MessageRedirects/Api/MessageRedirectsController.cs index b5f4dff114..fc52d6fbc0 100644 --- a/src/ServiceControl/MessageRedirects/Api/MessageRedirectsController.cs +++ b/src/ServiceControl/MessageRedirects/Api/MessageRedirectsController.cs @@ -25,7 +25,8 @@ public class MessageRedirectsController( IMessageSession session, IMessageRedirectsDataStore store, - IDomainEvents events) + IDomainEvents events, + TimeProvider timeProvider) : ControllerBase { [Authorize(Policy = Permissions.ErrorRedirectsManage)] @@ -38,11 +39,13 @@ public async Task NewRedirects(MessageRedirectRequest request, Ca return BadRequest(); } + var now = timeProvider.GetUtcNow().UtcDateTime; + var messageRedirect = new MessageRedirect { FromPhysicalAddress = request.FromPhysicalAddress, ToPhysicalAddress = request.ToPhysicalAddress, - LastModified = DateTime.UtcNow + LastModified = now }; var redirects = await store.GetRedirects(cancellationToken); @@ -87,7 +90,7 @@ await session.SendLocal(new RetryPendingMessages { QueueAddress = messageRedirect.FromPhysicalAddress, PeriodFrom = DateTime.MinValue, - PeriodTo = DateTime.UtcNow + PeriodTo = now }, cancellationToken); } @@ -129,7 +132,7 @@ public async Task UpdateRedirect(Guid messageRedirectId, MessageR ToPhysicalAddress = messageRedirect.ToPhysicalAddress = request.ToPhysicalAddress }; - messageRedirect.LastModified = DateTime.UtcNow; + messageRedirect.LastModified = timeProvider.GetUtcNow().UtcDateTime; await store.UpdateRedirect(messageRedirect, cancellationToken); diff --git a/src/ServiceControl/Persistence/PersistenceFactory.cs b/src/ServiceControl/Persistence/PersistenceFactory.cs index bc1c655cd8..892920cf18 100644 --- a/src/ServiceControl/Persistence/PersistenceFactory.cs +++ b/src/ServiceControl/Persistence/PersistenceFactory.cs @@ -10,6 +10,11 @@ public static IPersistence Create(Settings settings, bool maintenanceMode = fals { var persistenceConfiguration = CreatePersistenceConfiguration(settings); + if (maintenanceMode && !persistenceConfiguration.SupportsMaintenanceMode) + { + throw new Exception($"Maintenance mode is not supported by the {settings.PersistenceType} persister. It is only available on RavenDB, where it starts the embedded database so RavenDB Studio can be used."); + } + //HINT: This is false when executed from acceptance tests settings.PersisterSpecificSettings ??= persistenceConfiguration.CreateSettings(Settings.SettingsRootNamespace); settings.PersisterSpecificSettings.MaintenanceMode = maintenanceMode; diff --git a/src/ServiceControl/Recoverability/API/FailureGroupsRetryController.cs b/src/ServiceControl/Recoverability/API/FailureGroupsRetryController.cs index c2efea71bb..9459d5879e 100644 --- a/src/ServiceControl/Recoverability/API/FailureGroupsRetryController.cs +++ b/src/ServiceControl/Recoverability/API/FailureGroupsRetryController.cs @@ -16,14 +16,15 @@ public class FailureGroupsRetryController( IMessageSession bus, RetryingManager retryingManager, ICurrentUserAccessor userAccessor, - IMessageActionAuditLog auditLog) : ControllerBase + IMessageActionAuditLog auditLog, + TimeProvider timeProvider) : ControllerBase { [Authorize(Policy = Permissions.ErrorRecoverabilityGroupsRetry)] [Route("recoverability/groups/{groupId:required:minlength(1)}/errors/retry")] [HttpPost] public async Task ArchiveGroupErrors(string groupId, CancellationToken cancellationToken = default) { - var started = DateTime.UtcNow; + var started = timeProvider.GetUtcNow().UtcDateTime; if (!retryingManager.IsOperationInProgressFor(groupId, RetryType.FailureGroup)) { diff --git a/src/ServiceControl/Recoverability/Retrying/Handlers/RetryAllInGroupHandler.cs b/src/ServiceControl/Recoverability/Retrying/Handlers/RetryAllInGroupHandler.cs index b32bfc3f7d..d2a9cfc247 100644 --- a/src/ServiceControl/Recoverability/Retrying/Handlers/RetryAllInGroupHandler.cs +++ b/src/ServiceControl/Recoverability/Retrying/Handlers/RetryAllInGroupHandler.cs @@ -9,7 +9,7 @@ namespace ServiceControl.Recoverability using ServiceControl.Persistence.Recoverability; [Handler] - class RetryAllInGroupHandler(RetriesGateway retries, RetryingManager retryingManager, IArchiveMessages archiver, IGroupsDataStore dataStore, ILogger logger) + class RetryAllInGroupHandler(RetriesGateway retries, RetryingManager retryingManager, IArchiveMessages archiver, IGroupsDataStore dataStore, ILogger logger, TimeProvider timeProvider) : IHandleMessages { public async Task Handle(RetryAllInGroup message, IMessageHandlerContext context) @@ -35,7 +35,7 @@ public async Task Handle(RetryAllInGroup message, IMessageHandlerContext context originator = group.Title; } - var started = message.Started ?? DateTime.UtcNow; + var started = message.Started ?? timeProvider.GetUtcNow().UtcDateTime; await retryingManager.Wait(message.GroupId, RetryType.FailureGroup, started, originator, group?.Type, group?.Last, context.CancellationToken); var (user, operationId) = AuditHeaders.Read(context.MessageHeaders); diff --git a/src/ServiceControl/Recoverability/Retrying/InMemoryRetry.cs b/src/ServiceControl/Recoverability/Retrying/InMemoryRetry.cs index 98d203bd8c..9aef908c44 100644 --- a/src/ServiceControl/Recoverability/Retrying/InMemoryRetry.cs +++ b/src/ServiceControl/Recoverability/Retrying/InMemoryRetry.cs @@ -10,7 +10,7 @@ public class InMemoryRetry { - public InMemoryRetry(string requestId, RetryType retryType, IDomainEvents domainEvents, RetryMetrics metrics, ILogger logger) + public InMemoryRetry(string requestId, RetryType retryType, IDomainEvents domainEvents, RetryMetrics metrics, ILogger logger, TimeProvider timeProvider) { RequestId = requestId; RetryType = retryType; @@ -18,6 +18,7 @@ public InMemoryRetry(string requestId, RetryType retryType, IDomainEvents domain this.metrics = metrics; this.logger = logger; operationStartTimestamp = metrics.GetTimestamp(); + this.timeProvider = timeProvider; } public string RequestId { get; } @@ -166,7 +167,7 @@ async Task CheckForCompletion(CancellationToken cancellationToken) } RetryState = RetryState.Completed; - CompletionTime = DateTime.UtcNow; + CompletionTime = timeProvider.GetUtcNow().UtcDateTime; metrics.RecordOperationCompleted(RetryType, operationStartTimestamp, Failed); await domainEvents.Raise(new RetryOperationCompleted @@ -226,5 +227,6 @@ public bool IsInProgress() IDomainEvents domainEvents; readonly RetryMetrics metrics; readonly ILogger logger; + readonly TimeProvider timeProvider; } } \ No newline at end of file diff --git a/src/ServiceControl/Recoverability/Retrying/RetriesGateway.cs b/src/ServiceControl/Recoverability/Retrying/RetriesGateway.cs index 3b183defc0..f749cdd8be 100644 --- a/src/ServiceControl/Recoverability/Retrying/RetriesGateway.cs +++ b/src/ServiceControl/Recoverability/Retrying/RetriesGateway.cs @@ -15,12 +15,13 @@ namespace ServiceControl.Recoverability class RetriesGateway { - public RetriesGateway(IRetryBatchStore store, RetryingManager operationManager, RetryMetrics metrics, ILogger logger) + public RetriesGateway(IRetryBatchStore store, RetryingManager operationManager, RetryMetrics metrics, ILogger logger, TimeProvider timeProvider) { this.store = store; this.operationManager = operationManager; this.metrics = metrics; this.logger = logger; + this.timeProvider = timeProvider; metrics.ObservePendingBulkRequests(() => bulkRequests.Count); } @@ -36,7 +37,7 @@ public async Task StartRetryForSingleMessage(string uniqueMessageId, AuditUser? using var preparation = metrics.BeginPreparation(retryType, cancellationToken); await operationManager.Preparing(requestId, retryType, numberOfMessages, cancellationToken); - await AssignMessagesToBatch(requestId, retryType, new[] { uniqueMessageId }, DateTime.UtcNow, cancellationToken, initiatedBy: initiatedBy, operationId: operationId); + await AssignMessagesToBatch(requestId, retryType, new[] { uniqueMessageId }, timeProvider.GetUtcNow().UtcDateTime, cancellationToken, initiatedBy: initiatedBy, operationId: operationId); await operationManager.PreparedBatch(requestId, retryType, numberOfMessages, cancellationToken); preparation.Complete(); @@ -53,7 +54,7 @@ public async Task StartRetryForMessageSelection(string[] uniqueMessageIds, Audit using var preparation = metrics.BeginPreparation(retryType, cancellationToken); await operationManager.Preparing(requestId, retryType, numberOfMessages, cancellationToken); - await AssignMessagesToBatch(requestId, retryType, uniqueMessageIds, DateTime.UtcNow, cancellationToken, initiatedBy: initiatedBy, operationId: operationId); + await AssignMessagesToBatch(requestId, retryType, uniqueMessageIds, timeProvider.GetUtcNow().UtcDateTime, cancellationToken, initiatedBy: initiatedBy, operationId: operationId); await operationManager.PreparedBatch(requestId, retryType, numberOfMessages, cancellationToken); preparation.Complete(); @@ -132,21 +133,21 @@ static string GetBatchName(int pageNum, int totalPages, string context) public void StartRetryForAllMessages(AuditUser? initiatedBy = null, string operationId = null) { - var item = new RetryForAllMessages(initiatedBy, operationId); + var item = new RetryForAllMessages(timeProvider.GetUtcNow().UtcDateTime, initiatedBy, operationId); logger.LogInformation("Enqueuing index based bulk retry '{Item}'", item); bulkRequests.Enqueue(item); } public void StartRetryForEndpoint(string endpoint, AuditUser? initiatedBy = null, string operationId = null) { - var item = new RetryForEndpoint(endpoint, initiatedBy, operationId); + var item = new RetryForEndpoint(endpoint, timeProvider.GetUtcNow().UtcDateTime, initiatedBy, operationId); logger.LogInformation("Enqueuing index based bulk retry '{Item}'", item); bulkRequests.Enqueue(item); } public void StartRetryForFailedQueueAddress(string failedQueueAddress, FailedMessageStatus status, AuditUser? initiatedBy = null, string operationId = null) { - var item = new RetryForFailedQueueAddress(failedQueueAddress, status, initiatedBy, operationId); + var item = new RetryForFailedQueueAddress(failedQueueAddress, status, timeProvider.GetUtcNow().UtcDateTime, initiatedBy, operationId); logger.LogInformation("Enqueuing index based bulk retry '{Item}'", item); bulkRequests.Enqueue(item); } @@ -164,6 +165,7 @@ public void EnqueueRetryForFailureGroup(RetryForFailureGroup item) const int BatchSize = 1000; readonly ILogger logger; + readonly TimeProvider timeProvider; public abstract class BulkRetryRequest { @@ -234,7 +236,7 @@ Task Process(string uniqueMessageId, DateTime latestTimeOfFailure, CancellationT class RetryForAllMessages : BulkRetryRequest { - public RetryForAllMessages(AuditUser? initiatedBy = null, string operationId = null) : base(requestId: "All", RetryType.All, DateTime.UtcNow, "all messages", initiatedBy, operationId) + public RetryForAllMessages(DateTime startTime, AuditUser? initiatedBy = null, string operationId = null) : base(requestId: "All", RetryType.All, startTime, "all messages", initiatedBy, operationId) { } @@ -248,7 +250,7 @@ class RetryForEndpoint : BulkRetryRequest { public string Endpoint { get; } - public RetryForEndpoint(string endpoint, AuditUser? initiatedBy = null, string operationId = null) : base(requestId: endpoint, RetryType.AllForEndpoint, DateTime.UtcNow, originator: $"all messages for endpoint {endpoint}", initiatedBy, operationId) + public RetryForEndpoint(string endpoint, DateTime startTime, AuditUser? initiatedBy = null, string operationId = null) : base(requestId: endpoint, RetryType.AllForEndpoint, startTime, originator: $"all messages for endpoint {endpoint}", initiatedBy, operationId) { Endpoint = endpoint; } @@ -287,9 +289,10 @@ class RetryForFailedQueueAddress : BulkRetryRequest public RetryForFailedQueueAddress( string failedQueueAddress, FailedMessageStatus status, + DateTime startTime, AuditUser? initiatedBy = null, string operationId = null - ) : base(requestId: failedQueueAddress, RetryType.ByQueueAddress, DateTime.UtcNow, originator: $"all messages for failed queue address '{failedQueueAddress}'", initiatedBy, operationId) + ) : base(requestId: failedQueueAddress, RetryType.ByQueueAddress, startTime, originator: $"all messages for failed queue address '{failedQueueAddress}'", initiatedBy, operationId) { FailedQueueAddress = failedQueueAddress; Status = status; diff --git a/src/ServiceControl/Recoverability/Retrying/RetryingManager.cs b/src/ServiceControl/Recoverability/Retrying/RetryingManager.cs index e766a6f982..f7697ed2ff 100644 --- a/src/ServiceControl/Recoverability/Retrying/RetryingManager.cs +++ b/src/ServiceControl/Recoverability/Retrying/RetryingManager.cs @@ -12,11 +12,12 @@ public class RetryingManager { - public RetryingManager(IDomainEvents domainEvents, RetryMetrics metrics, ILogger logger) + public RetryingManager(IDomainEvents domainEvents, RetryMetrics metrics, ILogger logger, TimeProvider timeProvider) { this.domainEvents = domainEvents; this.metrics = metrics; this.logger = logger; + this.timeProvider = timeProvider; metrics.ObserveOperationsInProgress(() => retryOperations.Values.Select(operation => (operation.RetryType, operation.RetryState))); } @@ -97,7 +98,7 @@ InMemoryRetry GetOrCreate(RetryType retryType, string requestId) ArgumentException.ThrowIfNullOrWhiteSpace(requestId); var key = InMemoryRetry.MakeOperationId(requestId, retryType); - return retryOperations.GetOrAdd(key, _ => new InMemoryRetry(requestId, retryType, domainEvents, metrics, logger)); + return retryOperations.GetOrAdd(key, _ => new InMemoryRetry(requestId, retryType, domainEvents, metrics, logger, timeProvider)); } public InMemoryRetry GetStatusForRetryOperation(string requestId, RetryType retryType) @@ -110,6 +111,7 @@ public InMemoryRetry GetStatusForRetryOperation(string requestId, RetryType retr IDomainEvents domainEvents; readonly RetryMetrics metrics; readonly ILogger logger; + readonly TimeProvider timeProvider; ConcurrentDictionary retryOperations = new ConcurrentDictionary(); } } \ No newline at end of file