using FrameworkBLL.ActionProcessor;
using FrameworkBLL.ActionProcessor.Enricher;
using FrameworkBLL.ActionProcessor.Handlers;
using GB5Shared.ActionProcessor;
using FrameworkDAL.CustomCode.ActionProcessor;
using FrameworkDAL.DTO.SchedulerTaskGenerator;
using GB5Shared.Connection;
using GB5Shared.DTO.Framework.CommonConfig;
using GB5Shared.DTO.Framework.Login;
using GB5Shared.DTO.Framework.ServerConfig;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using Npgsql;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
using GB5Shared.Telemetry;
using Serilog;
using System;
using System.Collections.Generic;
using System.Data;
using System.Diagnostics;
using Microsoft.Data.SqlClient;
using System.Linq;
using System.Text;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using static GB5Shared.GB5Constant.Constant;
namespace FrameworkSL.Controllers.ActionProcessor
{
///
/// BackgroundService that:
/// 1. Subscribes to Scheduler.Ready — deserializes SchedulerTaskDTO,
/// generates TEVENTACTIONRUN + TACTIONOUTBOX per action via ISchedulerActionGeneratorBLL.
/// 2. Subscribes to action-exec-p0..p4 — deserializes ActionEventDto,
/// resolves and runs the appropriate IActionHandler, updates TEVENTACTIONRUN.
///
/// Idempotency: each ActionRunId is checked against LPROCESSEDACTION before handling.
/// Retry: TransientFailure NACKs the message (RabbitMQ requeues / moves to DLQ).
///
public sealed class ActionProcessorWorker : BackgroundService, IAsyncDisposable
{
private const int MaxParallel = 4;
private const int PartitionCount = 5;
private const int MaxHandlerAttempts = 3;
private const int RetryDelayMs = 10000; // backoff before a TransientFailure redelivers
private readonly IServiceScopeFactory _scopeFactory;
private readonly ILogger _logger;
private readonly IConfiguration _config;
private readonly FrameworkBLL.SchedulerTaskGenerator.IJobEngineHubNotifier _hub;
private readonly SemaphoreSlim _throttle = new(MaxParallel, MaxParallel);
private IConnection? _rabbitConnection;
private readonly List _channels = new();
public ActionProcessorWorker(
IServiceScopeFactory scopeFactory,
ILogger logger,
IConfiguration config,
FrameworkBLL.SchedulerTaskGenerator.IJobEngineHubNotifier hub)
{
_scopeFactory = scopeFactory;
_logger = logger;
_config = config;
_hub = hub;
}
private readonly TaskCompletionSource _reconnectSignal = new(TaskCreationOptions.RunContinuationsAsynchronously);
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
_logger.LogInformation("ActionProcessorWorker starting");
while (!stoppingToken.IsCancellationRequested)
{
await CleanupRabbitAsync().ConfigureAwait(false);
var delay = TimeSpan.FromSeconds(5);
while (!stoppingToken.IsCancellationRequested)
{
try
{
await InitRabbitAsync(stoppingToken).ConfigureAwait(false);
break;
}
catch (OperationCanceledException)
{
return;
}
catch (Exception ex)
{
_logger.LogWarning(ex,
"ActionProcessorWorker: RabbitMQ connect failed — retrying in {Delay}s",
delay.TotalSeconds);
try { await Task.Delay(delay, stoppingToken).ConfigureAwait(false); }
catch (OperationCanceledException) { return; }
delay = TimeSpan.FromSeconds(Math.Min(delay.TotalSeconds * 2, 60));
}
}
if (stoppingToken.IsCancellationRequested) break;
// Wait until connection drops or service stops
var reconnect = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
_rabbitConnection!.ConnectionShutdownAsync += (_, args) =>
{
_logger.LogWarning(
"ActionProcessorWorker: RabbitMQ connection lost — Initiator={I} ReplyCode={C} ReplyText={T}",
args.Initiator, args.ReplyCode, args.ReplyText);
reconnect.TrySetResult(true);
return Task.CompletedTask;
};
using var reg = stoppingToken.Register(() => reconnect.TrySetCanceled());
try
{
await reconnect.Task.ConfigureAwait(false);
}
catch (OperationCanceledException)
{
break;
}
_logger.LogInformation("ActionProcessorWorker: reconnecting to RabbitMQ...");
}
_logger.LogInformation("ActionProcessorWorker stopping");
}
private async Task CleanupRabbitAsync()
{
foreach (var ch in _channels)
{
try { await ch.CloseAsync().ConfigureAwait(false); } catch { }
try { ch.Dispose(); } catch { }
}
_channels.Clear();
if (_rabbitConnection is not null)
{
try { await _rabbitConnection.CloseAsync().ConfigureAwait(false); } catch { }
try { await _rabbitConnection.DisposeAsync().ConfigureAwait(false); } catch { }
_rabbitConnection = null;
}
}
// ── RabbitMQ setup ────────────────────────────────────────────────────
private async Task InitRabbitAsync(CancellationToken ct)
{
var factory = new ConnectionFactory
{
HostName = _config["RabbitMQ:Host"] ?? "localhost",
UserName = _config["RabbitMQ:UserName"] ?? "guest",
Password = _config["RabbitMQ:Password"] ?? "guest",
VirtualHost = _config["RabbitMQ:VirtualHost"] ?? "/",
AutomaticRecoveryEnabled = true
};
_logger.LogInformation("ActionProcessorWorker: connecting to RabbitMQ at {Host}", factory.HostName);
_rabbitConnection = await factory.CreateConnectionAsync(ct).ConfigureAwait(false);
_logger.LogInformation("ActionProcessorWorker: RabbitMQ connected");
// Scheduler.Ready — declare with same DLX args as RabbitMqPublisher (idempotent)
var dlxArgs = new Dictionary { ["x-dead-letter-exchange"] = "Scheduler.DLX" };
await SubscribeAsync("Scheduler.Ready", OnSchedulerReadyAsync, ct, dlxArgs).ConfigureAwait(false);
// action-exec-p0..p{N-1} — no DLX args on the live queue itself; each partition
// gets a matching "-retry" delay queue instead (declared below) so a
// TransientFailure backs off for RetryDelayMs instead of instantly redelivering.
for (int i = 0; i < PartitionCount; i++)
{
var queueName = $"action-exec-p{i}";
await SubscribeAsync(
queueName,
(body, tag, ch, cta) => OnActionExecAsync(body, tag, ch, cta, queueName),
ct).ConfigureAwait(false);
await DeclareRetryQueueAsync(queueName, ct).ConfigureAwait(false);
}
_logger.LogInformation(
"ActionProcessorWorker: listening on {Count} action-exec queues (action-exec-p0..p{Max}), " +
"each with a {DelayMs}ms delayed-retry queue for transient failures",
PartitionCount, PartitionCount - 1, RetryDelayMs);
}
///
/// Declares a "{targetQueue}-retry" queue with a fixed message TTL that dead-letters
/// back to targetQueue via the default exchange once the TTL expires. A TransientFailure
/// publishes here instead of NACK-requeuing the live queue directly — without this, a
/// flaky/rate-limited downstream endpoint gets hammered with near-instant redelivery
/// (confirmed live: ~90 requests to a test endpoint in under 2 minutes).
///
private async Task DeclareRetryQueueAsync(string targetQueue, CancellationToken ct)
{
var channel = await _rabbitConnection!.CreateChannelAsync(cancellationToken: ct)
.ConfigureAwait(false);
await channel.QueueDeclareAsync(
queue: $"{targetQueue}-retry", durable: true, exclusive: false, autoDelete: false,
arguments: new Dictionary
{
["x-message-ttl"] = RetryDelayMs,
["x-dead-letter-exchange"] = "",
["x-dead-letter-routing-key"] = targetQueue
},
cancellationToken: ct).ConfigureAwait(false);
_channels.Add(channel);
}
private async Task SubscribeAsync(
string queue,
Func, ulong, IChannel, CancellationToken, Task> handler,
CancellationToken ct,
IDictionary? declareArgs = null)
{
var channel = await _rabbitConnection!.CreateChannelAsync(cancellationToken: ct)
.ConfigureAwait(false);
// Active declare with correct args — idempotent, no race condition with RabbitMqPublisher
await channel.QueueDeclareAsync(
queue: queue, durable: true, exclusive: false, autoDelete: false,
arguments: declareArgs, cancellationToken: ct).ConfigureAwait(false);
await channel.BasicQosAsync(0, 1, false, ct).ConfigureAwait(false);
var consumer = new AsyncEventingBasicConsumer(channel);
consumer.ReceivedAsync += async (_, ea) =>
{
await _throttle.WaitAsync(ct).ConfigureAwait(false);
try
{
await handler(ea.Body, ea.DeliveryTag, channel, ct).ConfigureAwait(false);
}
finally
{
_throttle.Release();
}
};
await channel.BasicConsumeAsync(queue, autoAck: false, consumer, ct)
.ConfigureAwait(false);
_channels.Add(channel);
_logger.LogInformation("Subscribed to RabbitMQ queue {Queue}", queue);
}
// ── Scheduler.Ready handler ───────────────────────────────────────────
private async Task OnSchedulerReadyAsync(
ReadOnlyMemory body,
ulong deliveryTag,
IChannel channel,
CancellationToken ct)
{
SchedulerTaskDTO? task;
try
{
task = JsonSerializer.Deserialize(body.Span);
}
catch (Exception ex)
{
_logger.LogError(ex, "Scheduler.Ready: failed to deserialize message");
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
if (task is null)
{
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
var tenants = await GetTenantsAsync().ConfigureAwait(false);
var tenant = tenants.FirstOrDefault(t => t.ClientId == task.TenantId);
if (tenant is null)
{
_logger.LogWarning("Scheduler.Ready: tenant {Id} not found — ACKing", task.TenantId);
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
var login = BuildLogin(tenant);
using var scope = _scopeFactory.CreateScope();
var generator = scope.ServiceProvider.GetRequiredService();
var runDal = scope.ServiceProvider.GetRequiredService();
foreach (var action in task.Actions)
{
try
{
await generator.GenerateAsync(
action, (int)task.JobExecutionId, tenant.DatabaseName, login, ct)
.ConfigureAwait(false);
}
catch (Exception ex)
{
_logger.LogError(ex, "GenerateAsync failed | ActionId={Id}", action.ActionId);
}
}
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
}
// ── action-exec-p* handler ────────────────────────────────────────────
private async Task OnActionExecAsync(
ReadOnlyMemory body,
ulong deliveryTag,
IChannel channel,
CancellationToken ct,
string sourceQueue)
{
_logger.LogInformation("━━━ action-exec: MESSAGE RECEIVED | Bytes={Len} DeliveryTag={Tag}", body.Length, deliveryTag);
// ── STEP 1: Deserialize ───────────────────────────────────────────
ActionEventDto? dto;
try
{
var raw = System.Text.Encoding.UTF8.GetString(body.Span);
_logger.LogInformation("action-exec [1/8] RAW BODY: {Body}", raw);
// Handle double-serialized payload: dispatcher may publish the DTO as a
// JSON string (outer quotes), e.g. "{\u0022ActionRunId\u0022:6,...}"
// Detect by checking if the trimmed body starts with '"'
string jsonToParse = raw.TrimStart();
if (jsonToParse.StartsWith("\"", StringComparison.Ordinal))
{
_logger.LogInformation("action-exec [1/8] Detected double-serialized body — unwrapping outer JSON string");
jsonToParse = JsonSerializer.Deserialize(raw)!;
}
dto = JsonSerializer.Deserialize(jsonToParse);
}
catch (Exception ex)
{
_logger.LogError(ex, "action-exec [1/8] FAILED deserialize | Body={Body}",
System.Text.Encoding.UTF8.GetString(body.Span));
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
if (dto is null)
{
_logger.LogWarning("action-exec [1/8] deserialized to null — ACKing and skipping");
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
_logger.LogInformation(
"action-exec [1/8] DESERIALIZED OK | ActionRunId={RunId} ActionType={Type} TenantId={TenantId} SendTo={SendTo} MailCc={Cc}",
dto.ActionRunId, dto.ActionType, dto.TenantId, dto.SendTo ?? "(null)", dto.MailCc ?? "(null)");
// Reconstruct the parent trace from dto.TraceParent/TraceState (set by whichever
// producer built this DTO — EventSubBLL's Dapr-subscribe span, most commonly) so
// this handler's work, including the actual send below, lands in the SAME trace
// as the originating Quartz fire / Dapr publish, not a disconnected new one.
var actionExecParentContext =
dto.TraceParent is not null
&& ActivityContext.TryParse(dto.TraceParent, dto.TraceState, isRemote: true, out var actionExecCtx)
? actionExecCtx
: default;
using var actionExecActivity = actionExecParentContext != default
? GB5ActivitySources.ActionProcessor.StartActivity(
"action-exec", ActivityKind.Consumer, actionExecParentContext)
: GB5ActivitySources.ActionProcessor.StartActivity(
"action-exec", ActivityKind.Consumer);
actionExecActivity?.SetTag("gb5.action.run_id", dto.ActionRunId);
actionExecActivity?.SetTag("gb5.action.type", dto.ActionType);
actionExecActivity?.SetTag("gb5.tenant.id", dto.TenantId);
// ── EMAIL FAST-PATH (ActionType=0) ────────────────────────────────
// SMTP send itself is global config — no tenant needed for that part.
// ✅ Result write-back (TEVENTACTIONRUN + TJOBEXECUTION finalize) still needs
// the owning tenant's DB, so we resolve it here — previously this path
// ack/nack'd the RabbitMQ message without ever recording a result, leaving
// TEVENTACTIONRUN stuck at Pending and the parent TJOBEXECUTION at InProgress.
if (dto.ActionType == 0)
{
_logger.LogInformation("action-exec ⚡ EMAIL FAST-PATH | ActionRunId={Id} SendTo={To}", dto.ActionRunId, dto.SendTo);
using var emailScope = _scopeFactory.CreateScope();
var emailResolver = emailScope.ServiceProvider.GetRequiredService();
var emailHandler = emailResolver.Resolve(0);
if (emailHandler is null)
{
_logger.LogError("action-exec ❌ EmailActionHandler not registered");
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
var emailStarted = DateTime.UtcNow;
ActionResult emailResult;
try
{
emailResult = await emailHandler.HandleAsync(dto, ct).ConfigureAwait(false);
}
catch (Exception ex)
{
_logger.LogError(ex, "action-exec ❌ EmailActionHandler threw | ActionRunId={Id}", dto.ActionRunId);
await RecordEmailFastPathResultAsync(dto, ActionResult.Permanent(ex.Message), emailStarted, ct)
.ConfigureAwait(false);
await channel.BasicNackAsync(deliveryTag, false, requeue: false, ct).ConfigureAwait(false);
return;
}
await RecordEmailFastPathResultAsync(dto, emailResult, emailStarted, ct).ConfigureAwait(false);
if (emailResult.Outcome == ActionOutcome.Success)
{
_logger.LogInformation("action-exec ✅ EMAIL SENT via fast-path | ActionRunId={Id} To={To}", dto.ActionRunId, dto.SendTo);
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
}
else
{
_logger.LogError("action-exec ❌ EMAIL FAILED | Outcome={O} Error={E}", emailResult.Outcome, emailResult.Error);
await channel.BasicNackAsync(deliveryTag, false, requeue: false, ct).ConfigureAwait(false);
}
return;
}
// ─────────────────────────────────────────────────────────────────
// ── STEP 2: Tenant resolution ─────────────────────────────────────
_logger.LogInformation("action-exec [2/8] Loading tenants from MSERVERCONFIG...");
List tenants;
try
{
tenants = await GetTenantsAsync().ConfigureAwait(false);
}
catch (Exception ex)
{
_logger.LogError(ex, "action-exec [2/8] FAILED loading tenants — ACKing");
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
_logger.LogInformation(
"action-exec [2/8] Tenants found: {Count} | ClientIds=[{Ids}]",
tenants.Count,
string.Join(", ", tenants.Select(t => t.ClientId.ToString())));
var tenant = tenants.FirstOrDefault(t => t.ClientId == dto.TenantId);
if (tenant is null)
{
_logger.LogError(
"action-exec [2/8] ❌ TENANT NOT FOUND | Looking for ClientId={TenantId} | Available=[{Ids}] — ACKing",
dto.TenantId,
string.Join(", ", tenants.Select(t => t.ClientId.ToString())));
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
_logger.LogInformation(
"action-exec [2/8] ✅ Tenant resolved | ClientId={ClientId} DatabaseName={Db} ConnectionName={Conn}",
tenant.ClientId, tenant.DatabaseName, tenant.ConnectionName);
var login = BuildLogin(tenant);
// ── STEP 3: Resolve services ──────────────────────────────────────
_logger.LogInformation("action-exec [3/8] Resolving scoped services...");
using var scope = _scopeFactory.CreateScope();
var runDal = scope.ServiceProvider.GetRequiredService();
var processedDal = scope.ServiceProvider.GetRequiredService();
var handlerResolver = scope.ServiceProvider.GetRequiredService();
var enricher = scope.ServiceProvider.GetRequiredService();
_logger.LogInformation("action-exec [3/8] ✅ Services resolved");
// ── STEP 4: Idempotency check ─────────────────────────────────────
_logger.LogInformation("action-exec [4/8] Idempotency check | ActionRunId={Id}", dto.ActionRunId);
bool alreadyProcessed;
try
{
alreadyProcessed = await processedDal.ExistsAsync(dto.ActionRunId, login, ct).ConfigureAwait(false);
}
catch (Exception ex)
{
_logger.LogError(ex, "action-exec [4/8] FAILED idempotency check | ActionRunId={Id}", dto.ActionRunId);
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
if (alreadyProcessed)
{
_logger.LogWarning("action-exec [4/8] ⚠️ DUPLICATE — ActionRunId={Id} already processed, skipping", dto.ActionRunId);
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
_logger.LogInformation("action-exec [4/8] ✅ Not a duplicate, proceeding");
// ── STEP 5: Mark InProgress ───────────────────────────────────────
_logger.LogInformation("action-exec [5/8] Marking InProgress | ActionRunId={Id}", dto.ActionRunId);
try
{
await runDal.UpdateRunStatusAsync(dto.ActionRunId, 1, login, ct).ConfigureAwait(false);
_logger.LogInformation("action-exec [5/8] ✅ Marked InProgress");
}
catch (Exception ex)
{
_logger.LogWarning(ex, "action-exec [5/8] ⚠️ UpdateRunStatus failed (row may not exist) — continuing anyway | ActionRunId={Id}", dto.ActionRunId);
}
// ── STEP 6: Enrich payload ────────────────────────────────────────
_logger.LogInformation("action-exec [6/8] Enriching payload | ActionRunId={Id}", dto.ActionRunId);
try
{
dto = await enricher.EnrichAsync(dto).ConfigureAwait(false);
_logger.LogInformation("action-exec [6/8] ✅ Payload enriched | SendTo={SendTo}", dto.SendTo ?? "(null)");
}
catch (Exception ex)
{
_logger.LogError(ex, "action-exec [6/8] FAILED enrichment | ActionRunId={Id}", dto.ActionRunId);
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
// ── STEP 7: Resolve & execute handler ─────────────────────────────
_logger.LogInformation("action-exec [7/8] Resolving handler for ActionType={Type}", dto.ActionType);
var handler = handlerResolver.Resolve(dto.ActionType);
if (handler is null)
{
_logger.LogError(
"action-exec [7/8] ❌ NO HANDLER for ActionType={T} | ActionRunId={Id}",
dto.ActionType, dto.ActionRunId);
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
return;
}
_logger.LogInformation(
"action-exec [7/8] ✅ Handler resolved: {H} | Calling HandleAsync... ActionRunId={Id}",
handler.GetType().Name, dto.ActionRunId);
var started = DateTime.UtcNow;
ActionResult result;
try
{
result = await handler.HandleAsync(dto, ct).ConfigureAwait(false);
}
catch (Exception ex)
{
_logger.LogError(ex, "action-exec [7/8] ❌ Handler THREW exception | ActionRunId={Id}", dto.ActionRunId);
result = ActionResult.Transient(ex.Message);
}
var completed = DateTime.UtcNow;
_logger.LogInformation(
"action-exec [7/8] Handler returned | Outcome={Outcome} Response={Response} Error={Error} | ActionRunId={Id}",
result.Outcome, result.Response ?? "(null)", result.Error ?? "(null)", dto.ActionRunId);
// ── STEP 8: Finalize ──────────────────────────────────────────────
_logger.LogInformation("action-exec [8/8] Finalizing | Outcome={Outcome} ActionRunId={Id}", result.Outcome, dto.ActionRunId);
switch (result.Outcome)
{
case ActionOutcome.Success:
try { await runDal.UpdateResultAsync(dto.ActionRunId, 2, 1, null, started, completed, result.Response, login, ct).ConfigureAwait(false); }
catch (Exception ex) { _logger.LogWarning(ex, "action-exec [8/8] UpdateResult failed (non-critical) | ActionRunId={Id}", dto.ActionRunId); }
try { await processedDal.InsertAsync(dto.ActionRunId, login, ct).ConfigureAwait(false); }
catch (Exception ex) { _logger.LogWarning(ex, "action-exec [8/8] InsertProcessed failed (non-critical) | ActionRunId={Id}", dto.ActionRunId); }
await TryFinalizeJobExecutionAsync(dto.JobExecutionId, login, ct).ConfigureAwait(false);
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
_logger.LogInformation("action-exec [8/8] ✅ SUCCESS — email sent and ACKed | ActionRunId={Id}", dto.ActionRunId);
break;
case ActionOutcome.TransientFailure:
_logger.LogWarning(
"action-exec [8/8] ⚠️ TRANSIENT FAILURE — routing to {RetryQueue} for {DelayMs}ms backoff | ActionRunId={Id} Error={E}",
$"{sourceQueue}-retry", RetryDelayMs, dto.ActionRunId, result.Error);
try { await runDal.UpdateRunStatusAsync(dto.ActionRunId, 0, login, ct).ConfigureAwait(false); }
catch (Exception ex) { _logger.LogWarning(ex, "action-exec [8/8] UpdateRunStatus failed (non-critical)"); }
// Not terminal yet — will be retried after backoff, so TJOBEXECUTION is NOT finalized here.
try
{
var retryProps = new BasicProperties { Persistent = true };
await channel.BasicPublishAsync(
exchange: "", routingKey: $"{sourceQueue}-retry", mandatory: false,
basicProperties: retryProps, body: body, cancellationToken: ct)
.ConfigureAwait(false);
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
}
catch (Exception ex)
{
// Delay-queue publish failed — fall back to the old instant-requeue
// behavior rather than lose the message entirely.
_logger.LogError(ex,
"action-exec [8/8] Failed to route to delayed retry — falling back to instant NACK-requeue | ActionRunId={Id}",
dto.ActionRunId);
await channel.BasicNackAsync(deliveryTag, false, requeue: true, ct).ConfigureAwait(false);
}
break;
case ActionOutcome.PermanentFailure:
default:
_logger.LogError("action-exec [8/8] ❌ PERMANENT FAILURE — ACKing | ActionRunId={Id} Error={E}", dto.ActionRunId, result.Error);
try { await runDal.UpdateResultAsync(dto.ActionRunId, 3, 1, result.Error, started, completed, null, login, ct).ConfigureAwait(false); }
catch (Exception ex) { _logger.LogWarning(ex, "action-exec [8/8] UpdateResult failed (non-critical)"); }
await TryFinalizeJobExecutionAsync(dto.JobExecutionId, login, ct).ConfigureAwait(false);
await channel.BasicAckAsync(deliveryTag, false, ct).ConfigureAwait(false);
break;
}
_logger.LogInformation("━━━ action-exec: DONE | ActionRunId={Id}", dto.ActionRunId);
}
// ── Email fast-path result write-back ─────────────────────────────────
///
/// Records the fast-path email result to TEVENTACTIONRUN and finalizes the parent
/// TJOBEXECUTION if this was the last outstanding action. Resolves the tenant purely
/// for this DB write — the SMTP send itself already completed without it.
/// Best-effort: failures here are logged, never thrown (must not affect ack/nack).
///
private async Task RecordEmailFastPathResultAsync(
ActionEventDto dto, ActionResult result, DateTime startedOn, CancellationToken ct)
{
try
{
var tenants = await GetTenantsAsync().ConfigureAwait(false);
var tenant = tenants.FirstOrDefault(t => t.ClientId == dto.TenantId);
if (tenant is null)
{
_logger.LogWarning(
"RecordEmailFastPathResult: tenant {Id} not found — skipping write-back | ActionRunId={Id}",
dto.TenantId, dto.ActionRunId);
return;
}
var login = BuildLogin(tenant);
using var scope = _scopeFactory.CreateScope();
var runDal = scope.ServiceProvider.GetRequiredService();
var completed = DateTime.UtcNow;
if (result.Outcome == ActionOutcome.Success)
await runDal.UpdateResultAsync(dto.ActionRunId, 2, 1, null, startedOn, completed, result.Response, login, ct)
.ConfigureAwait(false);
else
await runDal.UpdateResultAsync(dto.ActionRunId, 3, 1, result.Error, startedOn, completed, null, login, ct)
.ConfigureAwait(false);
await TryFinalizeJobExecutionAsync(dto.JobExecutionId, login, ct).ConfigureAwait(false);
}
catch (Exception ex)
{
_logger.LogWarning(ex,
"RecordEmailFastPathResultAsync failed (non-critical) | ActionRunId={Id}", dto.ActionRunId);
}
}
// ── TJOBEXECUTION finalize ──────────────────────────────────────────────
///
/// Once every fanned-out action for a scheduler JobExecutionId has reached a terminal
/// state, writes the aggregate Success/Failed result back to TJOBEXECUTION via
/// ISchedulerTaskServiceBLL — closing the gap where a completed job never left
/// STATUS=InProgress(2). Best-effort: exceptions are logged, never thrown, since this
/// runs after the action's own result is already durably recorded — failing here must
/// not cause the action message to be redelivered and reprocessed.
///
private async Task TryFinalizeJobExecutionAsync(int jobExecutionId, LoginDTO login, CancellationToken ct)
{
if (jobExecutionId <= 0) return; // Not scheduler-originated — nothing to finalize.
try
{
using var scope = _scopeFactory.CreateScope();
var runDal = scope.ServiceProvider.GetRequiredService();
var progress = await runDal.GetJobExecutionProgressAsync(jobExecutionId, login, ct)
.ConfigureAwait(false);
if (!progress.IsComplete) return;
var schedulerService = scope.ServiceProvider
.GetRequiredService();
var isSuccess = progress.FailedActions == 0;
var message = isSuccess
? $"{progress.TotalActions} action(s) completed successfully"
: $"{progress.FailedActions}/{progress.TotalActions} action(s) failed: {progress.LastError}";
await schedulerService.UpdateExecutionResultAsync(jobExecutionId, isSuccess, message, login, ct)
.ConfigureAwait(false);
_logger.LogInformation(
"TJOBEXECUTION finalized | JobExecutionId={Id} Success={Success} Total={Total} Failed={Failed}",
jobExecutionId, isSuccess, progress.TotalActions, progress.FailedActions);
// Push the terminal SignalR event here — this is the only place that knows
// every fanned-out action for this JobExecutionId has reached a terminal
// state. The Quartz executor (SchedulerJobExecutorQuartzJob) only ever pushes
// ReceiveJobStarted/ReceiveJobDispatched — it has no visibility into this
// async downstream completion.
try
{
var taskDal = scope.ServiceProvider
.GetRequiredService();
var jobDetails = await taskDal.LoadJobDetailsAsync(jobExecutionId, login, ct).ConfigureAwait(false);
if (isSuccess)
{
await _hub.PushJobCompleted(login.ClientId, new GB5Shared.DTO.JobEngine.JobCompletedDTO
{
JobId = jobDetails?.JobId ?? 0,
JobExecutionId = jobExecutionId,
CorrelationId = jobDetails?.CorrelationId,
SuccessCount = progress.CompletedActions - progress.FailedActions,
TotalCount = progress.TotalActions,
CompletedAtUtc = DateTime.UtcNow
}, ct).ConfigureAwait(false);
}
else
{
await _hub.PushJobFailed(login.ClientId, new GB5Shared.DTO.JobEngine.JobFailedDTO
{
JobId = jobDetails?.JobId ?? 0,
JobExecutionId = jobExecutionId,
CorrelationId = jobDetails?.CorrelationId,
ErrorMessage = progress.LastError ?? message,
FailureCount = progress.FailedActions,
TotalCount = progress.TotalActions,
FailedAtUtc = DateTime.UtcNow
}, ct).ConfigureAwait(false);
}
}
catch (Exception hubEx)
{
_logger.LogWarning(hubEx,
"TryFinalizeJobExecutionAsync: hub push failed (non-critical) | JobExecutionId={Id}", jobExecutionId);
}
}
catch (Exception ex)
{
_logger.LogWarning(ex,
"TryFinalizeJobExecutionAsync failed (non-critical) | JobExecutionId={Id}", jobExecutionId);
}
}
// ── Helpers ───────────────────────────────────────────────────────────
private async Task> GetTenantsAsync()
{
using var scope = _scopeFactory.CreateScope();
var appConnection = scope.ServiceProvider.GetRequiredService();
var databaseDTO = scope.ServiceProvider.GetRequiredService>();
var systemConn = await appConnection.Gb5SystemConnectionString().ConfigureAwait(false);
int dbType = databaseDTO.CurrentValue.DataBaseType;
const string sql = @"
SELECT
SERVERCONFIG1.CLIENTID AS ClientId,
SERVERCONFIG1.DATABASENAME AS DatabaseName,
SERVERCONFIG1.DATABASETYPE AS DbType,
SERVERCONFIG1.CONNECTIONNAME AS ConnectionName
FROM MSERVERCONFIG SERVERCONFIG1
JOIN MSERVER SERVER1
ON SERVERCONFIG1.SERVERID = SERVER1.SERVERID
WHERE SERVERCONFIG1.STATUS = 1
AND SERVERCONFIG1.CONNECTIONNAME <> 'ACTIVITI'";
using IDbConnection connection = dbType switch
{
DBTYPE.SQL => new SqlConnection(systemConn),
DBTYPE.POSTGRESQL => new NpgsqlConnection(systemConn),
_ => throw new NotSupportedException($"Unsupported DB type: {dbType}")
};
return (await Dapper.SqlMapper.QueryAsync(connection, sql)
.ConfigureAwait(false)).ToList();
}
// ✅ ApplicationConnection.DBConnectionStringCached resolves connection strings by
// treating LoginDTO.DatabaseName as a lookup key against MSERVERCONFIG.CONNECTIONNAME
// (confusingly named, but confirmed — the same pattern required fixing
// FrameworkSL/SchedulerRun.json's DatabaseName from "unisoftgb4" to "UNISOFTGB5").
// Using tenant.DatabaseName (the literal DB name) here meant every DB call made with
// this LoginDTO — TEVENTACTIONRUN inserts, TJOBEXECUTION updates, everything in
// OnSchedulerReadyAsync/OnActionExecAsync — failed with "No server configuration
// found for connection name: " for every tenant, not just this one.
private static LoginDTO BuildLogin(ServerConfigDTO tenant) => new()
{
UserId = -1,
ClientId = tenant.ClientId,
ConnectionDatabaseName = tenant.ConnectionName,
DatabaseName = tenant.ConnectionName,
DatabaseType = tenant.DbType
};
public async ValueTask DisposeAsync()
{
await CleanupRabbitAsync().ConfigureAwait(false);
base.Dispose();
GC.SuppressFinalize(this);
}
}
}