Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,9 @@

<!-- The file system body storage custom check is an EF Core persister feature; RavenDB stores bodies in the database. -->
<Compile Remove="..\ServiceControl.AcceptanceTests\Monitoring\CustomChecks\When_the_body_storage_check_is_reported.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;
Comment thread
warwickschroeder marked this conversation as resolved.

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;
Comment thread
warwickschroeder marked this conversation as resolved.

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