using AutomationBLL.DocumentFlow; using AutomationBLL.EventMap; using AutomationBLL.Run; using AutomationBLL.ScriptVersion; using AutomationBLL.Target; using GB5Shared.DTO.Framework.Login; using GB5Shared.Telemetry; using Microsoft.Extensions.Logging; namespace AutomationBLL.EventDispatch; // Consumes inbound business events (ORDER_ACCEPTED/ASN_CREATED/INVOICE_POSTED, etc. — see // AutomationSL/Subscriptions/AutomationBusinessEventSubscriber.cs) and turns them into a dispatched // run, applying the sequencing/idempotency guards the spec calls for in one place. public class EventDispatchBLL : IEventDispatchBLL { private readonly IEventMapBLL _eventMapBll; private readonly IDocumentFlowBLL _documentFlowBll; private readonly ITargetBLL _targetBll; private readonly IScriptVersionBLL _versionBll; private readonly IRunBLL _runBll; private readonly ILogger _logger; public EventDispatchBLL( IEventMapBLL eventMapBll, IDocumentFlowBLL documentFlowBll, ITargetBLL targetBll, IScriptVersionBLL versionBll, IRunBLL runBll, ILogger logger) { _eventMapBll = eventMapBll; _documentFlowBll = documentFlowBll; _targetBll = targetBll; _versionBll = versionBll; _runBll = runBll; _logger = logger; } public async Task HandleInboundEvent(string eventType, string orderId, int? customerId, string? payloadJson, LoginDTO login, CancellationToken ct) { GB5Trace.Step("handle-inbound-event", new { eventType, orderId, customerId }); var mapping = await _eventMapBll.Resolve(eventType, customerId, login, ct).ConfigureAwait(false); if (mapping is null) { _logger.LogInformation("No TAUTOMATIONEVENTMAP configured for EventType {EventType} TenantId {TenantId} — skipping", eventType, login.ClientId); return $"Skipped: no mapping configured for event type {eventType}"; } if (mapping.PrecedingStage is not null) { var preceding = await _documentFlowBll.GetByOrderStage(orderId, mapping.PrecedingStage, login, ct).ConfigureAwait(false); if (preceding is null || preceding.Status != "Success") { _logger.LogInformation( "Sequencing guard: OrderId {OrderId} stage {PrecedingStage} is not Success yet (current: {Status}) — deferring {Stage}", orderId, mapping.PrecedingStage, preceding?.Status ?? "(not started)", mapping.Stage); return $"Skipped: waiting for preceding stage '{mapping.PrecedingStage}' to succeed"; } } // Idempotency guard: a Pending or Success row for this (OrderId, Stage) means this event // was already dispatched or already completed — only a prior Failed row is retried. var existing = await _documentFlowBll.GetByOrderStage(orderId, mapping.Stage, login, ct).ConfigureAwait(false); if (existing is not null && existing.Status != "Failed") { _logger.LogInformation( "Idempotency guard: OrderId {OrderId} stage {Stage} already has status {Status} (RunId {RunId}) — skipping duplicate dispatch", orderId, mapping.Stage, existing.Status, existing.RunId); return $"Skipped: stage '{mapping.Stage}' already {existing.Status} (RunId {existing.RunId})"; } var target = await _targetBll.GetById(mapping.TargetId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"TAUTOMATIONEVENTMAP {mapping.MapId} points at TargetId {mapping.TargetId}, which no longer exists."); var version = await _versionBll.GetLatestReleased(target.ScriptId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"No released version found for ScriptId {target.ScriptId} — cannot dispatch event {eventType}."); long runId = await _runBll.EnqueueRunRaw(target.ScriptId, version.VersionId, $"Event:{eventType}", payloadJson, mapping.TargetId, login, ct) .ConfigureAwait(false); await _documentFlowBll.MarkPending(orderId, mapping.Stage, runId, mapping.AckCallbackUrl, login, ct).ConfigureAwait(false); return $"Dispatched RunId {runId} for OrderId {orderId} stage {mapping.Stage}"; } }