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