using Dapper; using GB5Shared.ActionProcessor; using FrameworkDAL.CustomCode.ActionProcessor; using FrameworkDAL.CustomCode.DirectAction; using FrameworkDAL.CustomCode.EventSub; using GB5Shared.Connection; using GB5Shared.DTO.Framework.CommonConfig; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ServerConfig; using GB5Shared.DTO.PubSub; using GB5Shared.RuleEngine; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using Npgsql; using System; using System.Data; using System.Diagnostics; using System.Linq; using Microsoft.Data.SqlClient; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using static GB5Shared.GB5Constant.Constant; namespace FrameworkBLL.EventSub { public interface IEventSubBLL { /// /// Receives a published outbox event, looks up all MACTION rows for the /// EventTypeId, and queues one TEVENTACTIONRUN + TACTIONOUTBOX row per action. /// The ActionOutboxDispatcherQuartzJob will pick them up and execute /// (Email / SMS / Webhook / etc.) without this method doing the actual send. /// /// The outbox event message to process. /// /// Scheme + host of the Framework service, captured from the incoming HTTP request /// (e.g. "http://192.168.0.112:5000"). Used to build email approval button URLs. /// Derived dynamically — same pattern as WIP CALLBACKENDPOINT from login.RequestUrl. /// /// Cancellation token for async operation. Task ProcessEventAsync(OutboxEventMessage message, string frameworkBaseUrl = "", CancellationToken ct = default); } public sealed class EventSubBLL : IEventSubBLL { private const int Partitions = 5; private readonly IEventSubDAL _dal; private readonly IEventActionRunDAL _runDal; private readonly IActionOutboxDAL _outboxDal; private readonly IDirectActionConfigDAL _directActionConfigDAL; private readonly IApplicationConnection _appConnection; private readonly IOptionsMonitor _systemDto; private readonly IConditionEvaluator _conditionEvaluator; private readonly ILogger _logger; public EventSubBLL( IEventSubDAL dal, IEventActionRunDAL runDal, IActionOutboxDAL outboxDal, IDirectActionConfigDAL directActionConfigDAL, IApplicationConnection appConnection, IOptionsMonitor systemDto, IConditionEvaluator conditionEvaluator, ILogger logger) { _dal = dal; _runDal = runDal; _outboxDal = outboxDal; _directActionConfigDAL = directActionConfigDAL; _appConnection = appConnection; _systemDto = systemDto; _conditionEvaluator = conditionEvaluator; _logger = logger; } public async Task ProcessEventAsync(OutboxEventMessage message, string frameworkBaseUrl = "", CancellationToken ct = default) { // ── 1. Resolve tenant login ─────────────────────────────────────── // Prefer the ConnectionName carried in the message (set by both EventHandler.cs // direct publish and OutBox.PublishPendingEventsAsync). This is the exact MSERVER // connection the publisher used, so it is always correct — no MSERVERCONFIG lookup // needed, and no ambiguity from multiple STATUS=1 entries for the same ClientId. LoginDTO? login; if (!string.IsNullOrWhiteSpace(message.ConnectionName)) { login = new LoginDTO { UserId = message.UserId, ClientId = message.TenantId, ConnectionDatabaseName = message.ConnectionName, DatabaseName = message.ConnectionName, DatabaseType = (byte)_systemDto.CurrentValue.DataBaseType }; _logger.LogInformation( "EventSub: using ConnectionName from message | ConnectionName={Conn} TenantId={TenantId} EventTypeId={EventTypeId}", message.ConnectionName, message.TenantId, message.EventTypeId); } else { // Fallback for legacy messages that do not carry ConnectionName. _logger.LogWarning( "EventSub: ConnectionName missing in message — falling back to MSERVERCONFIG lookup | TenantId={TenantId} EventTypeId={EventTypeId}", message.TenantId, message.EventTypeId); login = await BuildLoginAsync(message.TenantId, message.UserId).ConfigureAwait(false); if (login is null) { _logger.LogWarning( "EventSub: tenant {TenantId} not found in MSERVERCONFIG — skipping EventTypeId={EventTypeId}", message.TenantId, message.EventTypeId); return; } } // ── 2. Load MACTION rows for this EventType ─────────────────────── var actions = await _dal.GetActionsByEventTypeAsync(message.EventTypeId, message.ObjectId, login, ct) .ConfigureAwait(false); if (actions.Count == 0) { Activity.Current?.SetTag("gb5.actions.count", 0); Activity.Current?.SetTag("gb5.actions.types", "(none)"); _logger.LogInformation( "EventSub: no active actions for EventTypeId={EventTypeId} TenantId={TenantId} — nothing queued", message.EventTypeId, message.TenantId); return; } // Enrich the parent subscriber span with what was found var actionTypeNames = actions.Select(a => a.ActionType switch { 0 => "Email", 1 => "SMS", 2 => "Webhook", 3 => "Notification", 4 => "Report", _ => $"Type{a.ActionType}" }); Activity.Current?.SetTag("gb5.actions.count", actions.Count); Activity.Current?.SetTag("gb5.actions.types", string.Join(",", actionTypeNames)); Activity.Current?.SetTag("gb5.actions.tenant.db", login.DatabaseName); Activity.Current?.SetTag("gb5.actions.ids", string.Join(",", actions.Select(a => a.ActionId))); _logger.LogInformation( "EventSub: {Count} action(s) found for EventTypeId={EventTypeId} TenantId={TenantId}", actions.Count, message.EventTypeId, message.TenantId); // ── 3. Queue one TACTIONOUTBOX row per action ───────────────────── foreach (var action in actions) { // ── Condition guard ──────────────────────────────────────────── // CONDITIONEXPRESSION is null/empty for existing actions → pass-through. // Evaluator logs and returns true on parse error (fail-open). if (!_conditionEvaluator.Evaluate(action.ConditionExpression, message)) { _logger.LogInformation( "EventSub: condition false — skipping ActionId={ActionId} EventTypeId={EventTypeId} Condition='{Condition}'", action.ActionId, message.EventTypeId, action.ConditionExpression); continue; } // Dynamic URL — captured from the live HTTP request by EventActionSubscribeController // (scheme + host, e.g. "http://192.168.0.112:5000"). Same pattern as WIP which // derives CALLBACKENDPOINT from login.RequestUrl. Falls back to appsettings only // if the controller did not forward the value (legacy path). var approvalBaseUrl = !string.IsNullOrWhiteSpace(frameworkBaseUrl) ? frameworkBaseUrl.TrimEnd('/') : (_systemDto.CurrentValue.ServiceBaseUrl?.TrimEnd('/') ?? string.Empty); var partition = Math.Abs(message.TenantId % Partitions); var destinationTopic = $"action-exec-p{partition}"; // ── Resolve DirectAction context fields ─────────────────────────── // If MACTION.DIRECTACTIONID > -1, load the MDIRECTACTION header to // obtain the field-name mapping, then extract ContextId and AssigneeUserId // from the outbox event payload using those configured field names. int directActionId = action.DirectActionId; int contextId = message.ObjectId; // sensible default: entity PK int assigneeUserId = -1; if (directActionId > -1) { try { var header = await _directActionConfigDAL .GetHeaderByIdAsync(directActionId, login, ct) .ConfigureAwait(false); if (header is not null && !string.IsNullOrWhiteSpace(message.Payload)) { using var doc = JsonDocument.Parse(message.Payload); var root = doc.RootElement; if (!string.IsNullOrWhiteSpace(header.ContextIdField) && root.TryGetProperty(header.ContextIdField, out var ctxEl) && ctxEl.TryGetInt32(out var ctxVal)) contextId = ctxVal; if (!string.IsNullOrWhiteSpace(header.AssigneeUserIdField) && root.TryGetProperty(header.AssigneeUserIdField, out var auidEl) && auidEl.TryGetInt32(out var auidVal)) assigneeUserId = auidVal; } } catch (Exception ex) { _logger.LogWarning(ex, "EventSub: failed to resolve DirectAction context | DirectActionId={Id} ActionId={ActionId}", directActionId, action.ActionId); } // Fallback: payload is an entity DTO (e.g. TLeaveDTO), not a WF task // notification — so ContextIdField was not found in the payload and // contextId is still the entity ObjectId. // Query TWORKFLOWTASK directly to get the pending task for this entity. // // The direct Dapr publish fires before the save transaction commits, so // TWORKFLOWTASK may not yet be visible on the first attempt. Retry with // increasing delays to let the commit propagate. if (contextId == message.ObjectId) { int[] retryDelaysMs = { 200, 500, 1000 }; foreach (var delayMs in retryDelaysMs) { try { await Task.Delay(delayMs, ct).ConfigureAwait(false); var wfTask = await _directActionConfigDAL .GetPendingWorkflowTaskAsync(message.ObjectId, login, ct) .ConfigureAwait(false); if (wfTask.HasValue) { contextId = wfTask.Value.WorkflowTaskId; if (assigneeUserId == -1) assigneeUserId = wfTask.Value.AssigneeUserId; _logger.LogInformation( "EventSub: TWORKFLOWTASK resolved after {Delay}ms | ObjectId={ObjectId} TaskId={TaskId} Assignee={Assignee}", delayMs, message.ObjectId, contextId, assigneeUserId); break; } _logger.LogWarning( "EventSub: TWORKFLOWTASK fallback — no pending task yet after {Delay}ms | ObjectId={ObjectId}", delayMs, message.ObjectId); } catch (OperationCanceledException) { throw; } catch (Exception ex) { _logger.LogWarning(ex, "EventSub: TWORKFLOWTASK fallback failed (delay={Delay}ms) | ObjectId={ObjectId}", delayMs, message.ObjectId); } } } } // ── Deduplication guard ─────────────────────────────────────────────── // The save pipeline publishes TWICE for every entity save: // 1. Direct Dapr publish (EventHandler.cs) — immediate, carries ContextJson // 2. TOUTBOX outbox relay (OutBox.PublishPendingEventsAsync) — reliability fallback // Both carry the same CorrelationKey. If we already inserted TEVENTACTIONRUN // for this (ActionId, CorrelationKey), skip the relay's arrival — no second email. if (!string.IsNullOrWhiteSpace(message.CorrelationKey)) { bool isDuplicate; try { isDuplicate = await _runDal.ExistsByCorrelationKeyAsync( action.ActionId, message.CorrelationKey, login, ct).ConfigureAwait(false); } catch (Exception exDup) { // If dedup check fails, proceed — better to send a duplicate than to miss isDuplicate = false; _logger.LogWarning(exDup, "EventSub: dedup check failed (non-fatal, will proceed) | ActionId={ActionId} CorrelationKey={CKey}", action.ActionId, message.CorrelationKey); } if (isDuplicate) { Activity.Current?.SetTag("gb5.event.duplicate", true); Activity.Current?.SetTag("gb5.event.skip_reason", "duplicate_correlation_key"); Activity.Current?.SetTag("gb5.event.correlation_key", message.CorrelationKey); Activity.Current?.SetTag("gb5.event.action_id", action.ActionId); _logger.LogInformation( "EventSub: ⚡ DUPLICATE skipped — already processed | ActionId={ActionId} CorrelationKey={CKey} EventTypeId={EvtId} TenantId={TenantId}", action.ActionId, message.CorrelationKey, message.EventTypeId, message.TenantId); continue; } } _logger.LogInformation( "EventSub: inserting TEVENTACTIONRUN | ActionId={ActionId} PayloadLength={Len}", action.ActionId, message.Payload?.Length ?? 0); int actionRunId; try { actionRunId = await _runDal.InsertAsync(new EventActionRunDTO { ActionId = action.ActionId, JobExecutionId = -1, EventTypeId = message.EventTypeId, TenantId = message.TenantId, Payload = message.Payload, CorrelationKey = message.CorrelationKey }, login, ct).ConfigureAwait(false); Activity.Current?.SetTag("gb5.run.id", actionRunId); _logger.LogInformation("EventSub: TEVENTACTIONRUN inserted | ActionRunId={RunId}", actionRunId); } catch (Exception ex) { Activity.Current?.SetStatus(ActivityStatusCode.Error, ex.Message); _logger.LogError(ex, "EventSub: TEVENTACTIONRUN INSERT FAILED | ActionId={ActionId} EventTypeId={EventTypeId} TenantId={TenantId} CorrelationKey={CorrelationKey} DB={DB}", action.ActionId, message.EventTypeId, message.TenantId, message.CorrelationKey, login.DatabaseName); throw; } var payloadJson = JsonSerializer.Serialize(new ActionEventDto { ActionRunId = actionRunId, ActionId = action.ActionId, ActionType = action.ActionType, TenantId = message.TenantId, DatabaseName = login.DatabaseName, SendTo = action.SendTo, ToDeliveryType = action.ToDeliveryType, ToContentType = action.ToContentType, TemplateId = action.TemplateId, MailCc = action.MailCc, MailCcDeliveryType = action.MailCcDeliveryType, MailBcc = action.MailBcc, MailBccDeliveryType = action.MailBccDeliveryType, ReplyTo = action.ReplyTo, WebServiceId = action.WebServiceId, UriParameterValue = action.UriParameterValue, ReportId = action.ReportId, CorrelationKey = message.CorrelationKey, // Captures whatever span is active right now — normally the Dapr-subscribe // span EventActionSubscribeController opened, itself parented from the // original publisher (e.g. SchedulerActionEventPublishHandler) — so // ActionProcessorWorker can continue the SAME trace after the RabbitMQ hop. TraceParent = System.Diagnostics.Activity.Current?.Id, TraceState = System.Diagnostics.Activity.Current?.TraceStateString, IncludeApprovalActions = action.ActionType == 0 && message.ObjectId > 0, WorkflowTaskId = message.ObjectId, ApprovalBaseUrl = approvalBaseUrl, DirectActionId = directActionId, ContextId = contextId, AssigneeUserId = assigneeUserId, ContextJson = message.ContextJson, Payload = JsonSerializer.SerializeToElement(new { EventTypeId = message.EventTypeId, ObjectTypeId = message.ObjectTypeId, ObjectId = message.ObjectId, EntityPayload = message.Payload }) }); _logger.LogInformation( "EventSub: inserting TACTIONOUTBOX | ActionRunId={RunId} Topic={Topic} PayloadLength={Len}", actionRunId, destinationTopic, payloadJson.Length); try { await _outboxDal.InsertAsync(new ActionOutboxDTO { ActionRunId = actionRunId, DestinationTopic = destinationTopic, Payload = payloadJson, TenantId = message.TenantId, CorrelationKey = message.CorrelationKey }, login, ct).ConfigureAwait(false); Activity.Current?.SetTag("gb5.outbox.topic", destinationTopic); _logger.LogInformation("EventSub: TACTIONOUTBOX inserted | ActionRunId={RunId} Topic={Topic}", actionRunId, destinationTopic); } catch (Exception ex) { Activity.Current?.SetStatus(ActivityStatusCode.Error, ex.Message); _logger.LogError(ex, "EventSub: TACTIONOUTBOX INSERT FAILED | ActionRunId={RunId} ActionId={ActionId} EventTypeId={EventTypeId} TenantId={TenantId} Topic={Topic} PayloadLength={Len} CorrelationKey={CorrelationKey}", actionRunId, action.ActionId, message.EventTypeId, message.TenantId, destinationTopic, payloadJson.Length, message.CorrelationKey); throw; } _logger.LogInformation( "EventSub: queued | EventTypeId={EventTypeId} ActionId={ActionId} ActionType={Type} ActionRunId={RunId} Topic={Topic}", message.EventTypeId, action.ActionId, action.ActionType, actionRunId, destinationTopic); } } // ── Helpers ─────────────────────────────────────────────────────────── private async Task BuildLoginAsync(int tenantId, int userId = -1) { var systemConn = await _appConnection.Gb5SystemConnectionString().ConfigureAwait(false); int dbType = _systemDto.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.CLIENTID = @TenantId"; using IDbConnection conn = dbType switch { DBTYPE.SQL => new SqlConnection(systemConn), DBTYPE.POSTGRESQL => new NpgsqlConnection(systemConn), _ => throw new NotSupportedException($"Unsupported DB type: {dbType}") }; var tenant = await conn.QueryFirstOrDefaultAsync( sql, new { TenantId = tenantId }).ConfigureAwait(false); if (tenant is null) return null; return new LoginDTO { UserId = userId, ClientId = tenant.ClientId, ConnectionDatabaseName = tenant.ConnectionName, DatabaseName = tenant.DatabaseName, DatabaseType = tenant.DbType }; } } }