From eac0c2358937c86ac0c0e1f72f940759b93dc219 Mon Sep 17 00:00:00 2001 From: Laurence Gillian Date: Mon, 20 Jul 2026 08:26:03 +0100 Subject: [PATCH] feat(postgresql): clear transport queue tables on message-store reset ClearAllAsync / RebuildAsync previously truncated only the envelope tables (incoming/outgoing/dead-letter/node/listener), leaving the PostgreSQL queue transport's queue and scheduled-message tables populated. Integration tests over the Postgres queue transport therefore carried rows between runs. Add a per-provider hook, truncateAdditionalTablesAsync(DbTransaction), invoked inside the existing reset transaction so the whole reset stays atomic. The PostgreSQL store overrides it to delete from the transport's own queue + scheduled tables. Scoped to the transport's table types on purpose: AddTable is a general registration path (SQL Server registers a rate-limit table through it that a reset must keep), so a blanket base-class loop would be wrong. Default behaviour for every other provider is unchanged. No public API added. Co-Authored-By: Claude Opus 4.8 (1M context) --- .../reset_clears_transport_queue_tables.cs | 78 +++++++++++++++++++ .../PostgresqlMessageStore.cs | 21 +++++ .../Wolverine.RDBMS/MessageDatabase.Admin.cs | 23 ++++++ 3 files changed, 122 insertions(+) create mode 100644 src/Persistence/PostgresqlTests/Transport/reset_clears_transport_queue_tables.cs diff --git a/src/Persistence/PostgresqlTests/Transport/reset_clears_transport_queue_tables.cs b/src/Persistence/PostgresqlTests/Transport/reset_clears_transport_queue_tables.cs new file mode 100644 index 000000000..2075bd623 --- /dev/null +++ b/src/Persistence/PostgresqlTests/Transport/reset_clears_transport_queue_tables.cs @@ -0,0 +1,78 @@ +using IntegrationTests; +using JasperFx.Core; +using Microsoft.Extensions.Hosting; +using Npgsql; +using Shouldly; +using Weasel.Postgresql; +using Wolverine; +using Wolverine.ComplianceTests; +using Wolverine.Persistence.Durability; +using Wolverine.Postgresql; +using Wolverine.Postgresql.Transport; +using Wolverine.Runtime; +using Wolverine.Tracking; + +namespace PostgresqlTests.Transport; + +public class reset_clears_transport_queue_tables : PostgresqlContext, IAsyncLifetime +{ + private IHost theHost = null!; + private PostgresqlQueue theQueue = null!; + private IMessageStore theMessageStore = null!; + + public async Task InitializeAsync() + { + using (var conn = new NpgsqlConnection(Servers.PostgresConnectionString)) + { + await conn.OpenAsync(); + await conn.DropSchemaAsync("reset_transports"); + await conn.CloseAsync(); + } + + theHost = await Host.CreateDefaultBuilder() + .UseWolverine(opts => + { + opts.UsePostgresqlPersistenceAndTransport(Servers.PostgresConnectionString, + schema: "reset_transports", transportSchema: "reset_transports"); + + // Neutralize the host's auto-started listener so it cannot drain the queue + // mid-test before the reset runs (mirrors the basic_functionality fixture). + opts.ListenToPostgresqlQueue("resetone").PollingInterval(1.Hours()); + }).StartAsync(); + + var transport = theHost.GetRuntime().Options.Transports.GetOrCreate(); + theQueue = transport.Queues["resetone"]; + theMessageStore = theHost.GetRuntime().Storage; + } + + public async Task DisposeAsync() + { + await theHost.StopAsync(); + theHost.Dispose(); + } + + [Fact] + public async Task clear_all_async_empties_the_transport_queue_tables() + { + // Row in the queue table + var immediate = ObjectMother.Envelope(); + immediate.DeliverBy = DateTimeOffset.UtcNow.AddHours(1); + await theQueue.SendAsync(immediate); + + // Row in the scheduled-message table + var scheduled = ObjectMother.Envelope(); + scheduled.ScheduleDelay = 1.Hours(); + scheduled.DeliverBy = DateTimeOffset.UtcNow.AddHours(1); + await theQueue.SendAsync(scheduled); + + (await theQueue.CountAsync()).ShouldBe(1); + (await theQueue.ScheduledCountAsync()).ShouldBe(1); + + // The message-store reset must also clear the registered transport queue tables, + // otherwise integration tests over the Postgres queue transport carry rows between runs. + await theMessageStore.Admin.ClearAllAsync(); + + (await theQueue.CountAsync()).ShouldBe(0); + (await theQueue.ScheduledCountAsync()).ShouldBe(0); + } +} diff --git a/src/Persistence/Wolverine.Postgresql/PostgresqlMessageStore.cs b/src/Persistence/Wolverine.Postgresql/PostgresqlMessageStore.cs index 3eb1e8e13..923076877 100644 --- a/src/Persistence/Wolverine.Postgresql/PostgresqlMessageStore.cs +++ b/src/Persistence/Wolverine.Postgresql/PostgresqlMessageStore.cs @@ -769,6 +769,27 @@ public void AddTable(Table table) _otherTables.Add(table); } + /// + /// Also clear the PostgreSQL queue transport's own tables (the per-queue message table and its + /// scheduled-message table) as part of a store reset. A reset truncates the envelope tables; the + /// queue transport keeps its own tables, so without this override a reset leaves queue rows behind + /// and integration tests over the Postgres queue transport carry rows between runs. Scoped to the + /// transport's own table types on purpose — AddTable is a general registration path, so we + /// clear only what the transport itself registered rather than every entry in _otherTables. + /// Runs inside the reset transaction (see the base truncateEnvelopeDataAsync) so the reset + /// stays atomic. + /// + protected override async Task truncateAdditionalTablesAsync(DbTransaction tx, CancellationToken token) + { + foreach (var table in _otherTables) + { + if (table is Transport.QueueTable or Transport.ScheduledMessageTable) + { + await tx.CreateCommand($"delete from {table.Identifier}").ExecuteNonQueryAsync(token); + } + } + } + public override DatabaseSagaSchema SagaSchemaFor() { if (_sagaStorage.TryFind(typeof(T), out var raw)) diff --git a/src/Persistence/Wolverine.RDBMS/MessageDatabase.Admin.cs b/src/Persistence/Wolverine.RDBMS/MessageDatabase.Admin.cs index 07d2b6e1f..df1837fa1 100644 --- a/src/Persistence/Wolverine.RDBMS/MessageDatabase.Admin.cs +++ b/src/Persistence/Wolverine.RDBMS/MessageDatabase.Admin.cs @@ -243,6 +243,14 @@ await tx.CreateCommand($"delete from {QuotedSchemaName}.{DatabaseConstants.Liste } } + // Let a provider clear its own extra tables inside the SAME reset transaction so + // nothing is left behind. This is deliberately a per-provider hook rather than a + // blanket loop over every table registered via AddTable: some providers register + // tables through that same path that a reset must preserve (e.g. SQL Server's + // rate-limit table). The PostgreSQL store overrides this to also empty the queue + // transport's queue + scheduled tables. + await truncateAdditionalTablesAsync(tx, _cancellation); + await tx.CommitAsync(_cancellation); await afterTruncateEnvelopeDataAsync(conn); @@ -262,4 +270,19 @@ protected virtual Task afterTruncateEnvelopeDataAsync(DbConnection conn) { return Task.CompletedTask; } + + /// + /// Hook to clear additional provider-specific tables within the reset transaction started by + /// / . The default is a no-op. Providers + /// whose transport registers extra tables that a reset should also empty override this — e.g. + /// the PostgreSQL queue transport registers its queue and scheduled-message tables on the store + /// and clears them here. This is deliberately a per-provider decision rather than a blanket loop + /// over every table registered via AddTable: some providers register tables through that + /// same path that a reset must keep intact (for example SQL Server's rate-limit table). Runs + /// inside the reset transaction so the whole reset stays atomic. + /// + protected virtual Task truncateAdditionalTablesAsync(DbTransaction tx, CancellationToken token) + { + return Task.CompletedTask; + } } \ No newline at end of file