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);
}
}