diff --git a/src/KArtSell.BuildingBlocks/Reliability/DapperOutboxMessageReader.cs b/src/KArtSell.BuildingBlocks/Reliability/DapperOutboxMessageReader.cs index 592cb54e..e4323610 100644 --- a/src/KArtSell.BuildingBlocks/Reliability/DapperOutboxMessageReader.cs +++ b/src/KArtSell.BuildingBlocks/Reliability/DapperOutboxMessageReader.cs @@ -17,7 +17,6 @@ public sealed class DapperOutboxMessageReader(IDbConnectionFactory connectionFac attempt as Attempt from building_blocks.outbox_message where published_at is null - and occurred_at >= @CutoffTime order by occurred_at asc limit @BatchSize """; @@ -36,7 +35,6 @@ public sealed class DapperOutboxMessageReader(IDbConnectionFactory connectionFac """; public async Task> GetUnpublishedAsync( - DateTimeOffset cutoffTime, int batchSize, CancellationToken cancellationToken) { @@ -44,7 +42,7 @@ public sealed class DapperOutboxMessageReader(IDbConnectionFactory connectionFac var messages = await connection.QueryAsync( new CommandDefinition( SelectUnpublishedSql, - new { CutoffTime = cutoffTime, BatchSize = batchSize }, + new { BatchSize = batchSize }, cancellationToken: cancellationToken)); return messages.ToList(); } diff --git a/src/KArtSell.Host/Jobs/OutboxPollerJob.cs b/src/KArtSell.Host/Jobs/OutboxPollerJob.cs index b626ed13..e98146ab 100644 --- a/src/KArtSell.Host/Jobs/OutboxPollerJob.cs +++ b/src/KArtSell.Host/Jobs/OutboxPollerJob.cs @@ -43,9 +43,10 @@ public sealed class OutboxPollerJob( public async Task ExecuteAsync(CancellationToken cancellationToken = default) { var now = clock.UtcNow; - var cutoffTime = now.AddMinutes(-5); - var messages = await reader.GetUnpublishedAsync(cutoffTime, DefaultBatchSize, cancellationToken); + // Process all unpublished messages ordered by occurred_at (oldest first). + // Monitoring: alert if any message pending > 5 min (see dashboard/alerts). + var messages = await reader.GetUnpublishedAsync(DefaultBatchSize, cancellationToken); var failureCount = 0; foreach (var message in messages) diff --git a/tests/KArtSell.Integration.Tests/OutboxPollerJobTests.cs b/tests/KArtSell.Integration.Tests/OutboxPollerJobTests.cs index f6a36126..f8b09e55 100644 --- a/tests/KArtSell.Integration.Tests/OutboxPollerJobTests.cs +++ b/tests/KArtSell.Integration.Tests/OutboxPollerJobTests.cs @@ -104,89 +104,6 @@ public sealed class OutboxPollerJobTests : IAsyncLifetime } } - [Fact] - public async Task ExecuteAsync_SkipsMessagesBeyondCutoff_OnlyProcessesRecentMessages() - { - // Arrange: Insert two messages - one old, one recent - var oldMessageId = Guid.NewGuid(); - var recentMessageId = Guid.NewGuid(); - var now = _clock.UtcNow; - const string payload = """{"data":"test"}"""; - var hash = ContentHasher.Sha256(payload); - - await using (var connection = await _connectionFactory.OpenAsync(CancellationToken.None)) - { - await connection.ExecuteAsync(""" - insert into building_blocks.outbox_message - (message_id, event_type, schema_version, payload_json, correlation_id, occurred_at, payload_hash) - values (@MessageId, @EventType, 1, cast(@PayloadJson as jsonb), @CorrelationId, @OccurredAt, @Hash) - """, - new - { - MessageId = oldMessageId, - EventType = "OldEvent", - PayloadJson = payload, - CorrelationId = Guid.NewGuid().ToString(), - OccurredAt = now.AddMinutes(-10), - Hash = hash - }); - - await connection.ExecuteAsync(""" - insert into building_blocks.outbox_message - (message_id, event_type, schema_version, payload_json, correlation_id, occurred_at, payload_hash) - values (@MessageId, @EventType, 1, cast(@PayloadJson as jsonb), @CorrelationId, @OccurredAt, @Hash) - """, - new - { - MessageId = recentMessageId, - EventType = "RecentEvent", - PayloadJson = payload, - CorrelationId = Guid.NewGuid().ToString(), - OccurredAt = now.AddMinutes(-2), - Hash = hash - }); - } - - var reader = new DapperOutboxMessageReader(_connectionFactory); - var job = new OutboxPollerJob(reader, _clock, _logger); - - // Act - await job.ExecuteAsync(CancellationToken.None); - - // Assert: Only recent message should be published (and have inbox entry) - await using (var connection = await _connectionFactory.OpenAsync(CancellationToken.None)) - { - var oldPublished = await connection.QuerySingleOrDefaultAsync( - "select published_at from building_blocks.outbox_message where message_id = @MessageId", - new { MessageId = oldMessageId }); - var recentPublished = await connection.QuerySingleOrDefaultAsync( - "select published_at from building_blocks.outbox_message where message_id = @MessageId", - new { MessageId = recentMessageId }); - - Assert.Null(oldPublished); - Assert.NotNull(recentPublished); - - // Verify inbox entry exists ONLY for recent message - var oldInbox = await connection.QuerySingleOrDefaultAsync( - "select message_id from building_blocks.inbox_message where message_id = @MessageId", - new { MessageId = oldMessageId }); - var recentInbox = await connection.QuerySingleOrDefaultAsync( - "select message_id from building_blocks.inbox_message where message_id = @MessageId", - new { MessageId = recentMessageId }); - - Assert.Null(oldInbox); - Assert.Equal(recentMessageId, recentInbox); - } - - // Cleanup - await using (var connection = await _connectionFactory.OpenAsync(CancellationToken.None)) - { - await connection.ExecuteAsync( - "delete from building_blocks.outbox_message where message_id in (@OldId, @RecentId)", - new { OldId = oldMessageId, RecentId = recentMessageId }); - } - } - [Fact] public async Task ExecuteAsync_SkipsMessagesExceedingMaxAttempts_LogsAsDeadLetter() {