using Dapper; using KArtSell.BuildingBlocks.Data; using Microsoft.Extensions.Logging; using Npgsql; using System.Text.Json; namespace KArtSell.Host.Jobs; /// /// Handles consumer errors: logs to dead-letter queue, tracks failure metrics. /// Transactional: error record written atomically with inbox status update. /// Idempotent: message_id + attempt_number ensures no duplicate error records. /// public sealed class ConsumerErrorHandler( IDbConnectionFactory connectionFactory, ILogger logger) { private const int MaxRetryAttempts = 3; public async Task HandleConsumerErrorAsync( Guid messageId, string eventType, string payloadJson, string correlationId, Exception exception, int attemptNumber, CancellationToken cancellationToken = default) { var conn = await connectionFactory.OpenAsync(cancellationToken) as NpgsqlConnection ?? throw new InvalidOperationException("Failed to open connection"); await using (conn) { await using var transaction = await conn.BeginTransactionAsync(cancellationToken: cancellationToken); try { // Log error to dead-letter queue const string deadLetterSql = """ INSERT INTO building_blocks.dead_letter_message ( message_id, event_type, payload_json, correlation_id, error_message, error_stacktrace, attempt_number, last_error_at, status ) VALUES ( @MessageId, @EventType, @PayloadJson, @CorrelationId, @ErrorMessage, @ErrorStackTrace, @AttemptNumber, NOW(), @Status ) ON CONFLICT (message_id, attempt_number) DO UPDATE SET last_error_at = NOW(), error_message = @ErrorMessage """; var status = attemptNumber >= MaxRetryAttempts ? "FAILED" : "RETRYING"; await conn.ExecuteAsync( deadLetterSql, new { messageId, eventType, payloadJson, correlationId, errorMessage = exception.Message, errorStackTrace = exception.StackTrace ?? string.Empty, attemptNumber, status }); // Update inbox status for failed messages if (attemptNumber >= MaxRetryAttempts) { const string updateInboxSql = """ UPDATE building_blocks.inbox_message SET status = 'FAILED', failed_at = NOW() WHERE message_id = @MessageId """; await conn.ExecuteAsync(updateInboxSql, new { messageId }); } await transaction.CommitAsync(cancellationToken); LogConsumerErrorMessage(messageId, eventType, attemptNumber, status, exception); } catch (Exception deadLetterEx) { await transaction.RollbackAsync(cancellationToken); logger.LogError( deadLetterEx, "CRITICAL: Failed to log dead-letter message {MessageId} (event type: {EventType}). " + "Original error: {OriginalError}", messageId, eventType, exception.Message); throw; } } } public async Task IsMessageFailedAsync(Guid messageId, CancellationToken cancellationToken = default) { var conn = await connectionFactory.OpenAsync(cancellationToken) as NpgsqlConnection ?? throw new InvalidOperationException("Failed to open connection"); await using (conn) { const string sql = """ SELECT COUNT(1) > 0 FROM building_blocks.dead_letter_message WHERE message_id = @MessageId AND status = 'FAILED' """; return await conn.QuerySingleAsync(sql, new { messageId }); } } private static readonly Action LogConsumerError = LoggerMessage.Define( LogLevel.Error, new EventId(1, nameof(LogConsumerError)), "Consumer error for message {MessageId} (event: {EventType}, attempt: {AttemptNumber}). Status: {Status}"); private void LogConsumerErrorMessage(Guid messageId, string eventType, int attemptNumber, string status, Exception ex) { LogConsumerError(logger, messageId, eventType, attemptNumber, status, ex); } }