using AutomationBLL.DocumentFlow; using AutomationBLL.ScriptSecretRef; using AutomationBLL.ScriptVersion; using AutomationDAL.CustomCode.Run; using AutomationDAL.DTO.Run; using GB5Shared.DTO.Framework.Login; using GB5Shared.Resource.Response; using GB5Shared.Telemetry; using Microsoft.Extensions.Logging; namespace AutomationBLL.Run; public class RunBLL : IRunBLL { private static readonly string[] TerminalStatuses = { "Success", "Failed", "Cancelled" }; // Agents long-poll for up to this long before the caller gets an empty response and retries โ€” // matches the 30s long-poll the spec describes for GET .../agents/{agentId}/next-job. private static readonly TimeSpan PollInterval = TimeSpan.FromSeconds(2); private readonly IRunDAL _dal; private readonly IScriptVersionBLL _versionBll; private readonly IScriptSecretRefBLL _secretRefBll; private readonly IAutomationRunNotifier _notifier; private readonly IDocumentFlowBLL _documentFlowBll; private readonly IAckCallbackClient _ackCallbackClient; private readonly ILogger _logger; public RunBLL( IRunDAL dal, IScriptVersionBLL versionBll, IScriptSecretRefBLL secretRefBll, IAutomationRunNotifier notifier, IDocumentFlowBLL documentFlowBll, IAckCallbackClient ackCallbackClient, ILogger logger) { _dal = dal; _versionBll = versionBll; _secretRefBll = secretRefBll; _notifier = notifier; _documentFlowBll = documentFlowBll; _ackCallbackClient = ackCallbackClient; _logger = logger; } public async Task TriggerRun(int scriptId, int versionId, string? parametersJson, LoginDTO login, CancellationToken ct) => await EnqueueRun(scriptId, versionId, $"Manual:{login.UserId}", parametersJson, null, login, ct).ConfigureAwait(false); // Shared by manual trigger and scheduled dispatch (AutomationBLL.Schedule.ScheduledDispatchBLL) โ€” // only TriggeredBy/TargetId differ between callers. public async Task EnqueueRun(int scriptId, int versionId, string triggeredBy, string? parametersJson, int? targetId, LoginDTO login, CancellationToken ct) { long runId = await EnqueueRunRaw(scriptId, versionId, triggeredBy, parametersJson, targetId, login, ct).ConfigureAwait(false); return $"{SuccessResponse.SaveSuccessMessage} {runId}"; } public async Task EnqueueRunRaw(int scriptId, int versionId, string triggeredBy, string? parametersJson, int? targetId, LoginDTO login, CancellationToken ct) { try { GB5Trace.Step("validate-enqueue-run", new { scriptId, versionId, triggeredBy }); var version = await _versionBll.GetById(versionId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"ScriptVersion {versionId} not found."); if (version.ScriptId != scriptId) throw new InvalidOperationException("VersionId does not belong to the given ScriptId."); GB5Trace.Step("enqueue-run", new { scriptId, versionId, triggeredBy, targetId }); return await _dal.InsertExecution(new RunExecutionDTO { ScriptId = scriptId, VersionId = versionId, TargetId = targetId, TriggeredBy = triggeredBy, ParametersJson = parametersJson }, login, ct).ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("enqueue-run-failed", ex); _logger.LogError(ex, "EnqueueRun failed for ScriptId {ScriptId} VersionId {VersionId} TriggeredBy {TriggeredBy}", scriptId, versionId, triggeredBy); throw; } } // Long-polls up to waitSeconds, retrying the atomic claim every PollInterval โ€” avoids hammering // the DB with a tight loop while still returning promptly the moment work appears. public async Task GetNextJob(string supportedRuntimes, int agentId, int waitSeconds, LoginDTO login, CancellationToken ct) { var deadline = DateTime.UtcNow.AddSeconds(waitSeconds); do { var runId = await _dal.ClaimNextJob(supportedRuntimes, agentId, login, ct).ConfigureAwait(false); if (runId.HasValue) { GB5Trace.Step("claim-next-job", new { runId, agentId }); var job = await _dal.GetJobDetail(runId.Value, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"Claimed RunId {runId} but its job detail could not be resolved."); var secretRefs = await _secretRefBll.GetByScript(job.ScriptId, login, ct).ConfigureAwait(false); job.SecretRefs = secretRefs .Select(s => new RunSecretRefDTO { SecretName = s.SecretName, VaultReference = s.VaultReference }) .ToList(); await _notifier.PushRunStarted(job.TenantId, job.RunId, job.TriggeredBy, ct).ConfigureAwait(false); return job; } if (DateTime.UtcNow >= deadline) return null; await Task.Delay(PollInterval, ct).ConfigureAwait(false); } while (!ct.IsCancellationRequested); return null; } public async Task SaveArtifact(long runId, string artifactType, string blobUri, LoginDTO login, CancellationToken ct) { _ = await _dal.GetById(runId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"Run {runId} not found."); long artifactId = await _dal.InsertArtifact(new RunArtifactDTO { RunId = runId, ArtifactType = artifactType, BlobUri = blobUri }, login, ct).ConfigureAwait(false); return $"{SuccessResponse.SaveSuccessMessage} {artifactId}"; } public async Task> GetArtifacts(long runId, LoginDTO login, CancellationToken ct) => await _dal.GetArtifacts(runId, login, ct).ConfigureAwait(false); public async Task GetLatestStatusForScript(int scriptId, LoginDTO login, CancellationToken ct) => await _dal.GetLatestStatusForScript(scriptId, login, ct).ConfigureAwait(false); public async Task AppendLog(long runId, string logLevel, string message, LoginDTO login, CancellationToken ct) { _ = await _dal.GetById(runId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"Run {runId} not found."); await _dal.AppendLog(runId, logLevel, message, login, ct).ConfigureAwait(false); await _notifier.PushLogAppended(runId, logLevel, message, ct).ConfigureAwait(false); return SuccessResponse.SaveSuccess; } public async Task CompleteRun(long runId, string status, string? summary, string? errorMessage, LoginDTO login, CancellationToken ct) { try { if (!TerminalStatuses.Contains(status)) throw new InvalidOperationException($"Status must be one of: {string.Join(", ", TerminalStatuses)}"); GB5Trace.Step("complete-run", new { runId, status }); int rows = await _dal.CompleteExecution(runId, status, summary, errorMessage, login, ct).ConfigureAwait(false); if (rows == 0) throw new InvalidOperationException($"Run {runId} not found or is not in Running status."); if (status == "Failed") await _notifier.PushRunFailed(login.ClientId, runId, errorMessage, ct).ConfigureAwait(false); else await _notifier.PushRunCompleted(login.ClientId, runId, status, summary, ct).ConfigureAwait(false); // No-ops for runs that aren't event-driven (no TDOCUMENTFLOW row points at this RunId). await _documentFlowBll.UpdateStatusByRunId(runId, status, login, ct).ConfigureAwait(false); // Ack-writeback (spec ยง4.6): only for a successful event-driven outbound-upload run // whose mapping configured a callback, and only if the agent actually registered an // AckReference artifact. if (status == "Success") { var flow = await _documentFlowBll.GetByRunId(runId, login, ct).ConfigureAwait(false); if (flow is { AckCallbackUrl: not null }) { var artifacts = await _dal.GetArtifacts(runId, login, ct).ConfigureAwait(false); var ackArtifact = artifacts.FirstOrDefault(a => a.ArtifactType == "AckReference"); if (ackArtifact is not null) await _ackCallbackClient.SendAsync(flow.AckCallbackUrl, flow.OrderId, ackArtifact.BlobUri, ct).ConfigureAwait(false); } } return SuccessResponse.UpdateSuccess; } catch (Exception ex) { GB5Trace.MarkFailed("complete-run-failed", ex); _logger.LogError(ex, "CompleteRun failed for RunId {RunId}", runId); throw; } } public async Task> GetHistory(int? scriptId, string? status, LoginDTO login, CancellationToken ct) => await _dal.GetHistory(scriptId, status, login, ct).ConfigureAwait(false); public async Task GetDetail(long runId, LoginDTO login, CancellationToken ct) { var execution = await _dal.GetById(runId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"Run {runId} not found."); var logs = await _dal.GetLogs(runId, login, ct).ConfigureAwait(false); var artifacts = await _dal.GetArtifacts(runId, login, ct).ConfigureAwait(false); return new RunDetailDTO { Execution = execution, Logs = logs, Artifacts = artifacts }; } }