using System; using System.Collections.Generic; using System.Threading.Tasks; using FrameworkDAL.CustomCode.GOP; using GB5Shared.Deployment; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.GOP; using GB5Shared.EventLogPublish; using GB5Shared.Telemetry; using GB5Shared.Validation; using Newtonsoft.Json.Linq; using static GB5Shared.GB5Constant.Constant; namespace FrameworkBLL.GOP { // ============================================================ // GopQueueBLL — Gateway / Submit layer only. // // Responsibilities: // 1. Validate the incoming request. // 2. Resolve the source binding → FlowId + SnapshotVersion. // 3. Atomically insert TGOPEXECUTIONQUEUE + TGOPEXECUTIONHEADER. // 4. Return 202-style DTO immediately (QueueId + ExecutionId). // // NOT responsible for: // - Executing steps (owned by GopWorkerService BackgroundService). // - Acquiring or releasing locks (owned by GopExecutionLockService). // - Retrying failed nodes (owned by GopExecutionPipeline). // // Environment MUST come from server-side IGB5Environment.GopEnvironmentCode // (GB5:Environment config, see GB5Shared/Deployment). // It is NEVER accepted from the client request payload (security violation). // ============================================================ public class GopQueueBLL : IGopQueueBLL { private readonly IGopFlowDAL _GopFlowDAL; private readonly IGopQueueDAL _GopQueueDAL; private readonly IValidation _Validation; private readonly IGB5Environment _GB5Environment; private readonly EventLogPublish _EventLog; public GopQueueBLL( IGopFlowDAL GopFlowDAL, IGopQueueDAL GopQueueDAL, IValidation Validation, IGB5Environment GB5Environment, EventLogPublish EventLog) { _GopFlowDAL = GopFlowDAL; _GopQueueDAL = GopQueueDAL; _Validation = Validation; _GB5Environment = GB5Environment; _EventLog = EventLog; } // ═══════════════════════════════════════════════════════════ // APPROVAL GATE // Resumes/rejects an execution paused at an Approval-type flow step // (see GopExecutionPipeline). Approving sets Status back to 'Queued' so // GopWorkerService re-picks it up and continues from CurrentNodeId using // LastOutputPayloadJson; rejecting fails the execution and dead-letters it. // ═══════════════════════════════════════════════════════════ public async Task ApproveGopExecutionStep(int ExecutionId, string? ApproverNotes, LoginDTO LoginDTO) { GB5Trace.Step("approve-gop-execution-step", new { ExecutionId }); var header = await _GopQueueDAL.GetExecutionById(LoginDTO.ClientId, ExecutionId, LoginDTO) ?? throw new InvalidOperationException($"Execution {ExecutionId} not found."); if (header.Status != "PendingApproval") throw new InvalidOperationException( $"Execution {ExecutionId} is not awaiting approval (Status='{header.Status}')."); bool resumed = await _GopQueueDAL.ResumeExecutionAfterApproval(LoginDTO.ClientId, ExecutionId, LoginDTO); if (!resumed) throw new InvalidOperationException( $"Execution {ExecutionId} could not be resumed — it may no longer be pending approval."); await _EventLog.PublishEventLogAsync( "GOP Execution Step Approved", new { ExecutionId, header.CurrentNodeCode, ApproverNotes }, EventTypeConstant.GOPEXECAPPROVEDEVENTTYPEID, ExecutionId, LoginDTO); return $"Execution {ExecutionId} approved — resumed for processing."; } public async Task RejectGopExecutionStep(int ExecutionId, string Reason, LoginDTO LoginDTO) { if (string.IsNullOrWhiteSpace(Reason)) throw new ArgumentException("Reason is required."); GB5Trace.Step("reject-gop-execution-step", new { ExecutionId }); var header = await _GopQueueDAL.GetExecutionById(LoginDTO.ClientId, ExecutionId, LoginDTO) ?? throw new InvalidOperationException($"Execution {ExecutionId} not found."); if (header.Status != "PendingApproval") throw new InvalidOperationException( $"Execution {ExecutionId} is not awaiting approval (Status='{header.Status}')."); string errorMessage = $"Rejected at approval gate: {Reason}"; await _GopQueueDAL.UpdateExecutionStatus(LoginDTO.ClientId, ExecutionId, "Failed", errorMessage, null, LoginDTO); await _GopQueueDAL.InsertDeadLetter(new GopDeadLetterDTO { ClientId = LoginDTO.ClientId, ExecutionId = ExecutionId, QueueId = header.QueueId, FlowId = header.FlowId, FlowCode = header.FlowCode, SnapshotVersion = header.SnapshotVersion, FailedNodeId = header.CurrentNodeId ?? 0, FailedNodeCode = header.CurrentNodeCode ?? string.Empty, PayloadJson = header.LastOutputPayloadJson ?? header.SourcePayloadJson, ErrorMessage = errorMessage, FailureCategory = "BusinessLogic" }, LoginDTO); await _EventLog.PublishEventLogAsync( "GOP Execution Step Rejected", new { ExecutionId, header.CurrentNodeCode, Reason }, EventTypeConstant.GOPEXECREJECTEDEVENTTYPEID, ExecutionId, LoginDTO); return $"Execution {ExecutionId} rejected and moved to dead-letter."; } // ═══════════════════════════════════════════════════════════ // SUBMIT EXECUTION // 1. Validate request fields. // 2. Resolve Environment from server config (never from client). // 3. Resolve SourceBinding → FlowId + SnapshotVersion. // 4. Verify active snapshot exists. // 5. Atomically insert queue row + execution header. // 6. Return DTO with QueueId + ExecutionId (worker picks up from queue). // ═══════════════════════════════════════════════════════════ public async Task SubmitExecution(GopSubmitRequestDTO RequestDTO, LoginDTO LoginDTO) { if (RequestDTO == null) throw new ArgumentNullException(nameof(RequestDTO)); if (string.IsNullOrWhiteSpace(RequestDTO.SourceCode)) throw new ArgumentException("SourceCode is required."); if (string.IsNullOrWhiteSpace(RequestDTO.PayloadJson)) throw new ArgumentException("PayloadJson is required."); // Environment resolved from server config — not from client payload. var environment = _GB5Environment.GopEnvironmentCode; var clientId = LoginDTO.ClientId; // Resolve binding var binding = await _GopFlowDAL.GetSourceBindingByCode( clientId, RequestDTO.SourceCode, environment, LoginDTO); if (binding == null) throw new InvalidOperationException( $"No active source binding found for SourceCode='{RequestDTO.SourceCode}', Environment='{environment}'."); // Verify snapshot exists var snapshot = await _GopFlowDAL.GetSnapshotByFlowVersion( clientId, binding.FlowId, binding.SnapshotVersion, LoginDTO); if (snapshot == null) throw new InvalidOperationException( $"Snapshot not found: FlowId={binding.FlowId}, Version='{binding.SnapshotVersion}'."); var flowCode = TryExtractFlowCode(snapshot.SnapshotJson); if (string.IsNullOrWhiteSpace(flowCode)) flowCode = RequestDTO.SourceCode; var queueDTO = new GopExecutionQueueDTO { ClientId = clientId, FlowId = binding.FlowId, FlowCode = flowCode, SnapshotVersion = binding.SnapshotVersion, SourceCode = RequestDTO.SourceCode, PayloadJson = RequestDTO.PayloadJson, IdempotencyKey = RequestDTO.IdempotencyKey, Priority = RequestDTO.Priority, MaxRetry = 5 }; var headerDTO = new GopExecutionHeaderDTO { ClientId = clientId, FlowId = binding.FlowId, FlowCode = flowCode, SnapshotVersion = binding.SnapshotVersion, SourcePayloadJson = RequestDTO.PayloadJson, IsReplay = false }; // Atomic insert (one transaction in DAL) var (queueId, execId) = await _GopQueueDAL.SubmitToQueue(queueDTO, headerDTO, LoginDTO); // GopWorkerService (BackgroundService) polls TGOPEXECUTIONQUEUE and picks this up. return new GopSubmitResponseDTO { QueueId = queueId, ExecutionId = execId, FlowId = binding.FlowId, FlowCode = flowCode, SnapshotVersion = binding.SnapshotVersion, Status = "Pending", Message = $"Execution queued. QueueId={queueId}, ExecutionId={execId}." }; } // ═══════════════════════════════════════════════════════════ // REPLAY // Re-submits a past execution's original payload as a brand-new execution. // Ported from the legacy ExecutionBLL.ReplayExecution idea — the modern stack // had no way to resubmit a dead-lettered/failed execution before this. // ═══════════════════════════════════════════════════════════ public async Task ReplayGopExecution(int ExecutionId, LoginDTO LoginDTO) { GB5Trace.Step("replay-gop-execution", new { ExecutionId }); var original = await _GopQueueDAL.GetExecutionById(LoginDTO.ClientId, ExecutionId, LoginDTO) ?? throw new InvalidOperationException($"Execution {ExecutionId} not found."); var queueDTO = new GopExecutionQueueDTO { ClientId = LoginDTO.ClientId, FlowId = original.FlowId, FlowCode = original.FlowCode, SnapshotVersion = original.SnapshotVersion, SourceCode = "REPLAY", PayloadJson = original.SourcePayloadJson, Priority = 5, MaxRetry = 5 }; var headerDTO = new GopExecutionHeaderDTO { ClientId = LoginDTO.ClientId, ParentExecutionId = original.ExecutionId, FlowId = original.FlowId, FlowCode = original.FlowCode, SnapshotVersion = original.SnapshotVersion, SourcePayloadJson = original.SourcePayloadJson, IsReplay = true }; var (queueId, execId) = await _GopQueueDAL.SubmitToQueue(queueDTO, headerDTO, LoginDTO); await _EventLog.PublishEventLogAsync( "GOP Execution Replayed", new { OriginalExecutionId = ExecutionId, NewExecutionId = execId }, EventTypeConstant.GOPEXECUTIONREPLAYEDEVENTTYPEID, execId, LoginDTO); return new GopSubmitResponseDTO { QueueId = queueId, ExecutionId = execId, FlowId = original.FlowId, FlowCode = original.FlowCode, SnapshotVersion = original.SnapshotVersion, Status = "Pending", Message = $"Execution {ExecutionId} replayed as new execution {execId} (QueueId={queueId})." }; } // ═══════════════════════════════════════════════════════════ // QUERY METHODS // ═══════════════════════════════════════════════════════════ public async Task GetExecutionStatus(int ExecutionId, LoginDTO LoginDTO) { var header = await _GopQueueDAL.GetExecutionById(LoginDTO.ClientId, ExecutionId, LoginDTO); var stateLogs = await _GopQueueDAL.GetStateLogsByExecution(LoginDTO.ClientId, ExecutionId, LoginDTO); return new { Header = header, StateLogs = stateLogs }; } public async Task> GetExecutionNodeLogs(int ExecutionId, LoginDTO LoginDTO) => await _GopQueueDAL.GetNodeLogsByExecution(LoginDTO.ClientId, ExecutionId, LoginDTO); public async Task> GetQueueDashboard(LoginDTO LoginDTO) => await _GopQueueDAL.GetQueueDashboard(LoginDTO.ClientId, LoginDTO); // ── Dead Letter ─────────────────────────────────────────── public async Task> GetDeadLetters(LoginDTO LoginDTO) => await _GopQueueDAL.GetDeadLetters(LoginDTO.ClientId, LoginDTO); public async Task ResolveDeadLetter(int DlqId, string ResolvedBy, LoginDTO LoginDTO) { if (string.IsNullOrWhiteSpace(ResolvedBy)) throw new ArgumentException("ResolvedBy is required."); await _GopQueueDAL.ResolveDeadLetter(LoginDTO.ClientId, DlqId, ResolvedBy, LoginDTO); } // ── Remediation ─────────────────────────────────────────── public async Task OpenRemediation(GopOpenRemediationRequestDTO RequestDTO, LoginDTO LoginDTO) { if (RequestDTO == null) throw new ArgumentNullException(nameof(RequestDTO)); var dto = new GopRemediationDTO { ClientId = LoginDTO.ClientId, ExecutionId = RequestDTO.ExecutionId, QueueId = RequestDTO.QueueId, FlowId = RequestDTO.FlowId, FlowCode = RequestDTO.FlowCode, SnapshotVersion = RequestDTO.SnapshotVersion, FailedNodeId = RequestDTO.FailedNodeId, FailedNodeCode = RequestDTO.FailedNodeCode, OriginalPayloadJson = RequestDTO.OriginalPayloadJson, TargetResponseJson = RequestDTO.TargetResponseJson, IssueListJson = RequestDTO.IssueListJson, AssignedTo = RequestDTO.AssignedTo, Status = "Open" }; return await _GopQueueDAL.InsertRemediation(dto, LoginDTO); } public async Task ResolveRemediation(GopResolveRemediationRequestDTO RequestDTO, LoginDTO LoginDTO) { if (RequestDTO == null) throw new ArgumentNullException(nameof(RequestDTO)); if (string.IsNullOrWhiteSpace(RequestDTO.ResolutionNotes)) throw new ArgumentException("ResolutionNotes are required."); await _GopQueueDAL.UpdateRemediation( LoginDTO.ClientId, RequestDTO.RemediationId, "Resolved", RequestDTO.ResolutionNotes, LoginDTO.UserCode, LoginDTO); } public async Task> GetOpenRemediations(LoginDTO LoginDTO) => await _GopQueueDAL.GetOpenRemediations(LoginDTO.ClientId, LoginDTO); // ── Helper ──────────────────────────────────────────────── private static string TryExtractFlowCode(string snapshotJson) { try { var obj = JObject.Parse(snapshotJson); return obj["flowCode"]?.ToString() ?? string.Empty; } catch { return string.Empty; } } } }