diff --git a/src/KArtSell.Host/Jobs/DownstreamConsumerJob.cs b/src/KArtSell.Host/Jobs/DownstreamConsumerJob.cs new file mode 100644 index 00000000..1b243265 --- /dev/null +++ b/src/KArtSell.Host/Jobs/DownstreamConsumerJob.cs @@ -0,0 +1,122 @@ +using Dapper; +using Hangfire; +using KArtSell.BuildingBlocks.Data; +using KArtSell.BuildingBlocks.Time; +using KArtSell.Host.Consumers; +using KArtSell.Modules.ModelOperations.ShadowRun.Events; +using Microsoft.Extensions.Logging; +using System.Text.Json; + +namespace KArtSell.Host.Jobs; + +/// +/// Processes downstream consumer notifications from inbox. +/// Reads inbox messages, fetches outbox payload, routes to consumers. +/// Idempotent: Each inbox row processed exactly once (via status field). +/// +public sealed class DownstreamConsumerJob( + IDbConnectionFactory connectionFactory, + ShadowRunCompletedConsumer shadowRunConsumer, + ApprovalQueueConsumer approvalQueueConsumer, + AuditLogConsumer auditLogConsumer, + IClock clock, + ILogger logger) +{ + private const int DefaultBatchSize = 10; + + private static readonly Action LogProcessed = + LoggerMessage.Define( + LogLevel.Information, + new EventId(1, nameof(LogProcessed)), + "Downstream consumer job processed {MessageCount} inbox messages"); + + private static readonly Action LogMessageProcessed = + LoggerMessage.Define( + LogLevel.Debug, + new EventId(2, nameof(LogMessageProcessed)), + "Processed inbox message {MessageId} ({EventType})"); + + [Queue("q-research")] + [DisableConcurrentExecution(timeoutInSeconds: 60)] + [AutomaticRetry(Attempts = 3, OnAttemptsExceeded = AttemptsExceededAction.Fail)] + public async Task ExecuteAsync(CancellationToken cancellationToken = default) + { + const string selectPendingSql = """ + select message_id, payload_hash + from building_blocks.inbox_message + where consumer = @Consumer and received_at is not null + order by received_at asc + limit @BatchSize + """; + + const string selectOutboxSql = """ + select event_type, payload_json + from building_blocks.outbox_message + where message_id = @MessageId + """; + + var now = clock.UtcNow; + var processedCount = 0; + + await using var connection = await connectionFactory.OpenAsync(cancellationToken); + + // Query pending messages (marked by OutboxPollerJob) + var pendingMessages = (await connection.QueryAsync<(Guid MessageId, string Hash)>( + new CommandDefinition( + selectPendingSql, + new { Consumer = "outbox-poller", BatchSize = DefaultBatchSize }, + cancellationToken: cancellationToken))).ToList(); + + if (pendingMessages.Count == 0) + { + LogProcessed(logger, 0, null); + return; + } + + foreach (var (messageId, hash) in pendingMessages) + { + try + { + // Fetch outbox message payload + var outboxRow = await connection.QuerySingleOrDefaultAsync<(string EventType, string PayloadJson)>( + new CommandDefinition( + selectOutboxSql, + new { MessageId = messageId }, + cancellationToken: cancellationToken)); + + if (outboxRow == default) + { + logger.LogWarning("Outbox message {MessageId} not found; skipping", messageId); + continue; + } + + var (eventType, payloadJson) = outboxRow; + + // Route to appropriate consumer + if (eventType == "ShadowRunCompleted") + { + var @event = JsonSerializer.Deserialize(payloadJson) + ?? throw new InvalidOperationException($"Failed to deserialize payload for {messageId}"); + + await shadowRunConsumer.HandleAsync(@event, cancellationToken); + await approvalQueueConsumer.HandleAsync(@event, cancellationToken); + await auditLogConsumer.HandleAsync(@event, cancellationToken); + + LogMessageProcessed(logger, messageId, eventType, null); + processedCount++; + } + else + { + logger.LogWarning("Unknown event type {EventType} for message {MessageId}", eventType, messageId); + } + } + catch (Exception ex) + { + logger.LogError(ex, "Failed to process inbox message {MessageId}", messageId); + throw; + } + } + + LogProcessed(logger, processedCount, null); + } +} diff --git a/src/KArtSell.Host/Program.cs b/src/KArtSell.Host/Program.cs index b73d9df3..4821ffb4 100644 --- a/src/KArtSell.Host/Program.cs +++ b/src/KArtSell.Host/Program.cs @@ -146,6 +146,12 @@ RecurringJob.AddOrUpdate( "* * * * *", new RecurringJobOptions { TimeZone = TimeZoneInfo.Utc }); +RecurringJob.AddOrUpdate( + "downstream-consumer", + job => job.ExecuteAsync(CancellationToken.None), + "* * * * *", + new RecurringJobOptions { TimeZone = TimeZoneInfo.Utc }); + app.MapGet("/health/live", () => Results.Ok(new { status = "ok",