From 1395ee52104959ba0b63ec35862fdc309d935f1f Mon Sep 17 00:00:00 2001 From: Rhys Bevilaqua Date: Tue, 21 Jul 2026 16:40:11 +0800 Subject: [PATCH 1/4] Add endpoint settings --- .../Recoverability/QueueAddress.cs | 8 ---- .../DbContexts/ServiceControlDbContext.cs | 1 + .../Entities/EndpointSettingsEntity.cs | 7 +++ .../Entities/FailedMessageEntity.cs | 9 ++++ .../EndpointSettingsEntity.cs | 14 ++++++ .../Implementation/DataStoreBase.cs | 27 +++++++++++ .../Implementation/EndpointSettingsStore.cs | 45 ++++++++++++++++--- .../Implementation/QueueAddressStore.cs | 9 ++-- .../TrialLicenseDataProvider.cs | 5 ++- .../IEndpointSettingsStore.cs | 2 +- 10 files changed, 106 insertions(+), 21 deletions(-) delete mode 100644 src/ServiceControl.Audit/Recoverability/QueueAddress.cs create mode 100644 src/ServiceControl.Persistence.EFCore/Entities/EndpointSettingsEntity.cs create mode 100644 src/ServiceControl.Persistence.EFCore/Entities/FailedMessageEntity.cs create mode 100644 src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsEntity.cs diff --git a/src/ServiceControl.Audit/Recoverability/QueueAddress.cs b/src/ServiceControl.Audit/Recoverability/QueueAddress.cs deleted file mode 100644 index e0b89747b3..0000000000 --- a/src/ServiceControl.Audit/Recoverability/QueueAddress.cs +++ /dev/null @@ -1,8 +0,0 @@ -namespace ServiceControl.Audit.Recoverability -{ - public class QueueAddress - { - public string PhysicalAddress { get; set; } - public int FailedMessageCount { get; set; } - } -} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs b/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs index 6a7e0a6fc9..20bd199a6a 100644 --- a/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs +++ b/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs @@ -6,6 +6,7 @@ namespace ServiceControl.Persistence.EFCore.DbContexts; public abstract class ServiceControlDbContext(DbContextOptions options) : DbContext(options) { + public DbSet EndpointSettings { get; set; } public DbSet KnownEndpoints { get; set; } public DbSet KnownEndpointsInsertOnly { get; set; } diff --git a/src/ServiceControl.Persistence.EFCore/Entities/EndpointSettingsEntity.cs b/src/ServiceControl.Persistence.EFCore/Entities/EndpointSettingsEntity.cs new file mode 100644 index 0000000000..2316c07245 --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Entities/EndpointSettingsEntity.cs @@ -0,0 +1,7 @@ +namespace ServiceControl.Persistence.EFCore.Entities; + +public class EndpointSettingsEntity +{ + public required string Name { get; set; } + public bool TrackInstances { get; set; } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Entities/FailedMessageEntity.cs b/src/ServiceControl.Persistence.EFCore/Entities/FailedMessageEntity.cs new file mode 100644 index 0000000000..ae60b4e8f5 --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Entities/FailedMessageEntity.cs @@ -0,0 +1,9 @@ +namespace ServiceControl.Persistence.EFCore.Entities; + +public class FailedMessageEntity +{ + public Guid Id { get; set; } + public string EndpointAddress + + //todo: incomplete, needs the other feature finished... +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsEntity.cs b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsEntity.cs new file mode 100644 index 0000000000..5c135c48af --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsEntity.cs @@ -0,0 +1,14 @@ +namespace ServiceControl.Persistence.EFCore.EntityConfigurations; + +using Entities; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Metadata.Builders; + +public class EndpointSettingsConfiguration : IEntityTypeConfiguration +{ + public void Configure(EntityTypeBuilder builder) + { + builder.ToTable("EndpointSettings"); + builder.HasKey(x => x.Name); + } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/DataStoreBase.cs b/src/ServiceControl.Persistence.EFCore/Implementation/DataStoreBase.cs index 157450c54b..cdafb21f2d 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/DataStoreBase.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/DataStoreBase.cs @@ -1,6 +1,7 @@ namespace ServiceControl.Persistence.EFCore.Implementation; using System; +using System.Runtime.CompilerServices; using System.Threading.Tasks; using DbContexts; using Microsoft.Extensions.DependencyInjection; @@ -32,6 +33,32 @@ protected async Task ExecuteWithDbContext(Func op await operation(dbContext); } + /// + /// Executes an operation with a scoped DbContext, without returning a result + /// + protected async IAsyncEnumerable ExecuteWithDbContext(Func> operation) + { + await using var scope = scopeFactory.CreateAsyncScope(); + var dbContext = scope.ServiceProvider.GetRequiredService(); + await foreach (var row in operation(dbContext)) + { + yield return row; + } + } + + /// + /// Executes an operation with a scoped DbContext, without returning a result + /// + protected async IAsyncEnumerable ExecuteWithDbContext(Func> operation, [EnumeratorCancellation] CancellationToken cancellationToken) + { + await using var scope = scopeFactory.CreateAsyncScope(); + var dbContext = scope.ServiceProvider.GetRequiredService(); + await foreach (var row in operation(dbContext).WithCancellation(cancellationToken)) + { + yield return row; + } + } + /// /// Creates a scope for operations that need to manage their own scope lifecycle (e.g., managers) /// diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs index 6de13c0fa4..65842f6828 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs @@ -1,13 +1,44 @@ namespace ServiceControl.Persistence.EFCore.Implementation; -public class EndpointSettingsStore : IEndpointSettingsStore +using Entities; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; + +public class EndpointSettingsStore(IServiceScopeFactory scopeFactory) : DataStoreBase(scopeFactory), IEndpointSettingsStore { - public IAsyncEnumerable GetAllEndpointSettings(CancellationToken cancellationToken) => - throw new NotImplementedException(); + public IAsyncEnumerable GetAllEndpointSettings(CancellationToken cancellationToken) + => ExecuteWithDbContext(context => context.EndpointSettings.Select(row => new EndpointSettings + { + Name = row.Name, + TrackInstances = row.TrackInstances + }).AsAsyncEnumerable(), cancellationToken); + + public Task UpdateEndpointSettings(EndpointSettings settings, CancellationToken token) => ExecuteWithDbContext(async context => + { + var entity = context.EndpointSettings.Find(settings.Name); + if (entity == null) + { + entity = new EndpointSettingsEntity() { Name = settings.Name, TrackInstances = settings.TrackInstances }; + context.EndpointSettings.Add(entity); + try + { + await context.SaveChangesAsync(token); + return; + } + catch (someexception e) when (condition) + { + //this failed because of key conflict so it's a race! + } - public Task UpdateEndpointSettings(EndpointSettings settings, CancellationToken token) => - throw new NotImplementedException(); + await context.Entry(entity).ReloadAsync(token); + entity.TrackInstances = settings.TrackInstances; + await context.SaveChangesAsync(token); + } + }); - public Task Delete(string name, CancellationToken cancellationToken) => - throw new NotImplementedException(); + public Task Delete(string name, CancellationToken cancellationToken) => ExecuteWithDbContext(async context => + { + context.KnownEndpoints.RemoveRange(context.KnownEndpoints.Where(x => x.Name == name)); + await context.SaveChangesAsync(cancellationToken); + }); } diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs index b53673b630..9a8176054c 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs @@ -5,9 +5,12 @@ namespace ServiceControl.Persistence.EFCore.Implementation; public class QueueAddressStore : IQueueAddressStore { - public Task>> GetAddresses(PagingInfo pagingInfo) => - throw new NotImplementedException(); + public Task>> GetAddresses(PagingInfo pagingInfo) + { + await exe + var queryable + } public Task>> GetAddressesBySearchTerm(string search, PagingInfo pagingInfo) => throw new NotImplementedException(); -} +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/TrialLicenseDataProvider.cs b/src/ServiceControl.Persistence.EFCore/Implementation/TrialLicenseDataProvider.cs index e10739ee89..317e6d4835 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/TrialLicenseDataProvider.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/TrialLicenseDataProvider.cs @@ -2,8 +2,9 @@ namespace ServiceControl.Persistence.EFCore.Implementation; public class TrialLicenseDataProvider : ITrialLicenseDataProvider { - public Task GetTrialEndDate(CancellationToken cancellationToken) => - throw new NotImplementedException(); + public Task GetTrialEndDate(CancellationToken cancellationToken) + { + } public Task StoreTrialEndDate(DateOnly trialEndDate, CancellationToken cancellationToken) => throw new NotImplementedException(); diff --git a/src/ServiceControl.Persistence/IEndpointSettingsStore.cs b/src/ServiceControl.Persistence/IEndpointSettingsStore.cs index d41296dd3a..3e15b838a6 100644 --- a/src/ServiceControl.Persistence/IEndpointSettingsStore.cs +++ b/src/ServiceControl.Persistence/IEndpointSettingsStore.cs @@ -6,7 +6,7 @@ public interface IEndpointSettingsStore { - IAsyncEnumerable GetAllEndpointSettings(CancellationToken token); + IAsyncEnumerable GetAllEndpointSettings(CancellationToken cancellationToken); Task UpdateEndpointSettings(EndpointSettings settings, CancellationToken token); Task Delete(string name, CancellationToken cancellationToken); From 19027530e73c7e4a5dd75931c19f57b31ba84cd1 Mon Sep 17 00:00:00 2001 From: Rhys Bevilaqua Date: Wed, 22 Jul 2026 11:49:54 +0800 Subject: [PATCH 2/4] Implement settings & extend test coverage --- .../DbContexts/ServiceControlDbContext.cs | 9 ++- .../Entities/FailedMessageEntity.cs | 8 ++- .../Entities/TrialMetadataEntity.cs | 9 +++ ...ty.cs => EndpointSettingsConfiguration.cs} | 3 +- .../FailedMessageEntityConfiguration.cs | 20 +++++++ .../TrialMetadataEntityConfiguration.cs | 18 ++++++ .../Implementation/DataStoreBase.cs | 15 +---- .../Implementation/EndpointSettingsStore.cs | 11 ++-- .../Implementation/QueueAddressStore.cs | 30 +++++++--- .../TrialLicenseDataProvider.cs | 24 ++++++-- .../PersistenceTestsContext.cs | 15 +++++ .../PersistenceTestsContext.cs | 33 ++++++++++ .../PersistenceTestsContext.cs | 17 ++++++ .../EndpointSettingsStoreTests.cs | 55 +++++++++++++++++ .../IPersistenceTestsContext.cs | 2 + .../PersistenceTestBase.cs | 2 + .../QueueAddressStoreTests.cs | 60 +++++++++++++++++++ .../IEndpointSettingsStore.cs | 2 +- 18 files changed, 293 insertions(+), 40 deletions(-) create mode 100644 src/ServiceControl.Persistence.EFCore/Entities/TrialMetadataEntity.cs rename src/ServiceControl.Persistence.EFCore/EntityConfigurations/{EndpointSettingsEntity.cs => EndpointSettingsConfiguration.cs} (68%) create mode 100644 src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageEntityConfiguration.cs create mode 100644 src/ServiceControl.Persistence.EFCore/EntityConfigurations/TrialMetadataEntityConfiguration.cs create mode 100644 src/ServiceControl.Persistence.Tests/EndpointSettingsStoreTests.cs create mode 100644 src/ServiceControl.Persistence.Tests/QueueAddressStoreTests.cs diff --git a/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs b/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs index 20bd199a6a..c87c7adcf0 100644 --- a/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs +++ b/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs @@ -7,19 +7,22 @@ namespace ServiceControl.Persistence.EFCore.DbContexts; public abstract class ServiceControlDbContext(DbContextOptions options) : DbContext(options) { public DbSet EndpointSettings { get; set; } + public DbSet FailedMessages { get; set; } public DbSet KnownEndpoints { get; set; } public DbSet KnownEndpointsInsertOnly { get; set; } + public DbSet TrialMetadata { get; set; } protected override void OnConfiguring(DbContextOptionsBuilder optionsBuilder) - { - optionsBuilder.EnableDetailedErrors(); - } + => optionsBuilder.EnableDetailedErrors(); protected override void OnModelCreating(ModelBuilder modelBuilder) { base.OnModelCreating(modelBuilder); + modelBuilder.ApplyConfiguration(new EndpointSettingsConfiguration()); + modelBuilder.ApplyConfiguration(new FailedMessageEntityConfiguration()); modelBuilder.ApplyConfiguration(new KnownEndpointConfiguration()); modelBuilder.ApplyConfiguration(new KnownEndpointInsertOnlyConfiguration()); + modelBuilder.ApplyConfiguration(new TrialMetadataEntityConfiguration()); } } diff --git a/src/ServiceControl.Persistence.EFCore/Entities/FailedMessageEntity.cs b/src/ServiceControl.Persistence.EFCore/Entities/FailedMessageEntity.cs index ae60b4e8f5..2e4eca7cf0 100644 --- a/src/ServiceControl.Persistence.EFCore/Entities/FailedMessageEntity.cs +++ b/src/ServiceControl.Persistence.EFCore/Entities/FailedMessageEntity.cs @@ -2,8 +2,10 @@ namespace ServiceControl.Persistence.EFCore.Entities; public class FailedMessageEntity { + //todo: incomplete, needs the ingestion done... public Guid Id { get; set; } - public string EndpointAddress - - //todo: incomplete, needs the other feature finished... + + public required string OriginalMessageIdentifier { get; set; } + + public required string EndpointAddress { get; set; } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Entities/TrialMetadataEntity.cs b/src/ServiceControl.Persistence.EFCore/Entities/TrialMetadataEntity.cs new file mode 100644 index 0000000000..68299b4585 --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/Entities/TrialMetadataEntity.cs @@ -0,0 +1,9 @@ +namespace ServiceControl.Persistence.EFCore.Entities; + +public class TrialMetadataEntity +{ + public const int TrialMetadataId = 1; + + public int Id { get; set; } + public DateOnly? TrialEndDate { get; set; } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsEntity.cs b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsConfiguration.cs similarity index 68% rename from src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsEntity.cs rename to src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsConfiguration.cs index 5c135c48af..bb81d61c48 100644 --- a/src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsEntity.cs +++ b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsConfiguration.cs @@ -4,11 +4,12 @@ namespace ServiceControl.Persistence.EFCore.EntityConfigurations; using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore.Metadata.Builders; -public class EndpointSettingsConfiguration : IEntityTypeConfiguration +public class EndpointSettingsConfiguration : IEntityTypeConfiguration { public void Configure(EntityTypeBuilder builder) { builder.ToTable("EndpointSettings"); builder.HasKey(x => x.Name); + builder.Property(x => x.Name).HasMaxLength(200).IsRequired(); } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageEntityConfiguration.cs b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageEntityConfiguration.cs new file mode 100644 index 0000000000..8b6c6fc62b --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageEntityConfiguration.cs @@ -0,0 +1,20 @@ +namespace ServiceControl.Persistence.EFCore.EntityConfigurations; + +using Entities; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Metadata.Builders; + +public class FailedMessageEntityConfiguration : IEntityTypeConfiguration +{ + public void Configure(EntityTypeBuilder builder) + { + builder.ToTable("FailedMessages"); + builder.HasKey(e => e.Id); + builder.Property(e => e.Id).ValueGeneratedNever(); + //todo: validate column constraints + builder.Property(e => e.OriginalMessageIdentifier).HasMaxLength(200).IsRequired(); + builder.Property(e => e.EndpointAddress).HasMaxLength(200).IsRequired(); + + builder.HasIndex(e => new { e.OriginalMessageIdentifier, e.EndpointAddress }); + } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/EntityConfigurations/TrialMetadataEntityConfiguration.cs b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/TrialMetadataEntityConfiguration.cs new file mode 100644 index 0000000000..40ca0a0fef --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/TrialMetadataEntityConfiguration.cs @@ -0,0 +1,18 @@ +namespace ServiceControl.Persistence.EFCore.EntityConfigurations; + +using Entities; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Metadata.Builders; + +public class TrialMetadataEntityConfiguration : IEntityTypeConfiguration +{ + public void Configure(EntityTypeBuilder builder) + { + builder.HasKey(e => e.Id); + builder.HasData(new TrialMetadataEntity + { + Id = 1, + TrialEndDate = null + }); + } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/DataStoreBase.cs b/src/ServiceControl.Persistence.EFCore/Implementation/DataStoreBase.cs index cdafb21f2d..90299c070e 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/DataStoreBase.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/DataStoreBase.cs @@ -36,20 +36,7 @@ protected async Task ExecuteWithDbContext(Func op /// /// Executes an operation with a scoped DbContext, without returning a result /// - protected async IAsyncEnumerable ExecuteWithDbContext(Func> operation) - { - await using var scope = scopeFactory.CreateAsyncScope(); - var dbContext = scope.ServiceProvider.GetRequiredService(); - await foreach (var row in operation(dbContext)) - { - yield return row; - } - } - - /// - /// Executes an operation with a scoped DbContext, without returning a result - /// - protected async IAsyncEnumerable ExecuteWithDbContext(Func> operation, [EnumeratorCancellation] CancellationToken cancellationToken) + protected async IAsyncEnumerable ExecuteWithDbContext(Func> operation, [EnumeratorCancellation] CancellationToken cancellationToken = default) { await using var scope = scopeFactory.CreateAsyncScope(); var dbContext = scope.ServiceProvider.GetRequiredService(); diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs index 65842f6828..cb773ea3ca 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs @@ -25,20 +25,21 @@ public Task UpdateEndpointSettings(EndpointSettings settings, CancellationToken await context.SaveChangesAsync(token); return; } - catch (someexception e) when (condition) + catch (DbUpdateException) { - //this failed because of key conflict so it's a race! + //this probably failed because of key conflict so try again } await context.Entry(entity).ReloadAsync(token); - entity.TrackInstances = settings.TrackInstances; - await context.SaveChangesAsync(token); } + + entity.TrackInstances = settings.TrackInstances; + await context.SaveChangesAsync(token); }); public Task Delete(string name, CancellationToken cancellationToken) => ExecuteWithDbContext(async context => { - context.KnownEndpoints.RemoveRange(context.KnownEndpoints.Where(x => x.Name == name)); + context.EndpointSettings.RemoveRange(context.EndpointSettings.Where(x => x.Name == name)); await context.SaveChangesAsync(cancellationToken); }); } diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs index 9a8176054c..9dafbf3509 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs @@ -1,16 +1,32 @@ namespace ServiceControl.Persistence.EFCore.Implementation; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; using ServiceControl.MessageFailures; using ServiceControl.Persistence.Infrastructure; -public class QueueAddressStore : IQueueAddressStore +public class QueueAddressStore(IServiceScopeFactory scopeFactory) : DataStoreBase(scopeFactory), IQueueAddressStore { - public Task>> GetAddresses(PagingInfo pagingInfo) - { - await exe - var queryable - } + public Task>> GetAddresses(PagingInfo pagingInfo) => + GetAddressesBySearchTerm(string.Empty, pagingInfo); public Task>> GetAddressesBySearchTerm(string search, PagingInfo pagingInfo) => - throw new NotImplementedException(); + ExecuteWithDbContext(async context => + { + var query = context.FailedMessages + .Where(fm => string.IsNullOrWhiteSpace(search) || fm.EndpointAddress.StartsWith(search)) + .Select(fm => new { fm.OriginalMessageIdentifier, fm.EndpointAddress }) + .Distinct() + .GroupBy(failure => failure.EndpointAddress) + .OrderBy(failuresByEndpoint => failuresByEndpoint.Key) + .Select(failuresByEndpoint => new QueueAddress + { + PhysicalAddress = failuresByEndpoint.Key, + FailedMessageCount = failuresByEndpoint.Count() + }); + + var items = await query.Skip(pagingInfo.Offset).Take(pagingInfo.PageSize).ToListAsync(); + + return new QueryResult>(items, new QueryStatsInfo("", query.Count(), false)); + }); } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/TrialLicenseDataProvider.cs b/src/ServiceControl.Persistence.EFCore/Implementation/TrialLicenseDataProvider.cs index 317e6d4835..44a317d0bc 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/TrialLicenseDataProvider.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/TrialLicenseDataProvider.cs @@ -1,11 +1,23 @@ namespace ServiceControl.Persistence.EFCore.Implementation; -public class TrialLicenseDataProvider : ITrialLicenseDataProvider +using Entities; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; + +public class TrialLicenseDataProvider(IServiceScopeFactory scopeFactory) : DataStoreBase(scopeFactory), ITrialLicenseDataProvider { public Task GetTrialEndDate(CancellationToken cancellationToken) - { - } + => ExecuteWithDbContext(async context => + { + var trialMetadata = await context.TrialMetadata.SingleAsync(t => t.Id == TrialMetadataEntity.TrialMetadataId, cancellationToken); + return trialMetadata.TrialEndDate; + }); - public Task StoreTrialEndDate(DateOnly trialEndDate, CancellationToken cancellationToken) => - throw new NotImplementedException(); -} + public Task StoreTrialEndDate(DateOnly trialEndDate, CancellationToken cancellationToken) + => ExecuteWithDbContext(async context => + { + var trialMetadata = await context.TrialMetadata.SingleAsync(t => t.Id == TrialMetadataEntity.TrialMetadataId, cancellationToken); + trialMetadata.TrialEndDate = trialEndDate; + await (Task)context.SaveChangesAsync(cancellationToken); + }); +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs index a104f188eb..d29f22f75d 100644 --- a/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs +++ b/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs @@ -69,6 +69,21 @@ public async Task CompleteDatabaseOperation() } } + public async Task SeedFailedMessagesForQueueAddressStore(params (string OriginalMessageIdentifier, string EndpointAddress)[] failedMessages) + { + using var scope = host.Services.CreateScope(); + var db = scope.ServiceProvider.GetRequiredService(); + + db.FailedMessages.AddRange(failedMessages.Select(failedMessage => new FailedMessageEntity + { + Id = Guid.NewGuid(), + OriginalMessageIdentifier = failedMessage.OriginalMessageIdentifier, + EndpointAddress = failedMessage.EndpointAddress + })); + + await db.SaveChangesAsync(); + } + public PersistenceSettings PersistenceSettings { get; set; } public string GenerateFailedMessageRecordId(string messageId) => messageId; diff --git a/src/ServiceControl.Persistence.Tests.RavenDB/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.RavenDB/PersistenceTestsContext.cs index 63e455105c..fc2c3b4bc6 100644 --- a/src/ServiceControl.Persistence.Tests.RavenDB/PersistenceTestsContext.cs +++ b/src/ServiceControl.Persistence.Tests.RavenDB/PersistenceTestsContext.cs @@ -8,6 +8,8 @@ namespace ServiceControl.Persistence.Tests; using Microsoft.Extensions.Hosting; using NUnit.Framework; using Raven.Client.Documents; +using ServiceControl.Contracts.Operations; +using ServiceControl.MessageFailures; using ServiceControl.Persistence; using ServiceControl.Persistence.RavenDB; using ServiceControl.RavenDB; @@ -65,6 +67,37 @@ public Task CompleteDatabaseOperation() return Task.CompletedTask; } + public async Task SeedFailedMessagesForQueueAddressStore(params (string OriginalMessageIdentifier, string EndpointAddress)[] failedMessages) + { + using var session = await SessionProvider.OpenSession(); + + foreach (var (originalMessageIdentifier, endpointAddress) in failedMessages) + { + var documentId = $"FailedMessages/{originalMessageIdentifier}"; + var failedMessage = new FailedMessage + { + Id = documentId, + UniqueMessageId = originalMessageIdentifier, + Status = FailedMessageStatus.Unresolved, + ProcessingAttempts = + [ + new FailedMessage.ProcessingAttempt + { + AttemptedAt = DateTime.UtcNow, + FailureDetails = new FailureDetails + { + AddressOfFailingEndpoint = endpointAddress + } + } + ] + }; + + await session.StoreAsync(failedMessage, documentId); + } + + await session.SaveChangesAsync(); + } + [Conditional("DEBUG")] public void BlockToInspectDatabase() { diff --git a/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs index 9af5dc6de0..ad6b2199c6 100644 --- a/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs +++ b/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs @@ -10,6 +10,8 @@ namespace ServiceControl.Persistence.Tests; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Time.Testing; +using ServiceControl.Persistence.EFCore.DbContexts; +using ServiceControl.Persistence.EFCore.Entities; using ServiceControl.Persistence.EFCore.Infrastructure; public class PersistenceTestsContext : IPersistenceTestsContext @@ -75,6 +77,21 @@ public async Task CompleteDatabaseOperation() } } + public async Task SeedFailedMessagesForQueueAddressStore(params (string OriginalMessageIdentifier, string EndpointAddress)[] failedMessages) + { + using var scope = host.Services.CreateScope(); + var db = scope.ServiceProvider.GetRequiredService(); + + db.FailedMessages.AddRange(failedMessages.Select(failedMessage => new FailedMessageEntity + { + Id = Guid.NewGuid(), + OriginalMessageIdentifier = failedMessage.OriginalMessageIdentifier, + EndpointAddress = failedMessage.EndpointAddress + })); + + await db.SaveChangesAsync(); + } + public PersistenceSettings PersistenceSettings { get; set; } public string GenerateFailedMessageRecordId(string messageId) => messageId; diff --git a/src/ServiceControl.Persistence.Tests/EndpointSettingsStoreTests.cs b/src/ServiceControl.Persistence.Tests/EndpointSettingsStoreTests.cs new file mode 100644 index 0000000000..7705f1b9a8 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EndpointSettingsStoreTests.cs @@ -0,0 +1,55 @@ +namespace ServiceControl.Persistence.Tests; + +using System.Collections.Generic; +using System.Linq; +using System.Threading.Tasks; +using NUnit.Framework; +using ServiceControl.Persistence; + +class EndpointSettingsStoreTests : PersistenceTestBase +{ + [Test] + public async Task UpdateEndpointSettings_stores_and_updates_existing_setting() + { + await EndpointSettingsStore.UpdateEndpointSettings(new EndpointSettings { Name = "Sales", TrackInstances = false }, default); + await EndpointSettingsStore.UpdateEndpointSettings(new EndpointSettings { Name = "Sales", TrackInstances = true }, default); + + var settings = await GetAllEndpointSettings(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(settings, Has.Count.EqualTo(1)); + Assert.That(settings.Single().Name, Is.EqualTo("Sales")); + Assert.That(settings.Single().TrackInstances, Is.True); + } + } + + [Test] + public async Task Delete_removes_only_target_setting() + { + await EndpointSettingsStore.UpdateEndpointSettings(new EndpointSettings { Name = "Sales", TrackInstances = false }, default); + await EndpointSettingsStore.UpdateEndpointSettings(new EndpointSettings { Name = "Shipping", TrackInstances = true }, default); + + await EndpointSettingsStore.Delete("Sales", default); + + var settings = await GetAllEndpointSettings(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(settings, Has.Count.EqualTo(1)); + Assert.That(settings.Single().Name, Is.EqualTo("Shipping")); + Assert.That(settings.Single().TrackInstances, Is.True); + } + } + + async Task> GetAllEndpointSettings() + { + var settings = new List(); + await foreach (var setting in EndpointSettingsStore.GetAllEndpointSettings()) + { + settings.Add(setting); + } + + return settings; + } +} diff --git a/src/ServiceControl.Persistence.Tests/IPersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests/IPersistenceTestsContext.cs index bf59a57835..b52c419ec5 100644 --- a/src/ServiceControl.Persistence.Tests/IPersistenceTestsContext.cs +++ b/src/ServiceControl.Persistence.Tests/IPersistenceTestsContext.cs @@ -13,6 +13,8 @@ public interface IPersistenceTestsContext Task CompleteDatabaseOperation(); + Task SeedFailedMessagesForQueueAddressStore(params (string OriginalMessageIdentifier, string EndpointAddress)[] failedMessages); + PersistenceSettings PersistenceSettings { get; } string GenerateFailedMessageRecordId(string messageId); diff --git a/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs b/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs index e4175e2c4a..dff4fbd232 100644 --- a/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs +++ b/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs @@ -107,4 +107,6 @@ protected static async Task WaitUntil(Func> conditionChecker, string protected IEventLogDataStore EventLogDataStore => ServiceProvider.GetRequiredService(); protected IRetryDocumentDataStore RetryDocumentDataStore => ServiceProvider.GetRequiredService(); protected ILicensingDataStore LicensingDataStore => ServiceProvider.GetRequiredService(); + protected IQueueAddressStore QueueAddressStore => ServiceProvider.GetRequiredService(); + protected IEndpointSettingsStore EndpointSettingsStore => ServiceProvider.GetRequiredService(); } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.Tests/QueueAddressStoreTests.cs b/src/ServiceControl.Persistence.Tests/QueueAddressStoreTests.cs new file mode 100644 index 0000000000..7031626717 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/QueueAddressStoreTests.cs @@ -0,0 +1,60 @@ +namespace ServiceControl.Persistence.Tests; + +using System.Linq; +using System.Threading.Tasks; +using NUnit.Framework; +using ServiceControl.Persistence.Infrastructure; + +class QueueAddressStoreTests : PersistenceTestBase +{ + [Test] + public async Task GetAddresses_groups_failed_messages_by_queue_address() + { + await SeedFailedMessages(); + + await CompleteDatabaseOperation(); + + var result = await QueueAddressStore.GetAddresses(new PagingInfo(1, 10)); + var addresses = result.Results; + var physicalAddresses = addresses.Select(address => address.PhysicalAddress).OrderBy(address => address).ToArray(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(result.QueryStats.TotalCount, Is.EqualTo(4)); + Assert.That(physicalAddresses, Is.EqualTo(new[] { "alpha", "alpha-child", "beta", "gamma" })); + Assert.That(addresses.Single(address => address.PhysicalAddress == "alpha").FailedMessageCount, Is.EqualTo(2)); + Assert.That(addresses.Single(address => address.PhysicalAddress == "alpha-child").FailedMessageCount, Is.EqualTo(1)); + Assert.That(addresses.Single(address => address.PhysicalAddress == "beta").FailedMessageCount, Is.EqualTo(1)); + Assert.That(addresses.Single(address => address.PhysicalAddress == "gamma").FailedMessageCount, Is.EqualTo(1)); + } + } + + [Test] + public async Task GetAddressesBySearchTerm_filters_by_prefix() + { + await SeedFailedMessages(); + + await CompleteDatabaseOperation(); + + var result = await QueueAddressStore.GetAddressesBySearchTerm("alpha", new PagingInfo(1, 10)); + var addresses = result.Results; + var physicalAddresses = addresses.Select(address => address.PhysicalAddress).OrderBy(address => address).ToArray(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(result.QueryStats.TotalCount, Is.EqualTo(2)); + Assert.That(physicalAddresses, Is.EqualTo(new[] { "alpha", "alpha-child" })); + Assert.That(addresses.Sum(address => address.FailedMessageCount), Is.EqualTo(3)); + } + } + + async Task SeedFailedMessages() + { + await PersistenceTestsContext.SeedFailedMessagesForQueueAddressStore( + ("msg-1", "alpha"), + ("msg-2", "alpha"), + ("msg-3", "beta"), + ("msg-4", "alpha-child"), + ("msg-5", "gamma")); + } +} diff --git a/src/ServiceControl.Persistence/IEndpointSettingsStore.cs b/src/ServiceControl.Persistence/IEndpointSettingsStore.cs index 3e15b838a6..2b6c52e8a1 100644 --- a/src/ServiceControl.Persistence/IEndpointSettingsStore.cs +++ b/src/ServiceControl.Persistence/IEndpointSettingsStore.cs @@ -6,7 +6,7 @@ public interface IEndpointSettingsStore { - IAsyncEnumerable GetAllEndpointSettings(CancellationToken cancellationToken); + IAsyncEnumerable GetAllEndpointSettings(CancellationToken cancellationToken = default); Task UpdateEndpointSettings(EndpointSettings settings, CancellationToken token); Task Delete(string name, CancellationToken cancellationToken); From 27b9641dc2de4bbdc39c7dade066bc4798e5de2f Mon Sep 17 00:00:00 2001 From: Rhys Bevilaqua Date: Wed, 22 Jul 2026 13:37:34 +0800 Subject: [PATCH 3/4] Recreate the initial migration --- .../20260720230745_Initial.Designer.cs | 93 ---------- .../Migrations/20260720230745_Initial.cs | 62 ------- .../20260722052903_Initial.Designer.cs | 165 ++++++++++++++++++ .../Migrations/20260722052903_Initial.cs | 119 +++++++++++++ ...SqlServiceControlDbContextModelSnapshot.cs | 82 ++++++++- ....cs => 20260722052648_Initial.Designer.cs} | 60 ++++++- ...1_Initial.cs => 20260722052648_Initial.cs} | 57 ++++++ ...verServiceControlDbContextModelSnapshot.cs | 58 ++++++ .../FailedMessageEntityConfiguration.cs | 2 +- .../PersistenceTestsContext.cs | 1 + .../TrialLicenseDataProviderTests.cs | 33 ++++ 11 files changed, 570 insertions(+), 162 deletions(-) delete mode 100644 src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260720230745_Initial.Designer.cs delete mode 100644 src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260720230745_Initial.cs create mode 100644 src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260722052903_Initial.Designer.cs create mode 100644 src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260722052903_Initial.cs rename src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/{20260720230811_Initial.Designer.cs => 20260722052648_Initial.Designer.cs} (57%) rename src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/{20260720230811_Initial.cs => 20260722052648_Initial.cs} (50%) create mode 100644 src/ServiceControl.Persistence.Tests/TrialLicenseDataProviderTests.cs diff --git a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260720230745_Initial.Designer.cs b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260720230745_Initial.Designer.cs deleted file mode 100644 index 1a1383b3e2..0000000000 --- a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260720230745_Initial.Designer.cs +++ /dev/null @@ -1,93 +0,0 @@ -// -using System; -using Microsoft.EntityFrameworkCore; -using Microsoft.EntityFrameworkCore.Infrastructure; -using Microsoft.EntityFrameworkCore.Migrations; -using Microsoft.EntityFrameworkCore.Storage.ValueConversion; -using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; -using ServiceControl.Persistence.EFCore.PostgreSql; - -#nullable disable - -namespace ServiceControl.Persistence.EFCore.PostgreSql.Migrations -{ - [DbContext(typeof(PostgreSqlServiceControlDbContext))] - [Migration("20260720230745_Initial")] - partial class Initial - { - /// - protected override void BuildTargetModel(ModelBuilder modelBuilder) - { -#pragma warning disable 612, 618 - modelBuilder - .HasAnnotation("ProductVersion", "10.0.9") - .HasAnnotation("Relational:MaxIdentifierLength", 63); - - NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); - - modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.KnownEndpointEntity", b => - { - b.Property("Id") - .HasColumnType("uuid") - .HasColumnName("id"); - - b.Property("Host") - .IsRequired() - .HasColumnType("text") - .HasColumnName("host"); - - b.Property("HostId") - .HasColumnType("uuid") - .HasColumnName("host_id"); - - b.Property("Monitored") - .HasColumnType("boolean") - .HasColumnName("monitored"); - - b.Property("Name") - .IsRequired() - .HasColumnType("text") - .HasColumnName("name"); - - b.HasKey("Id"); - - b.ToTable("known_endpoints", (string)null); - }); - - modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.KnownEndpointInsertOnlyEntity", b => - { - b.Property("Id") - .ValueGeneratedOnAdd() - .HasColumnType("bigint") - .HasColumnName("id"); - - NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id")); - - b.Property("Host") - .IsRequired() - .HasColumnType("text") - .HasColumnName("host"); - - b.Property("HostId") - .HasColumnType("uuid") - .HasColumnName("host_id"); - - b.Property("KnownEndpointId") - .HasColumnType("uuid") - .HasColumnName("known_endpoint_id"); - - b.Property("Name") - .IsRequired() - .HasColumnType("text") - .HasColumnName("name"); - - b.HasKey("Id"); - - b.HasIndex("KnownEndpointId"); - - b.ToTable("known_endpoints_insert_only", (string)null); - }); -#pragma warning restore 612, 618 - } - } -} diff --git a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260720230745_Initial.cs b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260720230745_Initial.cs deleted file mode 100644 index d2caaa189a..0000000000 --- a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260720230745_Initial.cs +++ /dev/null @@ -1,62 +0,0 @@ -using System; -using Microsoft.EntityFrameworkCore.Migrations; -using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; - -#nullable disable - -namespace ServiceControl.Persistence.EFCore.PostgreSql.Migrations -{ - /// - public partial class Initial : Migration - { - /// - protected override void Up(MigrationBuilder migrationBuilder) - { - migrationBuilder.CreateTable( - name: "known_endpoints", - columns: table => new - { - id = table.Column(type: "uuid", nullable: false), - name = table.Column(type: "text", nullable: false), - host_id = table.Column(type: "uuid", nullable: false), - host = table.Column(type: "text", nullable: false), - monitored = table.Column(type: "boolean", nullable: false) - }, - constraints: table => - { - table.PrimaryKey("PK_known_endpoints", x => x.id); - }); - - migrationBuilder.CreateTable( - name: "known_endpoints_insert_only", - columns: table => new - { - id = table.Column(type: "bigint", nullable: false) - .Annotation("Npgsql:ValueGenerationStrategy", NpgsqlValueGenerationStrategy.IdentityByDefaultColumn), - known_endpoint_id = table.Column(type: "uuid", nullable: false), - name = table.Column(type: "text", nullable: false), - host_id = table.Column(type: "uuid", nullable: false), - host = table.Column(type: "text", nullable: false) - }, - constraints: table => - { - table.PrimaryKey("PK_known_endpoints_insert_only", x => x.id); - }); - - migrationBuilder.CreateIndex( - name: "IX_known_endpoints_insert_only_known_endpoint_id", - table: "known_endpoints_insert_only", - column: "known_endpoint_id"); - } - - /// - protected override void Down(MigrationBuilder migrationBuilder) - { - migrationBuilder.DropTable( - name: "known_endpoints"); - - migrationBuilder.DropTable( - name: "known_endpoints_insert_only"); - } - } -} diff --git a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260722052903_Initial.Designer.cs b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260722052903_Initial.Designer.cs new file mode 100644 index 0000000000..53725d994f --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260722052903_Initial.Designer.cs @@ -0,0 +1,165 @@ +// +using System; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; +using ServiceControl.Persistence.EFCore.PostgreSql; + +#nullable disable + +namespace ServiceControl.Persistence.EFCore.PostgreSql.Migrations +{ + [DbContext(typeof(PostgreSqlServiceControlDbContext))] + [Migration("20260722052903_Initial")] + partial class Initial + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("ProductVersion", "10.0.9") + .HasAnnotation("Relational:MaxIdentifierLength", 63); + + NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.EndpointSettingsEntity", b => + { + b.Property("Name") + .HasMaxLength(200) + .HasColumnType("character varying(200)") + .HasColumnName("name"); + + b.Property("TrackInstances") + .HasColumnType("boolean") + .HasColumnName("track_instances"); + + b.HasKey("Name") + .HasName("pk_endpoint_settings"); + + b.ToTable("EndpointSettings", (string)null); + }); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageEntity", b => + { + b.Property("Id") + .HasColumnType("uuid") + .HasColumnName("id"); + + b.Property("EndpointAddress") + .IsRequired() + .HasMaxLength(200) + .HasColumnType("character varying(200)") + .HasColumnName("endpoint_address"); + + b.Property("OriginalMessageIdentifier") + .IsRequired() + .HasMaxLength(200) + .HasColumnType("character varying(200)") + .HasColumnName("original_message_identifier"); + + b.HasKey("Id") + .HasName("pk_failed_messages"); + + b.HasIndex("OriginalMessageIdentifier", "EndpointAddress") + .HasDatabaseName("ix_failed_messages_original_message_identifier_endpoint_address"); + + b.ToTable("FailedMessages", (string)null); + }); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.KnownEndpointEntity", b => + { + b.Property("Id") + .HasColumnType("uuid") + .HasColumnName("id"); + + b.Property("Host") + .IsRequired() + .HasColumnType("text") + .HasColumnName("host"); + + b.Property("HostId") + .HasColumnType("uuid") + .HasColumnName("host_id"); + + b.Property("Monitored") + .HasColumnType("boolean") + .HasColumnName("monitored"); + + b.Property("Name") + .IsRequired() + .HasColumnType("text") + .HasColumnName("name"); + + b.HasKey("Id") + .HasName("pk_known_endpoints"); + + b.ToTable("KnownEndpoints", (string)null); + }); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.KnownEndpointInsertOnlyEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("bigint") + .HasColumnName("id"); + + NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id")); + + b.Property("Host") + .IsRequired() + .HasColumnType("text") + .HasColumnName("host"); + + b.Property("HostId") + .HasColumnType("uuid") + .HasColumnName("host_id"); + + b.Property("KnownEndpointId") + .HasColumnType("uuid") + .HasColumnName("known_endpoint_id"); + + b.Property("Name") + .IsRequired() + .HasColumnType("text") + .HasColumnName("name"); + + b.HasKey("Id") + .HasName("pk_known_endpoints_insert_only"); + + b.HasIndex("KnownEndpointId") + .HasDatabaseName("ix_known_endpoints_insert_only_known_endpoint_id"); + + b.ToTable("KnownEndpointsInsertOnly", (string)null); + }); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.TrialMetadataEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("integer") + .HasColumnName("id"); + + NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id")); + + b.Property("TrialEndDate") + .HasColumnType("date") + .HasColumnName("trial_end_date"); + + b.HasKey("Id") + .HasName("pk_trial_metadata"); + + b.ToTable("trial_metadata", (string)null); + + b.HasData( + new + { + Id = 1 + }); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260722052903_Initial.cs b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260722052903_Initial.cs new file mode 100644 index 0000000000..0059248dea --- /dev/null +++ b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260722052903_Initial.cs @@ -0,0 +1,119 @@ +using System; +using Microsoft.EntityFrameworkCore.Migrations; +using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; + +#nullable disable + +namespace ServiceControl.Persistence.EFCore.PostgreSql.Migrations +{ + /// + public partial class Initial : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.CreateTable( + name: "EndpointSettings", + columns: table => new + { + name = table.Column(type: "character varying(200)", maxLength: 200, nullable: false), + track_instances = table.Column(type: "boolean", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("pk_endpoint_settings", x => x.name); + }); + + migrationBuilder.CreateTable( + name: "FailedMessages", + columns: table => new + { + id = table.Column(type: "uuid", nullable: false), + original_message_identifier = table.Column(type: "character varying(200)", maxLength: 200, nullable: false), + endpoint_address = table.Column(type: "character varying(200)", maxLength: 200, nullable: false) + }, + constraints: table => + { + table.PrimaryKey("pk_failed_messages", x => x.id); + }); + + migrationBuilder.CreateTable( + name: "KnownEndpoints", + columns: table => new + { + id = table.Column(type: "uuid", nullable: false), + name = table.Column(type: "text", nullable: false), + host_id = table.Column(type: "uuid", nullable: false), + host = table.Column(type: "text", nullable: false), + monitored = table.Column(type: "boolean", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("pk_known_endpoints", x => x.id); + }); + + migrationBuilder.CreateTable( + name: "KnownEndpointsInsertOnly", + columns: table => new + { + id = table.Column(type: "bigint", nullable: false) + .Annotation("Npgsql:ValueGenerationStrategy", NpgsqlValueGenerationStrategy.IdentityByDefaultColumn), + known_endpoint_id = table.Column(type: "uuid", nullable: false), + name = table.Column(type: "text", nullable: false), + host_id = table.Column(type: "uuid", nullable: false), + host = table.Column(type: "text", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("pk_known_endpoints_insert_only", x => x.id); + }); + + migrationBuilder.CreateTable( + name: "trial_metadata", + columns: table => new + { + id = table.Column(type: "integer", nullable: false) + .Annotation("Npgsql:ValueGenerationStrategy", NpgsqlValueGenerationStrategy.IdentityByDefaultColumn), + trial_end_date = table.Column(type: "date", nullable: true) + }, + constraints: table => + { + table.PrimaryKey("pk_trial_metadata", x => x.id); + }); + + migrationBuilder.InsertData( + table: "trial_metadata", + columns: new[] { "id", "trial_end_date" }, + values: new object[] { 1, null }); + + migrationBuilder.CreateIndex( + name: "ix_failed_messages_original_message_identifier_endpoint_address", + table: "FailedMessages", + columns: new[] { "original_message_identifier", "endpoint_address" }); + + migrationBuilder.CreateIndex( + name: "ix_known_endpoints_insert_only_known_endpoint_id", + table: "KnownEndpointsInsertOnly", + column: "known_endpoint_id"); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropTable( + name: "EndpointSettings"); + + migrationBuilder.DropTable( + name: "FailedMessages"); + + migrationBuilder.DropTable( + name: "KnownEndpoints"); + + migrationBuilder.DropTable( + name: "KnownEndpointsInsertOnly"); + + migrationBuilder.DropTable( + name: "trial_metadata"); + } + } +} diff --git a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/PostgreSqlServiceControlDbContextModelSnapshot.cs b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/PostgreSqlServiceControlDbContextModelSnapshot.cs index 83b17187a4..5d94f7947b 100644 --- a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/PostgreSqlServiceControlDbContextModelSnapshot.cs +++ b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/PostgreSqlServiceControlDbContextModelSnapshot.cs @@ -22,6 +22,50 @@ protected override void BuildModel(ModelBuilder modelBuilder) NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.EndpointSettingsEntity", b => + { + b.Property("Name") + .HasMaxLength(200) + .HasColumnType("character varying(200)") + .HasColumnName("name"); + + b.Property("TrackInstances") + .HasColumnType("boolean") + .HasColumnName("track_instances"); + + b.HasKey("Name") + .HasName("pk_endpoint_settings"); + + b.ToTable("EndpointSettings", (string)null); + }); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageEntity", b => + { + b.Property("Id") + .HasColumnType("uuid") + .HasColumnName("id"); + + b.Property("EndpointAddress") + .IsRequired() + .HasMaxLength(200) + .HasColumnType("character varying(200)") + .HasColumnName("endpoint_address"); + + b.Property("OriginalMessageIdentifier") + .IsRequired() + .HasMaxLength(200) + .HasColumnType("character varying(200)") + .HasColumnName("original_message_identifier"); + + b.HasKey("Id") + .HasName("pk_failed_messages"); + + b.HasIndex("OriginalMessageIdentifier", "EndpointAddress") + .HasDatabaseName("ix_failed_messages_original_message_identifier_endpoint_address"); + + b.ToTable("FailedMessages", (string)null); + }); + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.KnownEndpointEntity", b => { b.Property("Id") @@ -46,9 +90,10 @@ protected override void BuildModel(ModelBuilder modelBuilder) .HasColumnType("text") .HasColumnName("name"); - b.HasKey("Id"); + b.HasKey("Id") + .HasName("pk_known_endpoints"); - b.ToTable("known_endpoints", (string)null); + b.ToTable("KnownEndpoints", (string)null); }); modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.KnownEndpointInsertOnlyEntity", b => @@ -78,11 +123,38 @@ protected override void BuildModel(ModelBuilder modelBuilder) .HasColumnType("text") .HasColumnName("name"); - b.HasKey("Id"); + b.HasKey("Id") + .HasName("pk_known_endpoints_insert_only"); + + b.HasIndex("KnownEndpointId") + .HasDatabaseName("ix_known_endpoints_insert_only_known_endpoint_id"); + + b.ToTable("KnownEndpointsInsertOnly", (string)null); + }); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.TrialMetadataEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("integer") + .HasColumnName("id"); + + NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id")); + + b.Property("TrialEndDate") + .HasColumnType("date") + .HasColumnName("trial_end_date"); + + b.HasKey("Id") + .HasName("pk_trial_metadata"); - b.HasIndex("KnownEndpointId"); + b.ToTable("trial_metadata", (string)null); - b.ToTable("known_endpoints_insert_only", (string)null); + b.HasData( + new + { + Id = 1 + }); }); #pragma warning restore 612, 618 } diff --git a/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260720230811_Initial.Designer.cs b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260722052648_Initial.Designer.cs similarity index 57% rename from src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260720230811_Initial.Designer.cs rename to src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260722052648_Initial.Designer.cs index 5901a2770e..d0129a6611 100644 --- a/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260720230811_Initial.Designer.cs +++ b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260722052648_Initial.Designer.cs @@ -12,7 +12,7 @@ namespace ServiceControl.Persistence.EFCore.SqlServer.Migrations { [DbContext(typeof(SqlServerServiceControlDbContext))] - [Migration("20260720230811_Initial")] + [Migration("20260722052648_Initial")] partial class Initial { /// @@ -25,6 +25,42 @@ protected override void BuildTargetModel(ModelBuilder modelBuilder) SqlServerModelBuilderExtensions.UseIdentityColumns(modelBuilder); + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.EndpointSettingsEntity", b => + { + b.Property("Name") + .HasMaxLength(200) + .HasColumnType("nvarchar(200)"); + + b.Property("TrackInstances") + .HasColumnType("bit"); + + b.HasKey("Name"); + + b.ToTable("EndpointSettings", (string)null); + }); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageEntity", b => + { + b.Property("Id") + .HasColumnType("uniqueidentifier"); + + b.Property("EndpointAddress") + .IsRequired() + .HasMaxLength(200) + .HasColumnType("nvarchar(200)"); + + b.Property("OriginalMessageIdentifier") + .IsRequired() + .HasMaxLength(200) + .HasColumnType("nvarchar(200)"); + + b.HasKey("Id"); + + b.HasIndex("OriginalMessageIdentifier", "EndpointAddress"); + + b.ToTable("FailedMessages", (string)null); + }); + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.KnownEndpointEntity", b => { b.Property("Id") @@ -77,6 +113,28 @@ protected override void BuildTargetModel(ModelBuilder modelBuilder) b.ToTable("KnownEndpointsInsertOnly", (string)null); }); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.TrialMetadataEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("int"); + + SqlServerPropertyBuilderExtensions.UseIdentityColumn(b.Property("Id")); + + b.Property("TrialEndDate") + .HasColumnType("date"); + + b.HasKey("Id"); + + b.ToTable("TrialMetadata"); + + b.HasData( + new + { + Id = 1 + }); + }); #pragma warning restore 612, 618 } } diff --git a/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260720230811_Initial.cs b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260722052648_Initial.cs similarity index 50% rename from src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260720230811_Initial.cs rename to src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260722052648_Initial.cs index d1af2ee2e1..b95bbd57e7 100644 --- a/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260720230811_Initial.cs +++ b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260722052648_Initial.cs @@ -11,6 +11,31 @@ public partial class Initial : Migration /// protected override void Up(MigrationBuilder migrationBuilder) { + migrationBuilder.CreateTable( + name: "EndpointSettings", + columns: table => new + { + Name = table.Column(type: "nvarchar(200)", maxLength: 200, nullable: false), + TrackInstances = table.Column(type: "bit", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("PK_EndpointSettings", x => x.Name); + }); + + migrationBuilder.CreateTable( + name: "FailedMessages", + columns: table => new + { + Id = table.Column(type: "uniqueidentifier", nullable: false), + OriginalMessageIdentifier = table.Column(type: "nvarchar(200)", maxLength: 200, nullable: false), + EndpointAddress = table.Column(type: "nvarchar(200)", maxLength: 200, nullable: false) + }, + constraints: table => + { + table.PrimaryKey("PK_FailedMessages", x => x.Id); + }); + migrationBuilder.CreateTable( name: "KnownEndpoints", columns: table => new @@ -42,6 +67,29 @@ protected override void Up(MigrationBuilder migrationBuilder) table.PrimaryKey("PK_KnownEndpointsInsertOnly", x => x.Id); }); + migrationBuilder.CreateTable( + name: "TrialMetadata", + columns: table => new + { + Id = table.Column(type: "int", nullable: false) + .Annotation("SqlServer:Identity", "1, 1"), + TrialEndDate = table.Column(type: "date", nullable: true) + }, + constraints: table => + { + table.PrimaryKey("PK_TrialMetadata", x => x.Id); + }); + + migrationBuilder.InsertData( + table: "TrialMetadata", + columns: new[] { "Id", "TrialEndDate" }, + values: new object[] { 1, null }); + + migrationBuilder.CreateIndex( + name: "IX_FailedMessages_OriginalMessageIdentifier_EndpointAddress", + table: "FailedMessages", + columns: new[] { "OriginalMessageIdentifier", "EndpointAddress" }); + migrationBuilder.CreateIndex( name: "IX_KnownEndpointsInsertOnly_KnownEndpointId", table: "KnownEndpointsInsertOnly", @@ -51,11 +99,20 @@ protected override void Up(MigrationBuilder migrationBuilder) /// protected override void Down(MigrationBuilder migrationBuilder) { + migrationBuilder.DropTable( + name: "EndpointSettings"); + + migrationBuilder.DropTable( + name: "FailedMessages"); + migrationBuilder.DropTable( name: "KnownEndpoints"); migrationBuilder.DropTable( name: "KnownEndpointsInsertOnly"); + + migrationBuilder.DropTable( + name: "TrialMetadata"); } } } diff --git a/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/SqlServerServiceControlDbContextModelSnapshot.cs b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/SqlServerServiceControlDbContextModelSnapshot.cs index df880e2f92..dd4216b82b 100644 --- a/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/SqlServerServiceControlDbContextModelSnapshot.cs +++ b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/SqlServerServiceControlDbContextModelSnapshot.cs @@ -22,6 +22,42 @@ protected override void BuildModel(ModelBuilder modelBuilder) SqlServerModelBuilderExtensions.UseIdentityColumns(modelBuilder); + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.EndpointSettingsEntity", b => + { + b.Property("Name") + .HasMaxLength(200) + .HasColumnType("nvarchar(200)"); + + b.Property("TrackInstances") + .HasColumnType("bit"); + + b.HasKey("Name"); + + b.ToTable("EndpointSettings", (string)null); + }); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageEntity", b => + { + b.Property("Id") + .HasColumnType("uniqueidentifier"); + + b.Property("EndpointAddress") + .IsRequired() + .HasMaxLength(200) + .HasColumnType("nvarchar(200)"); + + b.Property("OriginalMessageIdentifier") + .IsRequired() + .HasMaxLength(200) + .HasColumnType("nvarchar(200)"); + + b.HasKey("Id"); + + b.HasIndex("OriginalMessageIdentifier", "EndpointAddress"); + + b.ToTable("FailedMessages", (string)null); + }); + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.KnownEndpointEntity", b => { b.Property("Id") @@ -74,6 +110,28 @@ protected override void BuildModel(ModelBuilder modelBuilder) b.ToTable("KnownEndpointsInsertOnly", (string)null); }); + + modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.TrialMetadataEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("int"); + + SqlServerPropertyBuilderExtensions.UseIdentityColumn(b.Property("Id")); + + b.Property("TrialEndDate") + .HasColumnType("date"); + + b.HasKey("Id"); + + b.ToTable("TrialMetadata"); + + b.HasData( + new + { + Id = 1 + }); + }); #pragma warning restore 612, 618 } } diff --git a/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageEntityConfiguration.cs b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageEntityConfiguration.cs index 8b6c6fc62b..225464d93c 100644 --- a/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageEntityConfiguration.cs +++ b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageEntityConfiguration.cs @@ -14,7 +14,7 @@ public void Configure(EntityTypeBuilder builder) //todo: validate column constraints builder.Property(e => e.OriginalMessageIdentifier).HasMaxLength(200).IsRequired(); builder.Property(e => e.EndpointAddress).HasMaxLength(200).IsRequired(); - + builder.HasIndex(e => new { e.OriginalMessageIdentifier, e.EndpointAddress }); } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs index d29f22f75d..cc05843a14 100644 --- a/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs +++ b/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs @@ -10,6 +10,7 @@ namespace ServiceControl.Persistence.Tests; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Time.Testing; using Npgsql; +using ServiceControl.Persistence.EFCore.Entities; using ServiceControl.Persistence.EFCore.Infrastructure; public class PersistenceTestsContext : IPersistenceTestsContext diff --git a/src/ServiceControl.Persistence.Tests/TrialLicenseDataProviderTests.cs b/src/ServiceControl.Persistence.Tests/TrialLicenseDataProviderTests.cs new file mode 100644 index 0000000000..627ea19be2 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/TrialLicenseDataProviderTests.cs @@ -0,0 +1,33 @@ +namespace ServiceControl.Persistence.Tests; + +using System; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using NUnit.Framework; +using ServiceControl.Persistence; + +class TrialLicenseDataProviderTests : PersistenceTestBase +{ + [Test] + public async Task GetTrialEndDate_returns_null_by_default() + { + var trialLicenseDataProvider = ServiceProvider.GetRequiredService(); + + var trialEndDate = await trialLicenseDataProvider.GetTrialEndDate(default); + + Assert.That(trialEndDate, Is.Null); + } + + [Test] + public async Task StoreTrialEndDate_persists_value() + { + var trialLicenseDataProvider = ServiceProvider.GetRequiredService(); + var expectedEndDate = DateOnly.FromDateTime(DateTime.UtcNow.Date.AddDays(13)); + + await trialLicenseDataProvider.StoreTrialEndDate(expectedEndDate, default); + + var trialEndDate = await trialLicenseDataProvider.GetTrialEndDate(default); + + Assert.That(trialEndDate, Is.EqualTo(expectedEndDate)); + } +} From 5db80bacc90b8a31b0e0cff9a46b01c67769b1da Mon Sep 17 00:00:00 2001 From: Rhys Bevilaqua Date: Wed, 22 Jul 2026 16:08:12 +0800 Subject: [PATCH 4/4] revert change outside of EF persisters --- src/ServiceControl.Audit/Recoverability/QueueAddress.cs | 8 ++++++++ 1 file changed, 8 insertions(+) create mode 100644 src/ServiceControl.Audit/Recoverability/QueueAddress.cs diff --git a/src/ServiceControl.Audit/Recoverability/QueueAddress.cs b/src/ServiceControl.Audit/Recoverability/QueueAddress.cs new file mode 100644 index 0000000000..e0b89747b3 --- /dev/null +++ b/src/ServiceControl.Audit/Recoverability/QueueAddress.cs @@ -0,0 +1,8 @@ +namespace ServiceControl.Audit.Recoverability +{ + public class QueueAddress + { + public string PhysicalAddress { get; set; } + public int FailedMessageCount { get; set; } + } +} \ No newline at end of file