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