using GB5Shared.ActionProcessor; using FrameworkDAL.CustomCode.Playground; using FrameworkDAL.DTO.Playground; using FrameworkDAL.DTO.SchedulerTaskGenerator; using GB5Shared.DTO.Framework.Login; using GB5Shared.Resource.Response; using Microsoft.Extensions.Logging; using System; using System.Diagnostics; using System.Linq; using System.Text.Json; using System.Threading; using System.Threading.Tasks; namespace FrameworkBLL.Playground { public sealed class PlaygroundBLL : IPlaygroundBLL { private readonly IPlaygroundDAL _dal; private readonly IEventActionRunDAL _runDal; private readonly IActionOutboxDAL _outboxDal; private readonly ILogger _logger; private const int Partitions = 5; public PlaygroundBLL( IPlaygroundDAL dal, IEventActionRunDAL runDal, IActionOutboxDAL outboxDal, ILogger logger) { _dal = dal; _runDal = runDal; _outboxDal = outboxDal; _logger = logger; } public async Task TriggerEventAsync( TriggerEventDTO dto, LoginDTO login, CancellationToken ct = default) { var playgroundRunId = Guid.NewGuid().ToString("N")[..16]; var actionsToQueue = dto.MockMode ? BuildMockActions(dto, login) : await _dal.GetActionsByEventTypeAsync(dto.EventTypeId, login, ct).ConfigureAwait(false); var limited = actionsToQueue.Take(dto.RecordLimit).ToList(); foreach (var action in limited) { var partition = login.ClientId % Partitions; var destinationTopic = $"action-exec-p{partition}"; // Insert run first — IDENTITY generates ActionRunId int actionRunId; try { actionRunId = await _runDal.InsertAsync(new EventActionRunDTO { ActionId = action.ActionId, JobExecutionId = -1, EventTypeId = dto.MockMode ? -1 : dto.EventTypeId, TenantId = login.ClientId, PlaygroundRunId = playgroundRunId }, login, ct).ConfigureAwait(false); Activity.Current?.SetTag("gb5.run.id", actionRunId); } catch (Exception ex) { Activity.Current?.SetStatus(ActivityStatusCode.Error, ex.Message); _logger.LogError(ex, "Playground: TEVENTACTIONRUN INSERT FAILED | ActionId={ActionId} RunId={RunId} TenantId={TenantId} MockMode={MockMode}", action.ActionId, playgroundRunId, login.ClientId, dto.MockMode); throw; } var payload = JsonSerializer.Serialize(new { actionRunId, action.ActionId, action.ActionType, login.ClientId, login.DatabaseName, action.SendTo, action.TemplateId, action.MailCc, action.MailBcc, action.ReplyTo, action.WebServiceId, action.UriParameterValue, action.ReportId, TriggerPayload = dto.Payload, ScenarioName = dto.ScenarioName ?? string.Empty }); try { await _outboxDal.InsertAsync(new ActionOutboxDTO { ActionRunId = actionRunId, DestinationTopic = destinationTopic, Payload = payload, TenantId = login.ClientId }, login, ct).ConfigureAwait(false); Activity.Current?.SetTag("gb5.outbox.topic", destinationTopic); } catch (Exception ex) { Activity.Current?.SetStatus(ActivityStatusCode.Error, ex.Message); _logger.LogError(ex, "Playground: TACTIONOUTBOX INSERT FAILED | ActionRunId={ActionRunId} RunId={RunId} ActionId={ActionId} TenantId={TenantId} Topic={Topic}", actionRunId, playgroundRunId, action.ActionId, login.ClientId, destinationTopic); throw; } _logger.LogInformation( "Playground action queued | RunId={RunId} ActionRunId={ActionRunId} ActionType={Type} Topic={Topic}", playgroundRunId, actionRunId, action.ActionType, destinationTopic); } _logger.LogInformation( "Playground trigger complete | RunId={RunId} ActionsQueued={Count} MockMode={Mock}", playgroundRunId, limited.Count, dto.MockMode); return JsonSerializer.Serialize(new TriggerEventResultDTO { RunId = playgroundRunId, TriggeredAt = DateTime.UtcNow, ActionsQueued = limited.Count, ScenarioName = dto.ScenarioName ?? string.Empty }); } public async Task GetPipelineStatusAsync( string runId, LoginDTO login, CancellationToken ct = default) { var taskCounts = await _dal.GetTaskStatusCountsAsync(runId, login, ct).ConfigureAwait(false); var actionCounts = await _dal.GetActionStatusCountsAsync(runId, login, ct).ConfigureAwait(false); var tq = new TaskQueueStatsDTO(); foreach (var row in taskCounts) { switch (row.RunStatus) { case 0: tq.Pending = row.ItemCount; break; case 1: tq.InProgress = row.ItemCount; break; case 2: tq.Completed = row.ItemCount; break; case 3: tq.Failed = row.ItemCount; break; } } var aq = new ActionQueueStatsDTO(); foreach (var row in actionCounts) { switch (row.SendStatus) { case 0: aq.Pending = row.ItemCount; break; case 1: aq.Sent = row.ItemCount; break; case 2: aq.Failed = row.ItemCount; break; } } bool isComplete = tq.Total > 0 && tq.Pending == 0 && tq.InProgress == 0; return new PipelineStatusDTO { RunId = runId, IsComplete = isComplete, TaskQueue = tq, ActionQueue = aq }; } public async Task GetActionQueueAsync( string runId, LoginDTO login, CancellationToken ct = default) { var rows = await _dal.GetActionQueueAsync(runId, login, ct).ConfigureAwait(false); return JsonSerializer.Serialize(rows); } public async Task RetryActionAsync( int actionRunId, LoginDTO login, CancellationToken ct = default) { await _dal.ResetActionRunAsync(actionRunId, login, ct).ConfigureAwait(false); _logger.LogInformation("Playground retry queued | ActionRunId={ActionRunId} UserId={UserId}", actionRunId, login.UserId); return SuccessResponse.SaveSuccessMessage; } // ── Helpers ─────────────────────────────────────────────────────────── private static System.Collections.Generic.List BuildMockActions( TriggerEventDTO dto, LoginDTO login) { // In mock mode, generate synthetic actions (up to RecordLimit) of the requested type. var list = new System.Collections.Generic.List(); for (int i = 0; i < dto.RecordLimit; i++) { list.Add(new PlaygroundActionDTO { ActionId = -(i + 1), // negative IDs for mock rows ActionType = dto.MockActionType, SendTo = dto.MockActionType == 1 ? "+1000000000" : $"demo-{i + 1}@playground.local", TemplateId = 0, WebServiceId = 0, ReportId = 0 }); } return list; } } }