using System; using System.Collections.Generic; using System.Diagnostics; using System.Linq; using System.Threading; using System.Threading.Tasks; using FrameworkDAL.CustomCode.GOP; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.GOP; using GB5Shared.Telemetry; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; namespace FrameworkBLL.GOP.Worker { // ============================================================ // GopExecutionPipeline — graph-traversal execution engine. // // Execution flow for a single GopExecution: // 1. Load TGOPEXECUTIONHEADER to get FlowId + SnapshotVersion. // 2. Load snapshot steps (MGOPFLOWSNAPSHOTSTEP) + snapshot edges // (MGOPFLOWSNAPSHOTSTEPEDGE) — the immutable graph resolved at publish time. // 3. Starting at the single root step (or, on resume, the step after the // Approval gate that paused it), dispatch each step to its keyed // INodeExecutor and follow outgoing edges whose EdgeCondition matches the // step's outcome (Always / OnSuccess / OnFailure) — supporting conditional // branching and fan-out to multiple next-steps. A fan-in (join) step is only // enqueued once every one of its incoming edges has been resolved (fired, or // determined to never fire because its source took the other branch) — see // ResolveOutgoingEdges — so a diamond-shaped graph runs the join step exactly // once with the payload from whichever branch actually completed, instead of // running early on the first-arriving branch and silently discarding the rest. // 4. An "Approval" step pauses the execution (Status='PendingApproval') and // returns immediately — GopQueueBLL.ApproveGopExecutionStep resumes it by // setting Status back to 'Queued'; this pipeline then resumes traversal // from the step(s) reachable via the gate's OnSuccess/Always edges, using // LastOutputPayloadJson instead of SourcePayloadJson. // 5. A step that fails after exhausting retries follows its OnFailure edges // if any exist (lets a flow route failures to a remediation branch instead // of the target); with no OnFailure edges, it dead-letters as before. // // Retry: configurable per step from GopFlowSnapshotStepDTO.MaxRetries. // Jitter: random ±20% applied to base delay to avoid thundering herd. // ============================================================ public class GopExecutionPipeline { private readonly IGopQueueDAL _QueueDal; private readonly IGopFlowDAL _FlowDal; private readonly IServiceProvider _ServiceProvider; private readonly IGopHubNotifier _HubNotifier; private readonly ILogger _Logger; private static readonly Random Jitter = new(); public GopExecutionPipeline( IGopQueueDAL queueDal, IGopFlowDAL flowDal, IServiceProvider serviceProvider, IGopHubNotifier hubNotifier, ILogger logger) { _QueueDal = queueDal; _FlowDal = flowDal; _ServiceProvider = serviceProvider; _HubNotifier = hubNotifier; _Logger = logger; } public async Task ExecuteAsync( int executionId, int clientId, LoginDTO loginDTO, CancellationToken ct) { var sw = Stopwatch.StartNew(); GB5Trace.Step("gop-pipeline-start", new { executionId, clientId }); // ── Load execution header ────────────────────────────────────────────── var header = await _QueueDal.GetExecutionById(clientId, executionId, loginDTO); if (header is null) { _Logger.LogError( "GopPipeline: execution {ExecutionId} not found — aborting", executionId); GB5Trace.MarkFailed("gop-pipeline-execution-not-found"); return; } await _QueueDal.UpdateExecutionStatus( clientId, executionId, "Running", null, null, loginDTO); _Logger.LogInformation( "GopPipeline: starting execution {ExecutionId} flow={FlowCode}", executionId, header.FlowCode); await _HubNotifier.PushExecutionUpdate(clientId, new GopExecutionStatusDTO { ExecutionId = executionId, FlowCode = header.FlowCode, SnapshotVersion = header.SnapshotVersion, Status = "Running", StartedAt = header.StartedAt }, ct); // ── Load snapshot graph (steps + edges) ──────────────────────────────── var steps = (await _FlowDal.GetSnapshotSteps(clientId, header.FlowId, header.SnapshotVersion, loginDTO)).ToList(); if (steps.Count == 0) { _Logger.LogWarning( "GopPipeline: no steps found for execution {ExecutionId} — marking complete", executionId); await MarkSuccessAsync(clientId, executionId, header.FlowCode, header.StartedAt, sw.ElapsedMilliseconds, loginDTO, ct, null); await InsertExecutionMetricsAsync( clientId, executionId, header.FlowId, header.FlowCode, sw.ElapsedMilliseconds, 0, 0, true, loginDTO); return; } var edges = (await _FlowDal.GetSnapshotStepEdges(clientId, header.FlowId, header.SnapshotVersion, loginDTO, ct)).ToList(); var stepsById = steps.ToDictionary(s => s.NodeId); // ── Determine the traversal starting point ───────────────────────────── // Resume-from-approval is detected by the *paused node's own type*, not by // execution Status — Status='Queued' with CurrentNodeId set also happens on // ordinary crash recovery (GopWorkerService resets stuck 'Running' rows back // to 'Queued' without clearing CurrentNodeId), and that case must NOT skip // re-execution of a step that may never have actually completed. Only a node // of type "Approval" can have been "completed" by an external action rather // than by the pipeline itself, so only that case is safe to resume past. var workQueue = new Queue<(GopFlowSnapshotStepDTO Step, string? InputPayload)>(); var visited = new HashSet(); // Fan-in (join) tracking: a node with N incoming edges must not execute until // ALL N have been "resolved" — either fired (its source took the matching // outcome) or determined to never fire (source took the other branch). // Resolving on every outgoing edge (not just the ones that fire) means a join // whose sibling branch was skipped still gets a resolution rather than waiting // forever, and a join fed by multiple branches only runs once every branch has // actually reported in — instead of running (and discarding data) on whichever // branch happens to arrive first, or never running for a genuine dead sibling. var incomingEdgeCounts = edges .GroupBy(e => e.ToNodeId) .ToDictionary(g => g.Key, g => g.Count()); var resolvedIncomingCount = new Dictionary(); var pendingJoinPayload = new Dictionary(); void ResolveOutgoingEdges(int fromNodeId, string outcome, string? payload) { foreach (var edge in edges.Where(e => e.FromNodeId == fromNodeId).OrderBy(e => e.ToNodeId)) { if (!stepsById.TryGetValue(edge.ToNodeId, out var targetStep)) continue; bool fires = outcome == "Success" ? edge.EdgeCondition is "Always" or "OnSuccess" : edge.EdgeCondition == "OnFailure"; if (fires) pendingJoinPayload[edge.ToNodeId] = payload; resolvedIncomingCount.TryGetValue(edge.ToNodeId, out var resolvedSoFar); resolvedIncomingCount[edge.ToNodeId] = resolvedSoFar + 1; int totalIncoming = incomingEdgeCounts.GetValueOrDefault(edge.ToNodeId, 1); if (resolvedIncomingCount[edge.ToNodeId] < totalIncoming) continue; // Every incoming edge for this node has now been resolved. Only enqueue // if at least one of them actually fired — a join whose every incoming // edge was skipped (all sources took the other branch) is correctly // never reached, matching normal DAG semantics. if (pendingJoinPayload.TryGetValue(edge.ToNodeId, out var finalPayload)) workQueue.Enqueue((targetStep, finalPayload)); } } bool isApprovalResume = header.CurrentNodeId.HasValue && stepsById.TryGetValue(header.CurrentNodeId.Value, out var pausedStep) && string.Equals(pausedStep.NodeType, "Approval", StringComparison.OrdinalIgnoreCase); if (isApprovalResume) { // Resuming after an Approval gate — treat the gate as having succeeded // and continue along its outgoing edges, using the payload captured // when it paused. visited.Add(header.CurrentNodeId!.Value); ResolveOutgoingEdges(header.CurrentNodeId.Value, "Success", header.LastOutputPayloadJson); } else { // Fresh execution, or a crash-recovered one — restart from the root using // the original source payload (safe, if wasteful for the latter case, // matching pre-existing crash-recovery behavior). var root = steps.FirstOrDefault(s => !edges.Any(e => e.ToNodeId == s.NodeId)); if (root is null) { // Should never happen — GopPublishService.ValidateAndSortDag requires // exactly one root before a flow can be published. const string msg = "No root step found in published snapshot — publish-time DAG validation should have caught this."; _Logger.LogError("GopPipeline: execution {ExecutionId} — {Message}", executionId, msg); GB5Trace.MarkFailed("gop-pipeline-no-root"); await _QueueDal.UpdateExecutionStatus(clientId, executionId, "Failed", msg, sw.ElapsedMilliseconds, loginDTO); return; } workQueue.Enqueue((root, header.SourcePayloadJson)); } // ── Traverse the graph ────────────────────────────────────────────────── string? lastOutputPayload = header.LastOutputPayloadJson; bool unhandledFailure = false; string? failureErrorMessage = null; int totalNodeCount = 0; int totalRetries = 0; while (workQueue.Count > 0 && !ct.IsCancellationRequested) { var (step, inputPayload) = workQueue.Dequeue(); if (!visited.Add(step.NodeId)) continue; await _QueueDal.UpdateExecutionCurrentNode( clientId, executionId, step.NodeId, step.NodeCode, loginDTO); if (string.Equals(step.NodeType, "Approval", StringComparison.OrdinalIgnoreCase)) { await _QueueDal.UpdateExecutionApprovalPending( clientId, executionId, step.NodeId, step.NodeCode, inputPayload, step.AssignedRole, step.AssignedToUserId, loginDTO); await _HubNotifier.PushExecutionUpdate(clientId, new GopExecutionStatusDTO { ExecutionId = executionId, FlowCode = header.FlowCode, SnapshotVersion = header.SnapshotVersion, Status = "PendingApproval", CurrentNodeCode = step.NodeCode, TotalDurationMs = sw.ElapsedMilliseconds, StartedAt = header.StartedAt, AssignedRole = step.AssignedRole, AssignedToUserId = step.AssignedToUserId }, ct); _Logger.LogInformation( "GopPipeline: execution {ExecutionId} paused at Approval gate {NodeCode}", executionId, step.NodeCode); GB5Trace.Step("gop-pipeline-approval-pending", new { executionId, nodeCode = step.NodeCode }); return; // Leave the execution paused — no Success/Failed transition. } var (stepOk, output, stepError, retryAttempts) = await ExecuteStepWithRetriesAsync( step, header, inputPayload, loginDTO, ct); totalNodeCount++; totalRetries += retryAttempts; await _HubNotifier.PushExecutionUpdate(clientId, new GopExecutionStatusDTO { ExecutionId = executionId, FlowCode = header.FlowCode, SnapshotVersion = header.SnapshotVersion, Status = stepOk ? "Running" : "Failed", CurrentNodeCode = step.NodeCode, TotalDurationMs = sw.ElapsedMilliseconds, StartedAt = header.StartedAt, ErrorMessage = stepOk ? null : stepError }, ct); if (stepOk) { lastOutputPayload = output ?? inputPayload; ResolveOutgoingEdges(step.NodeId, "Success", lastOutputPayload); } else { bool hasFailureEdge = edges.Any(e => e.FromNodeId == step.NodeId && e.EdgeCondition == "OnFailure"); if (hasFailureEdge) { // Route to a failure-handling branch instead of dead-lettering. ResolveOutgoingEdges(step.NodeId, "Failure", inputPayload); } else { unhandledFailure = true; failureErrorMessage = stepError; await InsertDeadLetterAsync(clientId, executionId, header, step, inputPayload, stepError, loginDTO); break; } } } sw.Stop(); if (unhandledFailure) { GB5Trace.MarkFailed("gop-pipeline-execution-failed"); await _QueueDal.UpdateExecutionStatus( clientId, executionId, "Failed", failureErrorMessage, sw.ElapsedMilliseconds, loginDTO); await InsertExecutionMetricsAsync( clientId, executionId, header.FlowId, header.FlowCode, sw.ElapsedMilliseconds, totalNodeCount, totalRetries, false, loginDTO); } else { await MarkSuccessAsync(clientId, executionId, header.FlowCode, header.StartedAt, sw.ElapsedMilliseconds, loginDTO, ct, lastOutputPayload); await InsertExecutionMetricsAsync( clientId, executionId, header.FlowId, header.FlowCode, sw.ElapsedMilliseconds, totalNodeCount, totalRetries, true, loginDTO); } } private async Task<(bool Success, string? Output, string? Error, int RetryAttempts)> ExecuteStepWithRetriesAsync( GopFlowSnapshotStepDTO step, GopExecutionHeaderDTO header, string? inputPayload, LoginDTO loginDTO, CancellationToken ct) { string? outputPayload = null; string? errorMessage = null; var executor = _ServiceProvider.GetKeyedService(step.NodeType); if (executor is null) { errorMessage = $"No INodeExecutor registered for NodeType '{step.NodeType}'."; _Logger.LogError("GopPipeline: {ErrorMessage}", errorMessage); return (false, null, errorMessage, 0); } int maxRetries = Math.Max(0, step.MaxRetries); for (int attempt = 1; attempt <= maxRetries + 1; attempt++) { var sw = Stopwatch.StartNew(); // Insert node log row (Running). var nodeLog = new GopExecutionNodeLogDTO { ClientId = header.ClientId, ExecutionId = header.ExecutionId, NodeId = step.NodeId, NodeCode = step.NodeCode, StepOrder = step.StepOrder, AttemptNo = attempt, Status = "Running", RequestPayloadJson = inputPayload }; int nodeLogId = await _QueueDal.InsertNodeLog(nodeLog, loginDTO); try { string? result = await executor.ExecuteAsync( step, header, inputPayload, loginDTO, ct); sw.Stop(); await _QueueDal.UpdateNodeLog( header.ClientId, nodeLogId, "Success", result, null, sw.ElapsedMilliseconds, null, null, loginDTO); await InsertNodeMetricsAsync( header.ClientId, header.ExecutionId, step.NodeId, step.NodeCode, sw.ElapsedMilliseconds, attempt - 1, true, loginDTO); outputPayload = result; return (true, outputPayload, null, attempt - 1); } catch (OperationCanceledException) { throw; // Don't swallow cancellation. } catch (Exception ex) { sw.Stop(); await _QueueDal.UpdateNodeLog( header.ClientId, nodeLogId, "Failed", null, null, sw.ElapsedMilliseconds, "STEP_FAILED", ex.Message, loginDTO); errorMessage = ex.Message; _Logger.LogWarning(ex, "GopPipeline: step {NodeCode} attempt {Attempt}/{Max} failed for execution {ExecutionId}", step.NodeCode, attempt, maxRetries + 1, header.ExecutionId); if (attempt > maxRetries) { await InsertNodeMetricsAsync( header.ClientId, header.ExecutionId, step.NodeId, step.NodeCode, sw.ElapsedMilliseconds, attempt - 1, false, loginDTO); break; } // Exponential backoff with ±20% jitter. double baseDelayMs = Math.Pow(2, attempt) * 1000; double jitter = baseDelayMs * 0.2 * (Jitter.NextDouble() * 2 - 1); await Task.Delay(TimeSpan.FromMilliseconds(baseDelayMs + jitter), ct); } } return (false, null, errorMessage, maxRetries); } private async Task InsertNodeMetricsAsync( int clientId, int executionId, int nodeId, string nodeCode, long durationMs, int retryAttempts, bool isSuccess, LoginDTO loginDTO) { try { await _QueueDal.InsertNodeMetrics(new GopExecutionNodeMetricsDTO { ClientId = clientId, ExecutionId = executionId, NodeId = nodeId, NodeCode = nodeCode, DurationMs = durationMs, RetryAttempts = retryAttempts, IsSuccess = isSuccess }, loginDTO); } catch (Exception ex) { _Logger.LogWarning(ex, "GopPipeline: failed to record node metrics for execution {ExecutionId} node {NodeCode}", executionId, nodeCode); } } private async Task InsertExecutionMetricsAsync( int clientId, int executionId, int flowId, string flowCode, long totalDurationMs, int nodeCount, int totalRetries, bool isSuccess, LoginDTO loginDTO) { try { await _QueueDal.InsertExecutionMetrics(new GopExecutionMetricsDTO { ClientId = clientId, ExecutionId = executionId, FlowId = flowId, FlowCode = flowCode, TotalDurationMs = totalDurationMs, NodeCount = nodeCount, TotalRetries = totalRetries, IsSuccess = isSuccess }, loginDTO); } catch (Exception ex) { _Logger.LogWarning(ex, "GopPipeline: failed to record execution metrics for execution {ExecutionId}", executionId); } } private async Task MarkSuccessAsync( int clientId, int executionId, string flowCode, DateTime startedAt, long durationMs, LoginDTO loginDTO, CancellationToken ct, string? lastOutputPayloadJson) { await _QueueDal.UpdateExecutionStatus( clientId, executionId, "Success", null, durationMs, loginDTO, lastOutputPayloadJson); _Logger.LogInformation( "GopPipeline: execution {ExecutionId} completed successfully in {DurationMs}ms", executionId, durationMs); await _HubNotifier.PushExecutionUpdate(clientId, new GopExecutionStatusDTO { ExecutionId = executionId, FlowCode = flowCode, Status = "Success", TotalDurationMs = durationMs, StartedAt = startedAt }, ct); } private async Task InsertDeadLetterAsync( int clientId, int executionId, GopExecutionHeaderDTO header, GopFlowSnapshotStepDTO step, string? failedPayload, string? errorMsg, LoginDTO loginDTO) { try { // QueueId/FlowId/FlowCode/SnapshotVersion/PayloadJson were previously left unset // (defaulting to 0/empty) — TGOPDEADLETTER.GOPFLOWID has a real FK to MGOPFLOW, so // FlowId=0 guaranteed an FK violation on every real dead-letter insert, silently // swallowed by this method's own catch block. Populated from `header`/the step's // actual input payload, both already in scope at the call site. await _QueueDal.InsertDeadLetter(new GopDeadLetterDTO { ClientId = clientId, ExecutionId = executionId, QueueId = header.QueueId, FlowId = header.FlowId, FlowCode = header.FlowCode, SnapshotVersion = header.SnapshotVersion, FailedNodeId = step.NodeId, FailedNodeCode = step.NodeCode, PayloadJson = failedPayload ?? header.SourcePayloadJson, ErrorMessage = errorMsg ?? "Unknown error", MovedToDlqAt = DateTime.UtcNow }, loginDTO); } catch (Exception ex) { _Logger.LogError(ex, "GopPipeline: failed to insert dead-letter for execution {ExecutionId}", executionId); } } } }