using GB5Shared.ActionProcessor; using FrameworkDAL.CustomCode.ActionProcessor; using FrameworkDAL.DTO.SchedulerTaskGenerator; using GB5Shared.DTO.Framework.Login; using Microsoft.Extensions.Logging; using System; using System.Diagnostics; using System.Text.Json; using System.Threading; using System.Threading.Tasks; namespace FrameworkBLL.ActionProcessor { /// /// Creates TEVENTACTIONRUN + TACTIONOUTBOX rows for each MACTION /// associated with a completed scheduler task. /// Called by ActionProcessorWorker after consuming from Scheduler.Ready. /// public interface ISchedulerActionGeneratorBLL { Task GenerateAsync( SchedulerActionDTO action, int jobExecutionId, string databaseName, LoginDTO login, CancellationToken ct = default); } public sealed class SchedulerActionGeneratorBLL : ISchedulerActionGeneratorBLL { private readonly IEventActionRunDAL _runDal; private readonly IActionOutboxDAL _outboxDal; private readonly ILogger _logger; // Partition count — messages spread across action-exec-p0..p{N-1} private const int Partitions = 5; public SchedulerActionGeneratorBLL( IEventActionRunDAL runDal, IActionOutboxDAL outboxDal, ILogger logger) { _runDal = runDal; _outboxDal = outboxDal; _logger = logger; } public async Task GenerateAsync( SchedulerActionDTO action, int jobExecutionId, string databaseName, LoginDTO login, CancellationToken ct = default) { // ── 1. Compute partition queue ─────────────────────────── // ✅ GB5 tenant/client IDs are commonly negative (e.g. -1399999958). // C#'s % keeps the dividend's sign, so a naive `TenantId % Partitions` // can yield a negative partition (e.g. -3), producing a queue name // ("action-exec-p-3") that ActionProcessorWorker never subscribes to // (it only listens on p0..p{Partitions-1}) — actions would silently // vanish into an unconsumed queue. Normalize to a non-negative index. var partition = ((action.TenantId % Partitions) + Partitions) % Partitions; var destinationTopic = $"action-exec-p{partition}"; // ── 2. Insert TEVENTACTIONRUN — IDENTITY generates ActionRunId ── int actionRunId; try { actionRunId = await _runDal.InsertAsync(new EventActionRunDTO { ActionId = action.ActionId, JobExecutionId = jobExecutionId, EventTypeId = -1, TenantId = action.TenantId }, login, ct).ConfigureAwait(false); Activity.Current?.SetTag("gb5.run.id", actionRunId); } catch (Exception ex) { Activity.Current?.SetStatus(ActivityStatusCode.Error, ex.Message); _logger.LogError(ex, "SchedulerActionGenerator: TEVENTACTIONRUN INSERT FAILED | ActionId={ActionId} JobExecutionId={JobExecId} TenantId={TenantId}", action.ActionId, jobExecutionId, action.TenantId); throw; } // ── 3. Build ActionEventDto with the real ActionRunId ──── var dto = new ActionEventDto { ActionRunId = actionRunId, ActionId = action.ActionId, ActionType = action.ActionType, TenantId = action.TenantId, DatabaseName = databaseName, // ✅ login.DatabaseName holds the tenant's CONNECTION NAME (set correctly by // ActionProcessorWorker.BuildLogin), not the literal database name above — // handlers that need to build their own fresh LoginDTO for DB access (e.g. // WebhookActionHandler → IWebhookEndpointResolver) need this, not DatabaseName. ConnectionName = login.DatabaseName, JobExecutionId = jobExecutionId, SendTo = action.SendTo, ReplyTo = action.ReplyTo, TemplateId = action.TemplateId, MailCc = action.MailCc, MailBcc = action.MailBcc, WebServiceId = action.WebServiceId, UriParameterValue = action.UriParameterValue, ReportId = action.ReportId, Payload = JsonSerializer.SerializeToElement(new { jobExecutionId, action.ActionId }) }; var payloadJson = JsonSerializer.Serialize(dto); // ── 4. Insert TACTIONOUTBOX (SENDSTATUS=0 Pending) ─────── try { await _outboxDal.InsertAsync(new ActionOutboxDTO { ActionRunId = actionRunId, DestinationTopic = destinationTopic, Payload = payloadJson, TenantId = action.TenantId }, login, ct).ConfigureAwait(false); Activity.Current?.SetTag("gb5.outbox.topic", destinationTopic); } catch (Exception ex) { Activity.Current?.SetStatus(ActivityStatusCode.Error, ex.Message); _logger.LogError(ex, "SchedulerActionGenerator: TACTIONOUTBOX INSERT FAILED | ActionRunId={RunId} ActionId={ActionId} JobExecutionId={JobExecId} TenantId={TenantId} Topic={Topic}", actionRunId, action.ActionId, jobExecutionId, action.TenantId, destinationTopic); throw; } _logger.LogInformation( "Action queued | ActionRunId={Id} ActionType={Type} Topic={Topic}", actionRunId, action.ActionType, destinationTopic); } } }