using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; using System; using System.Threading; using System.Threading.Tasks; namespace GB5Shared.ActionProcessor { public interface IEventActionRunDAL { Task InsertAsync(EventActionRunDTO dto, LoginDTO login, CancellationToken ct = default); /// /// Returns true if a TEVENTACTIONRUN row already exists for this ActionId + CorrelationKey. /// Used by EventSubBLL to skip duplicate events (direct Dapr publish + outbox relay /// both carry the same CorrelationKey for the same logical event). /// Task ExistsByCorrelationKeyAsync(int actionId, string correlationKey, LoginDTO login, CancellationToken ct = default); Task UpdateRunStatusAsync( int actionRunId, int runStatus, LoginDTO login, CancellationToken ct = default); Task UpdateResultAsync( int actionRunId, int runStatus, int attempts, string? errorMessage, DateTime? startedOn, DateTime? completedOn, string? result, LoginDTO login, CancellationToken ct = default); /// /// Returns action-completion counts for a scheduler JobExecutionId — used to detect /// when every fanned-out action has reached a terminal state so the parent /// TJOBEXECUTION row can be finalized (Success/Failed). /// Task GetJobExecutionProgressAsync( int jobExecutionId, LoginDTO login, CancellationToken ct = default); } /// Action-completion rollup for one scheduler JobExecutionId. public sealed class JobExecutionProgressDTO { public int TotalActions { get; set; } public int CompletedActions { get; set; } public int FailedActions { get; set; } public string? LastError { get; set; } /// True once every fanned-out action has reached a terminal state. public bool IsComplete => TotalActions > 0 && CompletedActions >= TotalActions; } public sealed class EventActionRunDAL : IEventActionRunDAL { private readonly IQueryExecutor _qe; public EventActionRunDAL(IQueryExecutor qe) => _qe = qe; public async Task InsertAsync(EventActionRunDTO dto, LoginDTO login, CancellationToken ct = default) { return await _qe.ExecuteIdentityAsync( login, EventActionRunQB.INSERT_RUN, new { dto.ActionId, dto.JobExecutionId, dto.EventTypeId, dto.Payload, dto.PlaygroundRunId, dto.CorrelationKey, CreatedById = login.UserId, TenantId = login.ClientId }).ConfigureAwait(false); } public async Task ExistsByCorrelationKeyAsync( int actionId, string correlationKey, LoginDTO login, CancellationToken ct = default) { var result = await _qe.ExecuteScalarAsync( login, EventActionRunQB.EXISTS_BY_CORRELATION_KEY, new { ActionId = actionId, CorrelationKey = correlationKey, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); return result == 1; } public async Task UpdateRunStatusAsync( int actionRunId, int runStatus, LoginDTO login, CancellationToken ct = default) { await _qe.ExecuteAsync( login, EventActionRunQB.UPDATE_RUNSTATUS, new { ActionRunId = actionRunId, RunStatus = runStatus, ModifiedById = login.UserId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); } public async Task UpdateResultAsync( int actionRunId, int runStatus, int attempts, string? errorMessage, DateTime? startedOn, DateTime? completedOn, string? result, LoginDTO login, CancellationToken ct = default) { await _qe.ExecuteAsync( login, EventActionRunQB.UPDATE_RESULT, new { ActionRunId = actionRunId, RunStatus = runStatus, Attempts = attempts, ErrorMessage = errorMessage, StartedOn = startedOn, CompletedOn = completedOn, Result = result, ModifiedById = login.UserId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); } public async Task GetJobExecutionProgressAsync( int jobExecutionId, LoginDTO login, CancellationToken ct = default) { var result = await _qe.QuerySingleAsync( login, EventActionRunQB.GET_JOB_EXECUTION_PROGRESS, new { JobExecutionId = jobExecutionId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); return result ?? new JobExecutionProgressDTO(); } } }