Skip to content
Open
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,9 @@

<!-- Error ingestion only mode is supported on SQL Server and PostgreSQL storage only. -->
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\When_hosting_error_ingestion_only.cs" />

<!-- RavenDB is the one persister that supports maintenance mode; StartupModeTests covers it. -->
<Compile Remove="..\ServiceControl.AcceptanceTests\MaintenanceModeTests.cs" />
</ItemGroup>

<ItemGroup>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<TimeProvider>(), Is.Not.Null);

await host.StopAsync();
}

Expand Down
55 changes: 55 additions & 0 deletions src/ServiceControl.AcceptanceTests/ClockRegistrationTests.cs
Original file line number Diff line number Diff line change
@@ -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<TimeProvider>(), Is.SameAs(hostClock));
}

sealed class StubTimeProvider : TimeProvider;
}
}
31 changes: 31 additions & 0 deletions src/ServiceControl.AcceptanceTests/MaintenanceModeTests.cs
Original file line number Diff line number Diff line change
@@ -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<Exception>(() => new ServiceCollection().AddPersistence(settings, maintenanceMode: true));

Assert.That(exception.Message, Does.Contain("Maintenance mode is not supported").And.Contain(StorageConfiguration.PersistenceType));
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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<MinimumRequiredStorageState>();

services.AddSingleton<IServiceControlSubscriptionStorage, SubscriptionStorage>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,8 @@ namespace ServiceControl.Persistence.EFCore.Implementation;
/// <summary>
/// Every operation here has to update both StatusChangedAt and LastModified.
/// </summary>
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)
{
Expand All @@ -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)
Expand All @@ -38,7 +39,7 @@ public async Task<bool> 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)
Expand All @@ -65,7 +66,7 @@ public async Task<string[]> UnArchiveMessages(IEnumerable<string> 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
Expand Down Expand Up @@ -93,7 +94,7 @@ public async Task<string[]> 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
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,14 +9,14 @@ namespace ServiceControl.Persistence.EFCore.Implementation;
/// EFCore equivalent of the RavenDB <see cref="ArchivingManager"/>. Wraps the shared
/// <see cref="OperationsManager"/> singleton to manage in-memory archive progress state.
/// </summary>
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;
}

Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,14 +9,14 @@ namespace ServiceControl.Persistence.EFCore.Implementation;
/// EFCore equivalent of the RavenDB <see cref="UnarchivingManager"/>. Wraps the shared
/// <see cref="OperationsManager"/> singleton to manage in-memory unarchive progress state.
/// </summary>
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;
}

Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,8 +31,8 @@ ILogger<MessageArchiver> 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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
10 changes: 6 additions & 4 deletions src/ServiceControl.Persistence.RavenDB/ExpirationManager.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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;
}
Expand All @@ -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);

Expand Down
4 changes: 4 additions & 0 deletions src/ServiceControl.Persistence.RavenDB/RavenPersistence.cs
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -81,6 +83,8 @@ public void AddInstaller(IServiceCollection services)

void ConfigureLifecycle(IServiceCollection services)
{
services.TryAddSingleton(TimeProvider.System);

services.AddSingleton<PersistenceSettings>(settings);
services.AddSingleton(settings);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<ArchiveOperation> LoadArchiveOperation(IAsyncDocumentSession session, string groupId, ArchiveType archiveType, CancellationToken cancellationToken = default) => session.LoadAsync<ArchiveOperation>(ArchiveOperation.MakeId(groupId, archiveType), cancellationToken);

Expand All @@ -26,7 +26,7 @@ public async Task<ArchiveOperation> 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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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;
}

Expand All @@ -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;
Expand Down Expand Up @@ -106,6 +107,7 @@ void RemoveArchiveOperation(string requestId, ArchiveType archiveType)
}

IDomainEvents domainEvents;
readonly TimeProvider timeProvider;

OperationsManager operationsManager;

Expand Down
Loading
Loading