diff --git a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260724055936_AddConfiguration.Designer.cs b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260724055936_AddConfiguration.Designer.cs
new file mode 100644
index 0000000000..1ac817af7f
--- /dev/null
+++ b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260724055936_AddConfiguration.Designer.cs
@@ -0,0 +1,315 @@
+//
+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("20260724055936_AddConfiguration")]
+ partial class AddConfiguration
+ {
+ ///
+ protected override void BuildTargetModel(ModelBuilder modelBuilder)
+ {
+#pragma warning disable 612, 618
+ modelBuilder
+ .HasAnnotation("ProductVersion", "10.0.10")
+ .HasAnnotation("Relational:MaxIdentifierLength", 63);
+
+ NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.EndpointSettingsEntity", b =>
+ {
+ b.Property("Name")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .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("UniqueMessageId")
+ .HasColumnType("uuid")
+ .HasColumnName("unique_message_id");
+
+ b.Property("BodyContentType")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("body_content_type");
+
+ b.Property("BodySize")
+ .HasColumnType("integer")
+ .HasColumnName("body_size");
+
+ b.Property("BodyStoredExternally")
+ .HasColumnType("boolean")
+ .HasColumnName("body_stored_externally");
+
+ b.Property("BodyText")
+ .HasColumnType("text")
+ .HasColumnName("body_text");
+
+ b.Property("ConversationId")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("conversation_id");
+
+ b.Property("ExceptionMessage")
+ .HasColumnType("text")
+ .HasColumnName("exception_message");
+
+ b.Property("ExceptionType")
+ .HasColumnType("text")
+ .HasColumnName("exception_type");
+
+ b.Property("FailingEndpointAddress")
+ .IsRequired()
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("failing_endpoint_address");
+
+ b.Property("FirstTimeOfFailure")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("first_time_of_failure");
+
+ b.Property("HeadersJson")
+ .IsRequired()
+ .HasColumnType("text")
+ .HasColumnName("headers_json");
+
+ b.Property("IsSystemMessage")
+ .HasColumnType("boolean")
+ .HasColumnName("is_system_message");
+
+ b.Property("LastAttemptedAt")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("last_attempted_at");
+
+ b.Property("LastModified")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("last_modified");
+
+ b.Property("LastTimeOfFailure")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("last_time_of_failure");
+
+ b.Property("MessageId")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("message_id");
+
+ b.Property("MessageType")
+ .HasColumnType("text")
+ .HasColumnName("message_type");
+
+ b.Property("NumberOfProcessingAttempts")
+ .HasColumnType("integer")
+ .HasColumnName("number_of_processing_attempts");
+
+ b.Property("QueueAddress")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("queue_address");
+
+ b.Property("ReceivingEndpointHost")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("receiving_endpoint_host");
+
+ b.Property("ReceivingEndpointHostId")
+ .HasColumnType("uuid")
+ .HasColumnName("receiving_endpoint_host_id");
+
+ b.Property("ReceivingEndpointName")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("receiving_endpoint_name");
+
+ b.Property("SendingEndpointHost")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("sending_endpoint_host");
+
+ b.Property("SendingEndpointHostId")
+ .HasColumnType("uuid")
+ .HasColumnName("sending_endpoint_host_id");
+
+ b.Property("SendingEndpointName")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("sending_endpoint_name");
+
+ b.Property("Status")
+ .HasColumnType("integer")
+ .HasColumnName("status");
+
+ b.Property("StatusChangedAt")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("status_changed_at");
+
+ b.Property("TimeSent")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("time_sent");
+
+ b.HasKey("UniqueMessageId")
+ .HasName("pk_failed_messages");
+
+ b.HasIndex("ConversationId")
+ .HasDatabaseName("ix_failed_messages_conversation_id");
+
+ b.HasIndex("FailingEndpointAddress")
+ .HasDatabaseName("ix_failed_messages_failing_endpoint_address");
+
+ b.HasIndex("QueueAddress")
+ .HasDatabaseName("ix_failed_messages_queue_address");
+
+ b.HasIndex("ReceivingEndpointName")
+ .HasDatabaseName("ix_failed_messages_receiving_endpoint_name");
+
+ b.HasIndex("StatusChangedAt")
+ .HasDatabaseName("ix_failed_messages_status_changed_at")
+ .HasFilter("status IN (2, 4)");
+
+ b.HasIndex("TimeSent")
+ .HasDatabaseName("ix_failed_messages_time_sent");
+
+ b.HasIndex("Status", "LastModified")
+ .HasDatabaseName("ix_failed_messages_status_last_modified");
+
+ b.ToTable("failed_messages", (string)null);
+ });
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageGroupEntity", b =>
+ {
+ b.Property("FailedMessageUniqueId")
+ .HasColumnType("uuid")
+ .HasColumnName("failed_message_unique_id");
+
+ b.Property("GroupId")
+ .HasMaxLength(64)
+ .HasColumnType("character varying(64)")
+ .HasColumnName("group_id");
+
+ b.Property("Title")
+ .IsRequired()
+ .HasColumnType("text")
+ .HasColumnName("title");
+
+ b.Property("Type")
+ .IsRequired()
+ .HasMaxLength(255)
+ .HasColumnType("character varying(255)")
+ .HasColumnName("type");
+
+ b.HasKey("FailedMessageUniqueId", "GroupId")
+ .HasName("pk_failed_message_groups");
+
+ b.HasIndex("GroupId")
+ .HasDatabaseName("ix_failed_message_groups_group_id");
+
+ b.ToTable("failed_message_groups", (string)null);
+ });
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageRetryEntity", b =>
+ {
+ b.Property("UniqueMessageId")
+ .HasColumnType("uuid")
+ .HasColumnName("unique_message_id");
+
+ b.Property("RetryId")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("retry_id");
+
+ b.HasKey("UniqueMessageId")
+ .HasName("pk_failed_message_retries");
+
+ b.ToTable("failed_message_retries", (string)null);
+ });
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.KnownEndpointEntity", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("uuid")
+ .HasColumnName("id");
+
+ b.Property("Host")
+ .IsRequired()
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("host");
+
+ b.Property("HostId")
+ .HasColumnType("uuid")
+ .HasColumnName("host_id");
+
+ b.Property("Monitored")
+ .HasColumnType("boolean")
+ .HasColumnName("monitored");
+
+ b.Property("Name")
+ .IsRequired()
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("name");
+
+ b.HasKey("Id")
+ .HasName("pk_known_endpoints");
+
+ b.ToTable("known_endpoints", (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
+ });
+ });
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageGroupEntity", b =>
+ {
+ b.HasOne("ServiceControl.Persistence.EFCore.Entities.FailedMessageEntity", null)
+ .WithMany()
+ .HasForeignKey("FailedMessageUniqueId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired()
+ .HasConstraintName("fk_failed_message_groups_failed_messages_failed_message_unique");
+ });
+#pragma warning restore 612, 618
+ }
+ }
+}
diff --git a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260724055936_AddConfiguration.cs b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260724055936_AddConfiguration.cs
new file mode 100644
index 0000000000..1f051f04da
--- /dev/null
+++ b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/20260724055936_AddConfiguration.cs
@@ -0,0 +1,77 @@
+using System;
+using Microsoft.EntityFrameworkCore.Migrations;
+using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata;
+
+#nullable disable
+
+namespace ServiceControl.Persistence.EFCore.PostgreSql.Migrations
+{
+ ///
+ public partial class AddConfiguration : Migration
+ {
+ ///
+ protected override void Up(MigrationBuilder migrationBuilder)
+ {
+ migrationBuilder.AddColumn(
+ name: "failing_endpoint_address",
+ table: "failed_messages",
+ type: "character varying(450)",
+ maxLength: 450,
+ nullable: false,
+ defaultValue: "");
+
+ migrationBuilder.CreateTable(
+ name: "EndpointSettings",
+ columns: table => new
+ {
+ name = table.Column(type: "character varying(450)", maxLength: 450, nullable: false),
+ track_instances = table.Column(type: "boolean", nullable: false)
+ },
+ constraints: table =>
+ {
+ table.PrimaryKey("pk_endpoint_settings", x => x.name);
+ });
+
+ 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_failing_endpoint_address",
+ table: "failed_messages",
+ column: "failing_endpoint_address");
+ }
+
+ ///
+ protected override void Down(MigrationBuilder migrationBuilder)
+ {
+ migrationBuilder.DropTable(
+ name: "EndpointSettings");
+
+ migrationBuilder.DropTable(
+ name: "trial_metadata");
+
+ migrationBuilder.DropIndex(
+ name: "ix_failed_messages_failing_endpoint_address",
+ table: "failed_messages");
+
+ migrationBuilder.DropColumn(
+ name: "failing_endpoint_address",
+ table: "failed_messages");
+ }
+ }
+}
diff --git a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/PostgreSqlServiceControlDbContextModelSnapshot.cs b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/PostgreSqlServiceControlDbContextModelSnapshot.cs
index fd525ada2c..1bfdf0f7ee 100644
--- a/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/PostgreSqlServiceControlDbContextModelSnapshot.cs
+++ b/src/ServiceControl.Persistence.EFCore.PostgreSql/Migrations/PostgreSqlServiceControlDbContextModelSnapshot.cs
@@ -22,6 +22,23 @@ protected override void BuildModel(ModelBuilder modelBuilder)
NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.EndpointSettingsEntity", b =>
+ {
+ b.Property("Name")
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .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("UniqueMessageId")
@@ -58,6 +75,12 @@ protected override void BuildModel(ModelBuilder modelBuilder)
.HasColumnType("text")
.HasColumnName("exception_type");
+ b.Property("FailingEndpointAddress")
+ .IsRequired()
+ .HasMaxLength(450)
+ .HasColumnType("character varying(450)")
+ .HasColumnName("failing_endpoint_address");
+
b.Property("FirstTimeOfFailure")
.HasColumnType("timestamp with time zone")
.HasColumnName("first_time_of_failure");
@@ -147,6 +170,9 @@ protected override void BuildModel(ModelBuilder modelBuilder)
b.HasIndex("ConversationId")
.HasDatabaseName("ix_failed_messages_conversation_id");
+ b.HasIndex("FailingEndpointAddress")
+ .HasDatabaseName("ix_failed_messages_failing_endpoint_address");
+
b.HasIndex("QueueAddress")
.HasDatabaseName("ix_failed_messages_queue_address");
@@ -246,6 +272,31 @@ protected override void BuildModel(ModelBuilder modelBuilder)
b.ToTable("known_endpoints", (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
+ });
+ });
+
modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageGroupEntity", b =>
{
b.HasOne("ServiceControl.Persistence.EFCore.Entities.FailedMessageEntity", null)
diff --git a/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260724055947_AddConfiguration.Designer.cs b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260724055947_AddConfiguration.Designer.cs
new file mode 100644
index 0000000000..13ee34cf70
--- /dev/null
+++ b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260724055947_AddConfiguration.Designer.cs
@@ -0,0 +1,257 @@
+//
+using System;
+using Microsoft.EntityFrameworkCore;
+using Microsoft.EntityFrameworkCore.Infrastructure;
+using Microsoft.EntityFrameworkCore.Metadata;
+using Microsoft.EntityFrameworkCore.Migrations;
+using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
+using ServiceControl.Persistence.EFCore.SqlServer;
+
+#nullable disable
+
+namespace ServiceControl.Persistence.EFCore.SqlServer.Migrations
+{
+ [DbContext(typeof(SqlServerServiceControlDbContext))]
+ [Migration("20260724055947_AddConfiguration")]
+ partial class AddConfiguration
+ {
+ ///
+ protected override void BuildTargetModel(ModelBuilder modelBuilder)
+ {
+#pragma warning disable 612, 618
+ modelBuilder
+ .HasAnnotation("ProductVersion", "10.0.10")
+ .HasAnnotation("Relational:MaxIdentifierLength", 128);
+
+ SqlServerModelBuilderExtensions.UseIdentityColumns(modelBuilder);
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.EndpointSettingsEntity", b =>
+ {
+ b.Property("Name")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("TrackInstances")
+ .HasColumnType("bit");
+
+ b.HasKey("Name");
+
+ b.ToTable("EndpointSettings", (string)null);
+ });
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageEntity", b =>
+ {
+ b.Property("UniqueMessageId")
+ .HasColumnType("uniqueidentifier");
+
+ b.Property("BodyContentType")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("BodySize")
+ .HasColumnType("int");
+
+ b.Property("BodyStoredExternally")
+ .HasColumnType("bit");
+
+ b.Property("BodyText")
+ .HasColumnType("nvarchar(max)");
+
+ b.Property("ConversationId")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("ExceptionMessage")
+ .HasColumnType("nvarchar(max)");
+
+ b.Property("ExceptionType")
+ .HasColumnType("nvarchar(max)");
+
+ b.Property("FailingEndpointAddress")
+ .IsRequired()
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("FirstTimeOfFailure")
+ .HasColumnType("datetime2");
+
+ b.Property("HeadersJson")
+ .IsRequired()
+ .HasColumnType("nvarchar(max)");
+
+ b.Property("IsSystemMessage")
+ .HasColumnType("bit");
+
+ b.Property("LastAttemptedAt")
+ .HasColumnType("datetime2");
+
+ b.Property("LastModified")
+ .HasColumnType("datetime2");
+
+ b.Property("LastTimeOfFailure")
+ .HasColumnType("datetime2");
+
+ b.Property("MessageId")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("MessageType")
+ .HasColumnType("nvarchar(max)");
+
+ b.Property("NumberOfProcessingAttempts")
+ .HasColumnType("int");
+
+ b.Property("QueueAddress")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("ReceivingEndpointHost")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("ReceivingEndpointHostId")
+ .HasColumnType("uniqueidentifier");
+
+ b.Property("ReceivingEndpointName")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("SendingEndpointHost")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("SendingEndpointHostId")
+ .HasColumnType("uniqueidentifier");
+
+ b.Property("SendingEndpointName")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("Status")
+ .HasColumnType("int");
+
+ b.Property("StatusChangedAt")
+ .HasColumnType("datetime2");
+
+ b.Property("TimeSent")
+ .HasColumnType("datetime2");
+
+ b.HasKey("UniqueMessageId");
+
+ b.HasIndex("ConversationId");
+
+ b.HasIndex("FailingEndpointAddress");
+
+ b.HasIndex("QueueAddress");
+
+ b.HasIndex("ReceivingEndpointName");
+
+ b.HasIndex("StatusChangedAt")
+ .HasFilter("[Status] IN (2, 4)");
+
+ b.HasIndex("TimeSent");
+
+ b.HasIndex("Status", "LastModified");
+
+ b.ToTable("FailedMessages");
+ });
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageGroupEntity", b =>
+ {
+ b.Property("FailedMessageUniqueId")
+ .HasColumnType("uniqueidentifier");
+
+ b.Property("GroupId")
+ .HasMaxLength(64)
+ .HasColumnType("nvarchar(64)");
+
+ b.Property("Title")
+ .IsRequired()
+ .HasColumnType("nvarchar(max)");
+
+ b.Property("Type")
+ .IsRequired()
+ .HasMaxLength(255)
+ .HasColumnType("nvarchar(255)");
+
+ b.HasKey("FailedMessageUniqueId", "GroupId");
+
+ b.HasIndex("GroupId");
+
+ b.ToTable("FailedMessageGroups");
+ });
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageRetryEntity", b =>
+ {
+ b.Property("UniqueMessageId")
+ .HasColumnType("uniqueidentifier");
+
+ b.Property("RetryId")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.HasKey("UniqueMessageId");
+
+ b.ToTable("FailedMessageRetries");
+ });
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.KnownEndpointEntity", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("uniqueidentifier");
+
+ b.Property("Host")
+ .IsRequired()
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("HostId")
+ .HasColumnType("uniqueidentifier");
+
+ b.Property("Monitored")
+ .HasColumnType("bit");
+
+ b.Property("Name")
+ .IsRequired()
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.HasKey("Id");
+
+ b.ToTable("KnownEndpoints");
+ });
+
+ 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
+ });
+ });
+
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageGroupEntity", b =>
+ {
+ b.HasOne("ServiceControl.Persistence.EFCore.Entities.FailedMessageEntity", null)
+ .WithMany()
+ .HasForeignKey("FailedMessageUniqueId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired();
+ });
+#pragma warning restore 612, 618
+ }
+ }
+}
diff --git a/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260724055947_AddConfiguration.cs b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260724055947_AddConfiguration.cs
new file mode 100644
index 0000000000..6f816c4835
--- /dev/null
+++ b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/20260724055947_AddConfiguration.cs
@@ -0,0 +1,76 @@
+using System;
+using Microsoft.EntityFrameworkCore.Migrations;
+
+#nullable disable
+
+namespace ServiceControl.Persistence.EFCore.SqlServer.Migrations
+{
+ ///
+ public partial class AddConfiguration : Migration
+ {
+ ///
+ protected override void Up(MigrationBuilder migrationBuilder)
+ {
+ migrationBuilder.AddColumn(
+ name: "FailingEndpointAddress",
+ table: "FailedMessages",
+ type: "nvarchar(450)",
+ maxLength: 450,
+ nullable: false,
+ defaultValue: "");
+
+ migrationBuilder.CreateTable(
+ name: "EndpointSettings",
+ columns: table => new
+ {
+ Name = table.Column(type: "nvarchar(450)", maxLength: 450, nullable: false),
+ TrackInstances = table.Column(type: "bit", nullable: false)
+ },
+ constraints: table =>
+ {
+ table.PrimaryKey("PK_EndpointSettings", x => x.Name);
+ });
+
+ 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_FailingEndpointAddress",
+ table: "FailedMessages",
+ column: "FailingEndpointAddress");
+ }
+
+ ///
+ protected override void Down(MigrationBuilder migrationBuilder)
+ {
+ migrationBuilder.DropTable(
+ name: "EndpointSettings");
+
+ migrationBuilder.DropTable(
+ name: "TrialMetadata");
+
+ migrationBuilder.DropIndex(
+ name: "IX_FailedMessages_FailingEndpointAddress",
+ table: "FailedMessages");
+
+ migrationBuilder.DropColumn(
+ name: "FailingEndpointAddress",
+ table: "FailedMessages");
+ }
+ }
+}
diff --git a/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/SqlServerServiceControlDbContextModelSnapshot.cs b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/SqlServerServiceControlDbContextModelSnapshot.cs
index cceec018f8..0a2cd26beb 100644
--- a/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/SqlServerServiceControlDbContextModelSnapshot.cs
+++ b/src/ServiceControl.Persistence.EFCore.SqlServer/Migrations/SqlServerServiceControlDbContextModelSnapshot.cs
@@ -22,6 +22,20 @@ protected override void BuildModel(ModelBuilder modelBuilder)
SqlServerModelBuilderExtensions.UseIdentityColumns(modelBuilder);
+ modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.EndpointSettingsEntity", b =>
+ {
+ b.Property("Name")
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
+ b.Property("TrackInstances")
+ .HasColumnType("bit");
+
+ b.HasKey("Name");
+
+ b.ToTable("EndpointSettings", (string)null);
+ });
+
modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageEntity", b =>
{
b.Property("UniqueMessageId")
@@ -50,6 +64,11 @@ protected override void BuildModel(ModelBuilder modelBuilder)
b.Property("ExceptionType")
.HasColumnType("nvarchar(max)");
+ b.Property("FailingEndpointAddress")
+ .IsRequired()
+ .HasMaxLength(450)
+ .HasColumnType("nvarchar(450)");
+
b.Property("FirstTimeOfFailure")
.HasColumnType("datetime2");
@@ -118,6 +137,8 @@ protected override void BuildModel(ModelBuilder modelBuilder)
b.HasIndex("ConversationId");
+ b.HasIndex("FailingEndpointAddress");
+
b.HasIndex("QueueAddress");
b.HasIndex("ReceivingEndpointName");
@@ -197,6 +218,28 @@ protected override void BuildModel(ModelBuilder modelBuilder)
b.ToTable("KnownEndpoints");
});
+ 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
+ });
+ });
+
modelBuilder.Entity("ServiceControl.Persistence.EFCore.Entities.FailedMessageGroupEntity", b =>
{
b.HasOne("ServiceControl.Persistence.EFCore.Entities.FailedMessageEntity", null)
diff --git a/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs b/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs
index 49990d4363..b65d8ce469 100644
--- a/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs
+++ b/src/ServiceControl.Persistence.EFCore/DbContexts/ServiceControlDbContext.cs
@@ -6,23 +6,25 @@ 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 FailedMessages { get; set; }
public DbSet FailedMessageGroups { get; set; }
public DbSet FailedMessageRetries { 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 KnownEndpointConfiguration());
+ modelBuilder.ApplyConfiguration(new EndpointSettingsConfiguration());
modelBuilder.ApplyConfiguration(new FailedMessageConfiguration());
modelBuilder.ApplyConfiguration(new FailedMessageGroupConfiguration());
modelBuilder.ApplyConfiguration(new FailedMessageRetryConfiguration());
+ modelBuilder.ApplyConfiguration(new KnownEndpointConfiguration());
+ modelBuilder.ApplyConfiguration(new TrialMetadataConfiguration());
}
}
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
index f9b7ee810a..14bdb2f27f 100644
--- a/src/ServiceControl.Persistence.EFCore/Entities/FailedMessageEntity.cs
+++ b/src/ServiceControl.Persistence.EFCore/Entities/FailedMessageEntity.cs
@@ -57,4 +57,6 @@ public class FailedMessageEntity
public int BodySize { get; set; }
public string? BodyContentType { get; set; }
+
+ public required string FailingEndpointAddress { get; set; }
}
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/EndpointSettingsConfiguration.cs b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsConfiguration.cs
new file mode 100644
index 0000000000..aaac494746
--- /dev/null
+++ b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/EndpointSettingsConfiguration.cs
@@ -0,0 +1,15 @@
+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);
+ builder.Property(x => x.Name).HasMaxLength(ColumnLengths.ShortTextLength).IsRequired();
+ }
+}
\ No newline at end of file
diff --git a/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageConfiguration.cs b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageConfiguration.cs
index 1702985c21..a378c6593b 100644
--- a/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageConfiguration.cs
+++ b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/FailedMessageConfiguration.cs
@@ -26,6 +26,7 @@ public void Configure(EntityTypeBuilder builder)
builder.Property(e => e.SendingEndpointHost).HasMaxLength(ColumnLengths.ShortTextLength);
builder.Property(e => e.ReceivingEndpointName).HasMaxLength(ColumnLengths.ShortTextLength);
builder.Property(e => e.ReceivingEndpointHost).HasMaxLength(ColumnLengths.ShortTextLength);
+ builder.Property(e => e.FailingEndpointAddress).HasMaxLength(ColumnLengths.ShortTextLength);
builder.Property(e => e.BodyContentType).HasMaxLength(ColumnLengths.ShortTextLength);
builder.Property(e => e.IsSystemMessage).IsRequired();
@@ -35,6 +36,7 @@ public void Configure(EntityTypeBuilder builder)
builder.HasIndex(e => new { e.Status, e.LastModified });
builder.HasIndex(e => e.ReceivingEndpointName);
+ builder.HasIndex(e => e.FailingEndpointAddress);
builder.HasIndex(e => e.ConversationId);
builder.HasIndex(e => e.TimeSent);
builder.HasIndex(e => e.QueueAddress);
diff --git a/src/ServiceControl.Persistence.EFCore/EntityConfigurations/TrialMetadataConfiguration.cs b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/TrialMetadataConfiguration.cs
new file mode 100644
index 0000000000..c6992ff3f0
--- /dev/null
+++ b/src/ServiceControl.Persistence.EFCore/EntityConfigurations/TrialMetadataConfiguration.cs
@@ -0,0 +1,18 @@
+namespace ServiceControl.Persistence.EFCore.EntityConfigurations;
+
+using Entities;
+using Microsoft.EntityFrameworkCore;
+using Microsoft.EntityFrameworkCore.Metadata.Builders;
+
+public class TrialMetadataConfiguration : 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 157450c54b..90299c070e 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,19 @@ 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, [EnumeratorCancellation] CancellationToken cancellationToken = default)
+ {
+ 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..cb773ea3ca 100644
--- a/src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs
+++ b/src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs
@@ -1,13 +1,45 @@
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 (DbUpdateException)
+ {
+ //this probably failed because of key conflict so try again
+ }
+
+ await context.Entry(entity).ReloadAsync(token);
+ }
- public Task UpdateEndpointSettings(EndpointSettings settings, CancellationToken token) =>
- throw new NotImplementedException();
+ 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.EndpointSettings.RemoveRange(context.EndpointSettings.Where(x => x.Name == name));
+ await context.SaveChangesAsync(cancellationToken);
+ });
}
diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/ErrorMessagesDataStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/ErrorMessagesDataStore.cs
index c081aa0674..692ce2b0c9 100644
--- a/src/ServiceControl.Persistence.EFCore/Implementation/ErrorMessagesDataStore.cs
+++ b/src/ServiceControl.Persistence.EFCore/Implementation/ErrorMessagesDataStore.cs
@@ -1,5 +1,8 @@
namespace ServiceControl.Persistence.EFCore.Implementation;
+using Entities;
+using Microsoft.Extensions.DependencyInjection;
+using Persistence.UnitOfWork;
using ServiceControl.CompositeViews.Messages;
using ServiceControl.EventLog;
using ServiceControl.MessageFailures;
@@ -8,7 +11,7 @@ namespace ServiceControl.Persistence.EFCore.Implementation;
using ServiceControl.Persistence.Infrastructure;
using ServiceControl.Recoverability;
-public class ErrorMessagesDataStore : IErrorMessageDataStore
+public class ErrorMessagesDataStore(IServiceScopeFactory scopeFactory) : DataStoreBase(scopeFactory), IErrorMessageDataStore
{
public Task>> GetAllMessages(PagingInfo pagingInfo, SortInfo sortInfo, bool includeSystemMessages, DateTimeRange? timeSentRange = null) =>
throw new NotImplementedException();
@@ -110,7 +113,4 @@ public Task FetchFromFailedMessage(string uniqueMessageId) =>
public Task StoreEventLogItem(EventLogItem logItem) =>
throw new NotImplementedException();
-
- public Task StoreFailedMessagesForTestsOnly(params FailedMessage[] failedMessages) =>
- throw new NotImplementedException();
}
diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs
index b53673b630..553d36f03e 100644
--- a/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs
+++ b/src/ServiceControl.Persistence.EFCore/Implementation/QueueAddressStore.cs
@@ -1,13 +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) =>
- throw new NotImplementedException();
+ 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.FailingEndpointAddress.StartsWith(search))
+ .Select(fm => new { fm.UniqueMessageId, fm.FailingEndpointAddress })
+ .Distinct()
+ .GroupBy(failure => failure.FailingEndpointAddress)
+ .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 e10739ee89..44a317d0bc 100644
--- a/src/ServiceControl.Persistence.EFCore/Implementation/TrialLicenseDataProvider.cs
+++ b/src/ServiceControl.Persistence.EFCore/Implementation/TrialLicenseDataProvider.cs
@@ -1,10 +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) =>
- throw new NotImplementedException();
+ 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.EFCore/Implementation/UnitOfWork/EFRecoverabilityIngestionUnitOfWork.cs b/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/EFRecoverabilityIngestionUnitOfWork.cs
index d33a539cdc..f6afd2530d 100644
--- a/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/EFRecoverabilityIngestionUnitOfWork.cs
+++ b/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/EFRecoverabilityIngestionUnitOfWork.cs
@@ -50,6 +50,7 @@ public Task RecordFailedProcessingAttempt(MessageContext context,
ReceivingEndpointHost = receivingEndpoint?.Host,
ExceptionType = processingAttempt.FailureDetails.Exception?.ExceptionType,
ExceptionMessage = processingAttempt.FailureDetails.Exception?.Message,
+ FailingEndpointAddress = processingAttempt.FailureDetails.AddressOfFailingEndpoint,
IsSystemMessage = GetMetadata(processingAttempt, "IsSystemMessage"),
BodyText = bodyText,
BodyStoredExternally = storeExternally,
diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/FailedMessageBatchWriter.cs b/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/FailedMessageBatchWriter.cs
index f2974c75b0..7ebc091c79 100644
--- a/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/FailedMessageBatchWriter.cs
+++ b/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/FailedMessageBatchWriter.cs
@@ -92,7 +92,8 @@ await strategy.ExecuteAsync(async ct =>
BodyText = last.BodyText,
BodyStoredExternally = last.BodyStoredExternally,
BodySize = last.BodySize,
- BodyContentType = last.BodyContentType
+ BodyContentType = last.BodyContentType,
+ FailingEndpointAddress = last.FailingEndpointAddress
});
groups.AddRange(last.Groups
diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/RecordedFailedProcessingAttempt.cs b/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/RecordedFailedProcessingAttempt.cs
index 82b6c047f4..b9bdf84228 100644
--- a/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/RecordedFailedProcessingAttempt.cs
+++ b/src/ServiceControl.Persistence.EFCore/Implementation/UnitOfWork/RecordedFailedProcessingAttempt.cs
@@ -29,4 +29,5 @@ sealed class RecordedFailedProcessingAttempt
public bool BodyStoredExternally { get; init; }
public int BodySize { get; init; }
public string? BodyContentType { get; init; }
+ public required string FailingEndpointAddress { get; set; }
}
diff --git a/src/ServiceControl.Persistence.EFCore/ServiceControl.Persistence.EFCore.csproj b/src/ServiceControl.Persistence.EFCore/ServiceControl.Persistence.EFCore.csproj
index f80428271a..6e18b274bb 100644
--- a/src/ServiceControl.Persistence.EFCore/ServiceControl.Persistence.EFCore.csproj
+++ b/src/ServiceControl.Persistence.EFCore/ServiceControl.Persistence.EFCore.csproj
@@ -24,4 +24,8 @@
+
+
+
+
diff --git a/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs b/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs
index df59b8fdf9..730de1cf57 100644
--- a/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs
+++ b/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs
@@ -663,16 +663,5 @@ public async Task StoreEventLogItem(EventLogItem logItem)
await session.SaveChangesAsync();
}
-
- public async Task StoreFailedMessagesForTestsOnly(params FailedMessage[] failedMessages)
- {
- using var session = await sessionProvider.OpenSession();
- foreach (var message in failedMessages)
- {
- await session.StoreAsync(message);
- }
-
- await session.SaveChangesAsync();
- }
}
}
diff --git a/src/ServiceControl.Persistence.Tests.InMemory/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.InMemory/PersistenceTestsContext.cs
index 2b17a1fa7e..e830dd8bbd 100644
--- a/src/ServiceControl.Persistence.Tests.InMemory/PersistenceTestsContext.cs
+++ b/src/ServiceControl.Persistence.Tests.InMemory/PersistenceTestsContext.cs
@@ -1,6 +1,7 @@
namespace ServiceControl.Persistence.Tests;
using System.Threading.Tasks;
+using MessageFailures;
using Microsoft.Extensions.Hosting;
using Particular.LicensingComponent.Persistence.InMemory;
using ServiceControl.Persistence.Tests.InMemory;
@@ -9,6 +10,7 @@ public class PersistenceTestsContext : IPersistenceTestsContext
{
public PersistenceSettings PersistenceSettings { get; private set; }
public string GenerateFailedMessageRecordId(string messageId) => throw new System.NotImplementedException();
+ public Task InsertFailedMessages(params FailedMessage[] messages) => throw new System.NotImplementedException();
public Task Setup(IHostApplicationBuilder hostBuilder)
{
diff --git a/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs
index a104f188eb..8053e89c41 100644
--- a/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs
+++ b/src/ServiceControl.Persistence.Tests.PostgreSql/PersistenceTestsContext.cs
@@ -1,10 +1,12 @@
// ReSharper disable once CheckNamespace
+
namespace ServiceControl.Persistence.Tests;
using System;
using System.Linq;
using System.Threading.Tasks;
using EFCore.PostgreSql;
+using MessageFailures;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
@@ -12,7 +14,7 @@ namespace ServiceControl.Persistence.Tests;
using Npgsql;
using ServiceControl.Persistence.EFCore.Infrastructure;
-public class PersistenceTestsContext : IPersistenceTestsContext
+public partial class PersistenceTestsContext : IPersistenceTestsContext
{
IHost host;
string databaseName;
@@ -23,15 +25,9 @@ public async Task Setup(IHostApplicationBuilder hostBuilder)
{
databaseName = $"sc_test_{Guid.NewGuid():n}";
- var connectionStringBuilder = new NpgsqlConnectionStringBuilder(await PostgreSqlSharedContainer.GetConnectionStringAsync())
- {
- Database = databaseName
- };
+ var connectionStringBuilder = new NpgsqlConnectionStringBuilder(await PostgreSqlSharedContainer.GetConnectionStringAsync()) { Database = databaseName };
- PersistenceSettings = new PostgreSqlPersisterSettings
- {
- ConnectionString = connectionStringBuilder.ConnectionString
- };
+ PersistenceSettings = new PostgreSqlPersisterSettings { ConnectionString = connectionStringBuilder.ConnectionString };
var persistence = new PostgreSqlPersistenceConfiguration().Create(PersistenceSettings);
@@ -72,4 +68,7 @@ public async Task CompleteDatabaseOperation()
public PersistenceSettings PersistenceSettings { get; set; }
public string GenerateFailedMessageRecordId(string messageId) => messageId;
-}
+
+ public Task InsertFailedMessages(params FailedMessage[] messages) => InsertFailedMessagesDirect(host.Services, messages);
+
+}
\ No newline at end of file
diff --git a/src/ServiceControl.Persistence.Tests.RavenDB/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.RavenDB/PersistenceTestsContext.cs
index 63e455105c..ebd1f288ab 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;
@@ -54,6 +56,15 @@ public async Task PostSetup(IHost host)
public PersistenceSettings PersistenceSettings { get; private set; }
public string GenerateFailedMessageRecordId(string messageId) => FailedMessageIdGenerator.MakeDocumentId(messageId);
+ public async Task InsertFailedMessages(params FailedMessage[] messages)
+ {
+ using var session = await SessionProvider.OpenSession();
+ foreach (var message in messages)
+ {
+ await session.StoreAsync(message);
+ }
+ await session.SaveChangesAsync();
+ }
public IDocumentStore DocumentStore { get; private set; }
diff --git a/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs
index 9af5dc6de0..17ff011e77 100644
--- a/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs
+++ b/src/ServiceControl.Persistence.Tests.SqlServer/PersistenceTestsContext.cs
@@ -5,6 +5,7 @@ namespace ServiceControl.Persistence.Tests;
using System.Linq;
using System.Threading.Tasks;
using EFCore.SqlServer;
+using MessageFailures;
using Microsoft.Data.SqlClient;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
@@ -12,7 +13,7 @@ namespace ServiceControl.Persistence.Tests;
using Microsoft.Extensions.Time.Testing;
using ServiceControl.Persistence.EFCore.Infrastructure;
-public class PersistenceTestsContext : IPersistenceTestsContext
+public partial class PersistenceTestsContext : IPersistenceTestsContext
{
IHost host;
string databaseName;
@@ -78,4 +79,6 @@ public async Task CompleteDatabaseOperation()
public PersistenceSettings PersistenceSettings { get; set; }
public string GenerateFailedMessageRecordId(string messageId) => messageId;
+
+ public Task InsertFailedMessages(params FailedMessage[] messages) => InsertFailedMessagesDirect(host.Services, messages);
}
\ No newline at end of file
diff --git a/src/ServiceControl.Persistence.Tests/EFCore/PersistenceTestsContext.cs b/src/ServiceControl.Persistence.Tests/EFCore/PersistenceTestsContext.cs
new file mode 100644
index 0000000000..d63e6a85e6
--- /dev/null
+++ b/src/ServiceControl.Persistence.Tests/EFCore/PersistenceTestsContext.cs
@@ -0,0 +1,73 @@
+#nullable enable
+namespace ServiceControl.Persistence.Tests;
+
+using System;
+using System.Collections.Generic;
+using System.Linq;
+using System.Text.Json;
+using System.Threading.Tasks;
+using EFCore.DbContexts;
+using EFCore.Entities;
+using EFCore.Implementation.UnitOfWork;
+using MessageFailures;
+using Microsoft.Extensions.DependencyInjection;
+using NServiceBus;
+using Operations;
+
+public partial class PersistenceTestsContext
+{
+ static async Task InsertFailedMessagesDirect(IServiceProvider serviceProvider, FailedMessage[] messages)
+ {
+ await using var scope = serviceProvider.CreateAsyncScope();
+ await using var db = scope.ServiceProvider.GetRequiredService();
+
+ static T? GetMetadata(FailedMessage.ProcessingAttempt processingAttempt, string key) =>
+ processingAttempt.MessageMetadata.TryGetValue(key, out var value) && value is T typed ? typed : default;
+
+ foreach (FailedMessage failedMessage in messages)
+ {
+ var now = DateTime.UtcNow;
+ var ordered = failedMessage.ProcessingAttempts
+ .OrderBy(x => x.AttemptedAt)
+ .ThenBy(x => x.MessageId)
+ .ToList();
+ var attempt = ordered.Last();
+ var sendingEndpoint = GetMetadata(attempt, "SendingEndpoint");
+ var receivingEndpoint = GetMetadata(attempt, "ReceivingEndpoint");
+ var contentType = attempt.Headers.GetValueOrDefault(Headers.ContentType, "text/plain");
+ db.FailedMessages.Add(new FailedMessageEntity
+ {
+ UniqueMessageId = Guid.Parse(failedMessage.UniqueMessageId),
+ FirstTimeOfFailure = ordered.Min(pa => pa.FailureDetails.TimeOfFailure),
+ LastTimeOfFailure = ordered.Max(pa => pa.FailureDetails.TimeOfFailure),
+ LastAttemptedAt = attempt.AttemptedAt,
+ HeadersJson = JsonSerializer.Serialize(attempt.Headers, HeadersJsonContext.Default.DictionaryStringString),
+ MessageId = attempt.MessageId,
+ MessageType = GetMetadata(attempt, "MessageType"),
+ TimeSent = GetMetadata(attempt, "TimeSent"),
+ ConversationId = GetMetadata(attempt, "ConversationId"),
+ QueueAddress = attempt.Headers.GetValueOrDefault(NServiceBus.Faults.FaultsHeaderKeys.FailedQ),
+ SendingEndpointName = sendingEndpoint?.Name,
+ SendingEndpointHostId = sendingEndpoint?.HostId,
+ SendingEndpointHost = sendingEndpoint?.Host,
+ ReceivingEndpointName = receivingEndpoint?.Name,
+ ReceivingEndpointHostId = receivingEndpoint?.HostId,
+ ReceivingEndpointHost = receivingEndpoint?.Host,
+ ExceptionType = attempt.FailureDetails.Exception?.ExceptionType,
+ ExceptionMessage = attempt.FailureDetails.Exception?.Message,
+ FailingEndpointAddress = attempt.FailureDetails.AddressOfFailingEndpoint,
+ IsSystemMessage = GetMetadata(attempt, "IsSystemMessage"),
+ BodyText = attempt.Body,
+ BodyStoredExternally = false,
+ BodySize = attempt.Body?.Length ?? 0,
+ BodyContentType = contentType,
+ Status = failedMessage.Status,
+ StatusChangedAt = now,
+ LastModified = now,
+ NumberOfProcessingAttempts = ordered.Select(pa => pa.AttemptedAt).Distinct().Count(),
+ });
+ }
+
+ await db.SaveChangesAsync();
+ }
+}
\ No newline at end of file
diff --git a/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs
index 0c9b614908..eb2337e595 100644
--- a/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs
+++ b/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs
@@ -134,7 +134,8 @@ await Store(new FailedMessageEntity
IsSystemMessage = false,
HeadersJson = "{}",
BodyStoredExternally = bodyStoredExternally,
- BodySize = 0
+ BodySize = 0,
+ FailingEndpointAddress = "Shipping"
});
return id;
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..abf74b0b33 100644
--- a/src/ServiceControl.Persistence.Tests/IPersistenceTestsContext.cs
+++ b/src/ServiceControl.Persistence.Tests/IPersistenceTestsContext.cs
@@ -1,6 +1,7 @@
namespace ServiceControl.Persistence.Tests;
using System.Threading.Tasks;
+using MessageFailures;
using Microsoft.Extensions.Hosting;
public interface IPersistenceTestsContext
@@ -16,4 +17,5 @@ public interface IPersistenceTestsContext
PersistenceSettings PersistenceSettings { get; }
string GenerateFailedMessageRecordId(string messageId);
+ Task InsertFailedMessages(params FailedMessage[] messages);
}
\ No newline at end of file
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..f4501923d8
--- /dev/null
+++ b/src/ServiceControl.Persistence.Tests/QueueAddressStoreTests.cs
@@ -0,0 +1,74 @@
+namespace ServiceControl.Persistence.Tests;
+
+using System;
+using System.Linq;
+using System.Threading.Tasks;
+using Contracts.Operations;
+using MessageFailures;
+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 SeedFailedMessages(
+ (new("7F21F22A-44B6-440C-851D-3524645FD083"), "alpha"),
+ (new("E5DF7B16-648E-47F2-A9FD-F2BC1BA5D53C"), "alpha"),
+ (new("8975FE11-DF21-438B-BB65-0006541CA73D"), "beta"),
+ (new("7ED8EEA2-66EA-496C-922F-87D53A699228"), "alpha-child"),
+ (new("D98BDA7B-B3D5-4CD8-A802-E38B57A58041"), "gamma"));
+
+ public Task SeedFailedMessages(params (Guid MessageId, string EndpointAddress)[] failedMessages) =>
+ PersistenceTestsContext.InsertFailedMessages(failedMessages.Select(f => new FailedMessage()
+ {
+ Id = f.MessageId.ToString(),
+ UniqueMessageId = f.MessageId.ToString(),
+ ProcessingAttempts = [
+ new ()
+ {
+ FailureDetails = new FailureDetails() { AddressOfFailingEndpoint = f.EndpointAddress,}
+ }
+ ]
+ }).ToArray());
+}
\ No newline at end of file
diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/EditHandlerAuditTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/EditHandlerAuditTests.cs
index df02f36532..f1622dca30 100644
--- a/src/ServiceControl.Persistence.Tests/Recoverability/EditHandlerAuditTests.cs
+++ b/src/ServiceControl.Persistence.Tests/Recoverability/EditHandlerAuditTests.cs
@@ -105,7 +105,7 @@ async Task CreateAndStoreFailedMessage(string failedMessageId = n
}
]
};
- await ErrorMessageDataStore.StoreFailedMessagesForTestsOnly(new[] { failedMessage });
+ await PersistenceTestsContext.InsertFailedMessages(failedMessage);
return failedMessage;
}
}
\ No newline at end of file
diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/EditMessageTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/EditMessageTests.cs
index 7836a00e0f..ae9f7298be 100644
--- a/src/ServiceControl.Persistence.Tests/Recoverability/EditMessageTests.cs
+++ b/src/ServiceControl.Persistence.Tests/Recoverability/EditMessageTests.cs
@@ -262,7 +262,7 @@ async Task CreateAndStoreFailedMessage(string failedMessageId = n
}
]
};
- await ErrorMessageDataStore.StoreFailedMessagesForTestsOnly(new[] { failedMessage });
+ await PersistenceTestsContext.InsertFailedMessages(failedMessage);
return failedMessage;
}
}
diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/RetryConfirmationProcessorTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/RetryConfirmationProcessorTests.cs
index e9d720a9e1..3fd4a136fb 100644
--- a/src/ServiceControl.Persistence.Tests/Recoverability/RetryConfirmationProcessorTests.cs
+++ b/src/ServiceControl.Persistence.Tests/Recoverability/RetryConfirmationProcessorTests.cs
@@ -25,7 +25,7 @@ public async Task Setup()
Handler = new LegacyMessageFailureResolvedHandler(ErrorMessageDataStore, domainEvents);
- await ErrorMessageDataStore.StoreFailedMessagesForTestsOnly(
+ await PersistenceTestsContext.InsertFailedMessages(
new FailedMessage
{
Id = MessageId,
diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs
index 5909b71b5d..d91b27da72 100644
--- a/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs
+++ b/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs
@@ -243,8 +243,6 @@ public Task>> ErrorsByEndpointName(string s
public Task FetchFromFailedMessage(string bodyId) => Task.FromResult(Encoding.UTF8.GetBytes(bodyId));
public Task StoreEventLogItem(EventLogItem logItem) => throw new NotImplementedException();
-
- public Task StoreFailedMessagesForTestsOnly(params FailedMessage[] failedMessages) => throw new NotImplementedException();
}
}
}
\ No newline at end of file
diff --git a/src/ServiceControl.Persistence.Tests/RetryStateTests.cs b/src/ServiceControl.Persistence.Tests/RetryStateTests.cs
index 406388eda7..4c11092471 100644
--- a/src/ServiceControl.Persistence.Tests/RetryStateTests.cs
+++ b/src/ServiceControl.Persistence.Tests/RetryStateTests.cs
@@ -254,7 +254,7 @@ public async Task When_a_selection_is_staged_each_message_is_audited_as_a_batch(
]
}).ToArray();
- await ErrorStore.StoreFailedMessagesForTestsOnly(messages);
+ await PersistenceTestsContext.InsertFailedMessages(messages);
await CompleteDatabaseOperation();
var gateway = new CustomRetriesGateway(true, RetryStore, retryManager);
@@ -342,7 +342,7 @@ async Task CreateAFailedMessageAndMarkAsPartOfRetryBatch(RetryingManager retryMa
]
}).ToArray();
- await ErrorStore.StoreFailedMessagesForTestsOnly(messages);
+ await PersistenceTestsContext.InsertFailedMessages(messages);
// Needs index FailedMessages_ByGroup
// Needs index FailedMessages_UniqueMessageIdAndTimeOfFailures
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));
+ }
+}
diff --git a/src/ServiceControl.Persistence/IEndpointSettingsStore.cs b/src/ServiceControl.Persistence/IEndpointSettingsStore.cs
index d41296dd3a..2b6c52e8a1 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 = default);
Task UpdateEndpointSettings(EndpointSettings settings, CancellationToken token);
Task Delete(string name, CancellationToken cancellationToken);
diff --git a/src/ServiceControl.Persistence/IErrorMessageDatastore.cs b/src/ServiceControl.Persistence/IErrorMessageDatastore.cs
index 0dd2c6aad0..28ad891425 100644
--- a/src/ServiceControl.Persistence/IErrorMessageDatastore.cs
+++ b/src/ServiceControl.Persistence/IErrorMessageDatastore.cs
@@ -71,7 +71,5 @@ public interface IErrorMessageDataStore
// AuditEventLogWriter
Task StoreEventLogItem(EventLogItem logItem);
-
- Task StoreFailedMessagesForTestsOnly(params FailedMessage[] failedMessages);
}
}
\ No newline at end of file
diff --git a/src/ServiceControl.UnitTests/MessageFailures/AsyncRangeAndQueueAuditTests.cs b/src/ServiceControl.UnitTests/MessageFailures/AsyncRangeAndQueueAuditTests.cs
index ffb2dcae62..9c26e50c2f 100644
--- a/src/ServiceControl.UnitTests/MessageFailures/AsyncRangeAndQueueAuditTests.cs
+++ b/src/ServiceControl.UnitTests/MessageFailures/AsyncRangeAndQueueAuditTests.cs
@@ -192,6 +192,5 @@ internal sealed class StubErrorMessageDataStore : IErrorMessageDataStore
public Task RevertRetry(string messageUniqueId) => throw new NotImplementedException();
public Task FetchFromFailedMessage(string uniqueMessageId) => throw new NotImplementedException();
public Task StoreEventLogItem(EventLogItem logItem) => throw new NotImplementedException();
- public Task StoreFailedMessagesForTestsOnly(params FailedMessage[] failedMessages) => throw new NotImplementedException();
}
}
diff --git a/src/ServiceControl.UnitTests/MessageFailures/EditFailedMessagesControllerAuditTests.cs b/src/ServiceControl.UnitTests/MessageFailures/EditFailedMessagesControllerAuditTests.cs
index 96cbb10995..4cb8aa0738 100644
--- a/src/ServiceControl.UnitTests/MessageFailures/EditFailedMessagesControllerAuditTests.cs
+++ b/src/ServiceControl.UnitTests/MessageFailures/EditFailedMessagesControllerAuditTests.cs
@@ -97,6 +97,5 @@ sealed class StubErrorMessageDataStore : IErrorMessageDataStore
public Task GetRetryPendingMessages(DateTime from, DateTime to, string queueAddress) => throw new NotImplementedException();
public Task FetchFromFailedMessage(string uniqueMessageId) => throw new NotImplementedException();
public Task StoreEventLogItem(EventLogItem logItem) => throw new NotImplementedException();
- public Task StoreFailedMessagesForTestsOnly(params FailedMessage[] failedMessages) => throw new NotImplementedException();
}
}