using System; using System.Collections.Generic; using System.Linq; using System.Security.Cryptography; using System.Text; using System.Text.Json; using System.Threading; 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.GenerateAutoNumber; using GB5Shared.QueryExecutor; using GB5Shared.Telemetry; using Microsoft.Extensions.Logging; using static GB5Shared.GB5Constant.Constant; namespace FrameworkBLL.GOP { // ============================================================ // GopPublishService — validates and publishes a flow snapshot. // // Publish algorithm: // 1. Load the design-time graph — MGOPFLOWSTEP + MGOPFLOWSTEPEDGE for the flow // (NOT MGOPFLOWSNAPSHOTSTEP — that table only holds *already-published* // snapshots and is empty for a first-time publish under a new label). // 2. DAG validation via Kahn's algorithm: every edge must reference a real // step, exactly one root (a step with no incoming edge) must exist, and // Kahn's algorithm must be able to consume every step — if it can't, the // remaining steps form a cycle. // 3. Resolve each ApiCall step against MGOPTARGETOPERATION (endpoint, method, // schemas, spec hash) for drift detection; Qualifier/Mapper/Approval steps // carry no external endpoint. // 4. Deactivate any currently Active snapshots for this flow // (UPDATE MGOPFLOWSNAPSHOT SET ISACTIVE=0) inside transaction. // 5. Insert new MGOPFLOWSNAPSHOT row inside same transaction. // 6. Insert MGOPFLOWSNAPSHOTSTEP rows in topological order, and // MGOPFLOWSNAPSHOTSTEPEDGE rows mirroring the design-time edges — this is // the immutable edge set GopExecutionPipeline branches on at runtime. // // Idempotency: if the server crashes mid-publish, a retry will: // - Detect the unique constraint on (GOPFLOWID, SNAPSHOTVERSION, ENVIRONMENT, TENANTID) // - Deactivate the partial snapshot (IsActive=0) first // - Then start fresh // // Environment: ALWAYS from IGB5Environment.GopEnvironmentCode (GB5:Environment config, // see GB5Shared/Deployment) — NEVER from the caller's request. // ============================================================ public class GopPublishService : IGopPublishService { private readonly IGopFlowDAL _FlowDal; private readonly IGopTargetOperationDAL _TargetOperationDal; private readonly IQueryExecutor _QueryExecutor; private readonly AutoNumber _AutoNumber; private readonly IGB5Environment _GB5Environment; private readonly EventLogPublish _EventLog; private readonly ILogger _Logger; public GopPublishService( IGopFlowDAL flowDal, IGopTargetOperationDAL targetOperationDal, IQueryExecutor queryExecutor, AutoNumber autoNumber, IGB5Environment gb5Environment, EventLogPublish eventLog, ILogger logger) { _FlowDal = flowDal; _TargetOperationDal = targetOperationDal; _QueryExecutor = queryExecutor; _AutoNumber = autoNumber; _GB5Environment = gb5Environment; _EventLog = eventLog; _Logger = logger; } public async Task PublishAsync( int flowId, string snapshotVersionLabel, LoginDTO loginDTO, CancellationToken ct = default) { if (flowId == 0) throw new ArgumentException("FlowId is required."); if (string.IsNullOrWhiteSpace(snapshotVersionLabel)) throw new ArgumentException("SnapshotVersionLabel is required."); var environment = _GB5Environment.GopEnvironmentCode; GB5Trace.Step("publish-gop-flow", new { flowId, snapshotVersionLabel, environment }); // ── Load flow ───────────────────────────────────────────────────────── var flow = await _FlowDal.GetFlowById(loginDTO.ClientId, flowId, loginDTO); if (flow is null) throw new InvalidOperationException($"Flow {flowId} not found."); // ── Load design-time graph (source of truth — NOT any existing snapshot) ── var steps = (await _FlowDal.GetFlowSteps(flowId, loginDTO, ct)).ToList(); var edges = (await _FlowDal.GetFlowStepEdges(flowId, loginDTO, ct)).ToList(); // ── DAG validation + topological order (Kahn's algorithm) ───────────── var dagResult = ValidateAndSortDag(steps, edges); if (dagResult.Errors.Count > 0) { GB5Trace.MarkFailed("publish-gop-flow-dag-invalid"); _Logger.LogWarning( "GopPublishService: flow {FlowId} DAG validation failed: {Errors}", flowId, string.Join("; ", dagResult.Errors)); return GopPublishResultDTO.Failure(dagResult.Errors); } // ── Resolve each step against its ApiCall target (if any) ───────────── var snapshotSteps = new List(dagResult.SortedSteps.Count); for (int i = 0; i < dagResult.SortedSteps.Count; i++) { var step = dagResult.SortedSteps[i]; var snapshotStep = new GopFlowSnapshotStepDTO { ClientId = loginDTO.ClientId, StepOrder = i + 1, NodeId = step.FlowStepId, NodeCode = step.NodeCode, NodeType = step.NodeType, LogicalServiceName = step.LogicalServiceName ?? string.Empty, TimeoutSeconds = step.TimeoutSeconds, MaxRetries = step.MaxRetries, SupportsRemediation = step.SupportsRemediation, AssignedRole = step.AssignedRole, AssignedToUserId = step.AssignedToUserId }; // TargetOperationId is a "-1 = not set" sentinel, not a "0 = unset" one -- // real AutoNumber-assigned IDs in this repo are large NEGATIVE numbers // (e.g. -1200009989), so a "> 0" check silently skipped every real one. if (string.Equals(step.NodeType, "ApiCall", StringComparison.OrdinalIgnoreCase) && step.TargetOperationId != -1) { var targetOp = await _TargetOperationDal.GetById(step.TargetOperationId, loginDTO, ct); if (targetOp is null) dagResult.Errors.Add($"Step '{step.NodeCode}' references TargetOperationId {step.TargetOperationId}, which does not exist."); else { // ResolvedEndpoint is the BASE URL only — ApiCallNodeExecutor joins it // with RelativePath at runtime. Putting the full path in both would // double it up (e.g. ".../UOM/SaveUOM/UOM/SaveUOM"). snapshotStep.ResolvedEndpoint = targetOp.BaseUrl.TrimEnd('/'); snapshotStep.HttpMethod = targetOp.HttpMethod; snapshotStep.RelativePath = targetOp.RelativePath; snapshotStep.RequestSchemaJson = targetOp.RequestSchemaJson; snapshotStep.ResponseSchemaJson = targetOp.ResponseSchemaJson; snapshotStep.SpecHash = targetOp.SpecHash ?? string.Empty; } } snapshotSteps.Add(snapshotStep); } if (dagResult.Errors.Count > 0) { GB5Trace.MarkFailed("publish-gop-flow-target-resolution-failed"); return GopPublishResultDTO.Failure(dagResult.Errors); } // ── Auto-number the snapshot ────────────────────────────────────────── var auto = await _AutoNumber.GetNumberAsync( 1, AUTONUMBERCONSTANT.GOPFLOWSNAPSHOT, loginDTO); int snapshotId = auto.StartNumber; var aggregatePayload = new { Steps = snapshotSteps, Edges = edges.Select(e => new { e.FromStepId, e.ToStepId, e.EdgeLabel, e.EdgeCondition }) }; string snapshotJson = JsonSerializer.Serialize(aggregatePayload); string specAggregateHash = ComputeHash(snapshotJson); var snapshot = new GopFlowSnapshotDTO { ClientId = loginDTO.ClientId, SnapshotId = snapshotId, FlowId = flowId, SnapshotVersion = snapshotVersionLabel, Environment = environment, IsActive = true, PublishedAt = DateTime.UtcNow, PublishedBy = loginDTO.UserId.ToString(), SpecAggregateHash = specAggregateHash, SnapshotJson = snapshotJson }; // ── Transaction: deactivate old → insert new ────────────────────────── var trans = await _QueryExecutor.BeginTransactionAsync(loginDTO); try { // Deactivate any currently active snapshots for this flow. await _FlowDal.DeactivateSnapshotsByFlow(flowId, loginDTO, trans); // Insert the new snapshot. await _FlowDal.InsertSnapshot(snapshot, loginDTO, trans); // Insert snapshot steps in topological order. int stepAutoBase = (await _AutoNumber.GetNumberAsync( snapshotSteps.Count, AUTONUMBERCONSTANT.GOPFLOWSNAPSHOTSTEP, loginDTO)).StartNumber; for (int i = 0; i < snapshotSteps.Count; i++) { var step = snapshotSteps[i]; step.SnapshotStepId = stepAutoBase + i; step.SnapshotId = snapshotId; step.SnapshotVersion = snapshotVersionLabel; await _FlowDal.InsertSnapshotStep(step, loginDTO, trans); } // Insert snapshot edges — the immutable routing GopExecutionPipeline traverses. if (edges.Count > 0) { int edgeAutoBase = (await _AutoNumber.GetNumberAsync( edges.Count, AUTONUMBERCONSTANT.GOPFLOWSNAPSHOTSTEPEDGE, loginDTO)).StartNumber; for (int i = 0; i < edges.Count; i++) { var edge = edges[i]; await _FlowDal.InsertSnapshotStepEdge(new GopFlowSnapshotStepEdgeDTO { ClientId = loginDTO.ClientId, SnapshotStepEdgeId = edgeAutoBase + i, SnapshotId = snapshotId, SnapshotVersion = snapshotVersionLabel, FromNodeId = edge.FromStepId, ToNodeId = edge.ToStepId, EdgeLabel = edge.EdgeLabel, EdgeCondition = edge.EdgeCondition }, loginDTO, trans); } } await _QueryExecutor.CommitAsync(trans); _Logger.LogInformation( "GopPublishService: flow {FlowId} published as snapshot {SnapshotId} " + "version='{Version}' env={Env} steps={StepCount} edges={EdgeCount}", flowId, snapshotId, snapshotVersionLabel, environment, snapshotSteps.Count, edges.Count); await _EventLog.PublishEventLogAsync( "GOP Flow Published", new { flowId, flow.FlowCode, snapshotId, snapshotVersionLabel, environment }, EventTypeConstant.GOPFLOWPUBLISHEDEVENTTYPEID, flowId, loginDTO, ct: ct); return GopPublishResultDTO.Ok(snapshotId, snapshotVersionLabel); } catch (Exception ex) { await _QueryExecutor.RollbackAsync(trans); GB5Trace.MarkFailed("publish-gop-flow-failed", ex); throw; } } // ── DAG validation + topological sort (Kahn's algorithm) ────────────────── private static DagValidationResult ValidateAndSortDag( List steps, List edges) { var errors = new List(); if (steps.Count == 0) { errors.Add("Flow has no steps — cannot publish an empty flow."); return new DagValidationResult { Errors = errors }; } var stepIds = steps.Select(s => s.FlowStepId).ToHashSet(); var duplicates = steps.GroupBy(s => s.FlowStepId).Where(g => g.Count() > 1).Select(g => g.Key); foreach (var dup in duplicates) errors.Add($"Duplicate step id {dup} in flow steps."); foreach (var approvalStep in steps.Where(s => string.Equals(s.NodeType, "Approval", StringComparison.OrdinalIgnoreCase))) if (string.IsNullOrWhiteSpace(approvalStep.AssignedRole) && approvalStep.AssignedToUserId is null) errors.Add($"Approval step '{approvalStep.NodeCode}' requires an AssignedRole or AssignedToUserId before it can be published."); foreach (var edge in edges) { if (!stepIds.Contains(edge.FromStepId)) errors.Add($"Edge references unknown FromStepId {edge.FromStepId}."); if (!stepIds.Contains(edge.ToStepId)) errors.Add($"Edge references unknown ToStepId {edge.ToStepId}."); } if (errors.Count > 0) return new DagValidationResult { Errors = errors }; // In-degree per step, adjacency FromStepId -> [ToStepId...] var inDegree = steps.ToDictionary(s => s.FlowStepId, _ => 0); var adjacency = steps.ToDictionary(s => s.FlowStepId, _ => new List()); foreach (var edge in edges) { adjacency[edge.FromStepId].Add(edge.ToStepId); inDegree[edge.ToStepId]++; } var roots = steps.Where(s => inDegree[s.FlowStepId] == 0).ToList(); if (roots.Count == 0) { errors.Add("Flow has no root step (every step has an incoming edge) — likely a cycle with no entry point."); return new DagValidationResult { Errors = errors }; } if (roots.Count > 1) { errors.Add( $"Flow has {roots.Count} root steps ({string.Join(", ", roots.Select(r => r.NodeCode))}) " + "— exactly one entry step is required."); return new DagValidationResult { Errors = errors }; } // Kahn's algorithm — also doubles as cycle detection: if it can't // consume every step, the leftover steps form a cycle. var stepsById = steps.ToDictionary(s => s.FlowStepId); var queue = new Queue(roots.Select(r => r.FlowStepId)); var sortedSteps = new List(steps.Count); var remaining = new Dictionary(inDegree); while (queue.Count > 0) { int current = queue.Dequeue(); sortedSteps.Add(stepsById[current]); foreach (var next in adjacency[current].OrderBy(n => n)) { remaining[next]--; if (remaining[next] == 0) queue.Enqueue(next); } } if (sortedSteps.Count != steps.Count) { var stuck = steps.Where(s => !sortedSteps.Contains(s)).Select(s => s.NodeCode); errors.Add($"Flow graph contains a cycle involving step(s): {string.Join(", ", stuck)}."); return new DagValidationResult { Errors = errors }; } return new DagValidationResult { SortedSteps = sortedSteps }; } private static string ComputeHash(string content) { var bytes = SHA256.HashData(Encoding.UTF8.GetBytes(content)); return Convert.ToHexString(bytes).ToLowerInvariant(); } private sealed class DagValidationResult { public List Errors { get; set; } = new(); public List SortedSteps { get; set; } = new(); } } // ───────────────────────────────────────────────────────────── // IGopPublishService // ───────────────────────────────────────────────────────────── public interface IGopPublishService { Task PublishAsync( int flowId, string snapshotVersionLabel, LoginDTO loginDTO, CancellationToken ct = default); } // ───────────────────────────────────────────────────────────── // GopPublishResultDTO // ───────────────────────────────────────────────────────────── public sealed class GopPublishResultDTO { public bool IsSuccess { get; init; } public int SnapshotId { get; init; } public string SnapshotVersion { get; init; } = string.Empty; public IReadOnlyList Errors { get; init; } = []; public static GopPublishResultDTO Ok(int snapshotId, string version) => new() { IsSuccess = true, SnapshotId = snapshotId, SnapshotVersion = version }; public static GopPublishResultDTO Failure(IReadOnlyList errors) => new() { IsSuccess = false, Errors = errors }; } }