using System.Data; using System.Data.Common; using Dapper; using KArtSell.BuildingBlocks.Data; using KArtSell.BuildingBlocks.Hashing; using KArtSell.BuildingBlocks.Time; using KArtSell.Host.Jobs; using Microsoft.Extensions.Logging; using Npgsql; using Xunit; namespace KArtSell.Integration.Tests; /// /// Outbox/Inbox Crash-Recovery & Audit Reconciliation Tests /// Covers: Process crashes, idempotency, consumer retries, audit trails /// Following AGENTS.md v16.0: Failure modes, recovery validation, evidence preservation /// public sealed class OutboxInboxCrashRecoveryTests : IAsyncLifetime { private readonly string _connectionString; private readonly NpgsqlDataSource _dataSource; private readonly IDbConnectionFactory _connectionFactory; private readonly IClock _clock = new SystemClock(); private readonly ILogger _logger = new NoOpLogger(); public OutboxInboxCrashRecoveryTests() { _connectionString = Environment.GetEnvironmentVariable("KARTSELL_POSTGRES") ?? "Host=localhost;Port=5432;Database=kartselldb;Username=kartsell;Password=kartsell4321@!"; _dataSource = new NpgsqlDataSourceBuilder(_connectionString).Build(); _connectionFactory = new NpgsqlConnectionFactory(_dataSource); } public async Task InitializeAsync() { // Verify database is accessible await using var connection = await _dataSource.OpenConnectionAsync(); await using var cmd = connection.CreateCommand(); cmd.CommandText = "SELECT 1"; await cmd.ExecuteScalarAsync(); } public async Task DisposeAsync() { await _dataSource.DisposeAsync(); } /// /// Gate 2.1: Outbox Durability - Messages survive process crash /// Scenario: Process crashes before marking message published. /// Recovery: Restart retrieves the unpublished message. /// [Fact] public async Task CrashRecovery_OutboxMessage_SurvivesProcessCrash() { // Arrange: Insert unpublished message simulating pre-crash state var messageId = Guid.NewGuid(); var now = _clock.UtcNow; const string payload = """{"decision_id":"crash-test"}"""; var hash = ContentHasher.Sha256(payload); var correlationId = Guid.NewGuid().ToString(); 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 = messageId, EventType = "CrashTestEvent", PayloadJson = payload, CorrelationId = correlationId, OccurredAt = now.AddMinutes(-2), Hash = hash }); } // Simulate process crash: Verify message is still unpublished await using (var connection = await _connectionFactory.OpenAsync(CancellationToken.None)) { var published = await connection.QuerySingleOrDefaultAsync( "SELECT published_at FROM building_blocks.outbox_message WHERE message_id = @MessageId", new { MessageId = messageId }); Assert.Null(published); } // Act: Restart process - query for unpublished messages (what job would retrieve) await using (var connection = await _connectionFactory.OpenAsync(CancellationToken.None)) { var unpublished = await connection.QueryAsync(""" SELECT message_id FROM building_blocks.outbox_message WHERE published_at IS NULL ORDER BY occurred_at ASC """); // Assert: Message is still retrievable for processing Assert.Contains(messageId, unpublished); } } /// /// Gate 2.2: Inbox Idempotency - Duplicate message handling via UNIQUE constraint /// Scenario: Same outbox message processed by multiple consumers (concurrent). /// Expected: UNIQUE(message_id, consumer) prevents duplicates. /// [Fact] public async Task CrashRecovery_InboxIdempotency_PreventsDuplicatesByConsumer() { // Arrange: Create outbox message var messageId = Guid.NewGuid(); var now = _clock.UtcNow; const string payload = """{"idempotency":"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, published_at) VALUES (@MessageId, @EventType, 1, CAST(@PayloadJson AS jsonb), @CorrelationId, @OccurredAt, @Hash, @Now) """, new { MessageId = messageId, EventType = "IdempotencyTestEvent", PayloadJson = payload, CorrelationId = Guid.NewGuid().ToString(), OccurredAt = now.AddMinutes(-2), Hash = hash, Now = now.DateTime }); // Insert first inbox record (simulating first consumer) await connection.ExecuteAsync(""" INSERT INTO building_blocks.inbox_message (message_id, consumer, status) VALUES (@MessageId, @Consumer1, 'Pending') """, new { MessageId = messageId, Consumer1 = "Consumer1" }); // Act: Try to insert duplicate (same message_id, same consumer) var exception = await Assert.ThrowsAsync(async () => { await connection.ExecuteAsync(""" INSERT INTO building_blocks.inbox_message (message_id, consumer, status) VALUES (@MessageId, @Consumer1, 'Pending') """, new { MessageId = messageId, Consumer1 = "Consumer1" }); }); // Assert: UNIQUE constraint violation (error code 23505 = unique_violation) Assert.Contains("23505", exception.SqlState); } /// /// Gate 2.3: Inbox Status Transitions - Invalid transitions blocked /// Scenario: Process marks message as Processed before processed_at is set. /// Expected: Trigger enforces processed_at must be set when status = Processed. /// [Fact] public async Task CrashRecovery_InboxStatus_EnforcesProcessedAtTimestamp() { // Arrange: Create outbox message var messageId = Guid.NewGuid(); var now = _clock.UtcNow; const string payload = """{"status":"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, published_at) VALUES (@MessageId, @EventType, 1, CAST(@PayloadJson AS jsonb), @CorrelationId, @OccurredAt, @Hash, @Now) """, new { MessageId = messageId, EventType = "StatusTestEvent", PayloadJson = payload, CorrelationId = Guid.NewGuid().ToString(), OccurredAt = now.AddMinutes(-2), Hash = hash, Now = now.DateTime }); // Insert inbox message in Pending state await connection.ExecuteAsync(""" INSERT INTO building_blocks.inbox_message (message_id, consumer, status) VALUES (@MessageId, @Consumer, 'Pending') """, new { MessageId = messageId, Consumer = "TestConsumer" }); // Act: Try to mark as Processed without setting processed_at var exception = await Assert.ThrowsAsync(async () => { await connection.ExecuteAsync(""" UPDATE building_blocks.inbox_message SET status = 'Processed' WHERE message_id = @MessageId """, new { MessageId = messageId }); }); // Assert: Trigger violation (processed_at must be set) Assert.Contains("processed_at must be set", exception.Message, StringComparison.OrdinalIgnoreCase); } /// /// Gate 2.4: Consumer Failure Retry - Failed messages are retrievable for retry /// Scenario: Consumer process crashes mid-processing, message marked Failed. /// Recovery: Retry job queries Failed messages and retries them. /// [Fact] public async Task CrashRecovery_ConsumerFailure_FailedMessagesRetrieval() { // Arrange: Create outbox message var messageId = Guid.NewGuid(); var now = _clock.UtcNow; const string payload = """{"retry":"test"}"""; var hash = ContentHasher.Sha256(payload); var errorMsg = "Consumer process crashed before completion"; 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, published_at) VALUES (@MessageId, @EventType, 1, CAST(@PayloadJson AS jsonb), @CorrelationId, @OccurredAt, @Hash, @Now) """, new { MessageId = messageId, EventType = "RetryTestEvent", PayloadJson = payload, CorrelationId = Guid.NewGuid().ToString(), OccurredAt = now.AddMinutes(-2), Hash = hash, Now = now.DateTime }); // Insert inbox message with Failed status await connection.ExecuteAsync(""" INSERT INTO building_blocks.inbox_message (message_id, consumer, status, error_message, attempted_at) VALUES (@MessageId, @Consumer, 'Failed', @ErrorMsg, @AttemptedAt) """, new { MessageId = messageId, Consumer = "RetryConsumer", ErrorMsg = errorMsg, AttemptedAt = now.DateTime }); // Act: Query for failed messages (to retry) var failedMessages = await connection.QueryAsync<(Guid MessageId, string Consumer, string ErrorMessage)>(""" SELECT message_id, consumer, error_message FROM building_blocks.inbox_message WHERE status = 'Failed' ORDER BY attempted_at DESC LIMIT 10 """); // Assert: Failed message is retrievable Assert.NotEmpty(failedMessages); var failed = failedMessages.First(); Assert.Equal(messageId, failed.MessageId); Assert.Equal("RetryConsumer", failed.Consumer); Assert.Contains("crashed", failed.ErrorMessage, StringComparison.OrdinalIgnoreCase); } /// /// Gate 2.5: Audit Reconciliation - All events are tracked with correlation ID /// Scenario: Process completes event handling. Audit log should have end-to-end trace. /// Expected: Every inbox message has correlation_id linking to outbox message. /// [Fact] public async Task CrashRecovery_AuditReconciliation_CorrelationIdTracing() { // Arrange: Create complete event flow var messageId = Guid.NewGuid(); var correlationId = Guid.NewGuid().ToString(); var now = _clock.UtcNow; const string payload = """{"audit":"trace"}"""; var hash = ContentHasher.Sha256(payload); await using var connection = await _connectionFactory.OpenAsync(CancellationToken.None); // Insert outbox message await connection.ExecuteAsync(""" INSERT INTO building_blocks.outbox_message (message_id, event_type, schema_version, payload_json, correlation_id, occurred_at, payload_hash, published_at) VALUES (@MessageId, @EventType, 1, CAST(@PayloadJson AS jsonb), @CorrelationId, @OccurredAt, @Hash, @Now) """, new { MessageId = messageId, EventType = "AuditTraceEvent", PayloadJson = payload, CorrelationId = correlationId, OccurredAt = now.AddMinutes(-2), Hash = hash, Now = now.DateTime }); // Insert inbox message (simulating successful processing) await connection.ExecuteAsync(""" INSERT INTO building_blocks.inbox_message (message_id, consumer, status, processed_at) VALUES (@MessageId, @Consumer, 'Processed', @Now) """, new { MessageId = messageId, Consumer = "AuditConsumer", Now = now.DateTime }); // Act: Reconcile - Verify inbox message links back to outbox via message_id var reconciliation = await connection.QuerySingleOrDefaultAsync<(string OutboxCorrelationId, string InboxConsumer, string InboxStatus)?>(""" SELECT o.correlation_id, i.consumer, i.status FROM building_blocks.outbox_message o INNER JOIN building_blocks.inbox_message i ON o.message_id = i.message_id WHERE o.message_id = @MessageId """, new { MessageId = messageId }); // Assert: End-to-end traceability Assert.NotNull(reconciliation); Assert.Equal(correlationId, reconciliation!.Value.OutboxCorrelationId); Assert.Equal("AuditConsumer", reconciliation!.Value.InboxConsumer); Assert.Equal("Processed", reconciliation!.Value.InboxStatus); } /// /// Gate 2.6: Multiple Consumers - Same outbox event routed to multiple consumers /// Scenario: One outbox message should create N inbox records (one per consumer). /// Expected: Each consumer processes independently (idempotent dedup per consumer). /// [Fact] public async Task CrashRecovery_MultipleConsumers_IndependentProcessing() { // Arrange: Create outbox message for multiple consumers var messageId = Guid.NewGuid(); var now = _clock.UtcNow; const string payload = """{"broadcast":"multi-consumer"}"""; 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, published_at) VALUES (@MessageId, @EventType, 1, CAST(@PayloadJson AS jsonb), @CorrelationId, @OccurredAt, @Hash, @Now) """, new { MessageId = messageId, EventType = "MultiConsumerEvent", PayloadJson = payload, CorrelationId = Guid.NewGuid().ToString(), OccurredAt = now.AddMinutes(-2), Hash = hash, Now = now.DateTime }); // Insert inbox records for three consumers var consumers = new[] { "Consumer1", "Consumer2", "Consumer3" }; foreach (var consumer in consumers) { await connection.ExecuteAsync(""" INSERT INTO building_blocks.inbox_message (message_id, consumer, status) VALUES (@MessageId, @Consumer, 'Pending') """, new { MessageId = messageId, Consumer = consumer }); } // Act: Query inbox records for this message var inboxRecords = await connection.QueryAsync<(string Consumer, string Status)>(""" SELECT consumer, status FROM building_blocks.inbox_message WHERE message_id = @MessageId ORDER BY consumer """, new { MessageId = messageId }); // Assert: All three consumers have independent records Assert.Equal(3, inboxRecords.Count()); Assert.Contains("Consumer1", inboxRecords.Select(r => r.Consumer)); Assert.Contains("Consumer2", inboxRecords.Select(r => r.Consumer)); Assert.Contains("Consumer3", inboxRecords.Select(r => r.Consumer)); Assert.All(inboxRecords, r => Assert.Equal("Pending", r.Status)); } // ========== Helper Classes ========== private sealed class NoOpLogger : ILogger { public IDisposable? BeginScope(TState state) where TState : notnull => null; public bool IsEnabled(LogLevel logLevel) => false; public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func formatter) { } } }