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