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