using System; using System.Collections.Generic; using System.Data.Common; using System.Threading; using System.Threading.Tasks; using FrameworkDAL.Query.GOP; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.GOP; using GB5Shared.QueryExecutor; using GB5Shared.Resource.Response; using GB5Shared.Validation; namespace FrameworkDAL.CustomCode.GOP { public class GopFlowDAL : IGopFlowDAL { private readonly IQueryExecutor _QueryExecutor; private readonly IValidation _Validation; public GopFlowDAL(IQueryExecutor QueryExecutor, IValidation Validation) { _QueryExecutor = QueryExecutor; _Validation = Validation; } // ── Source Binding ──────────────────────────────────────────────────────── // Every query in GopFlowQB.cs requires @TenantId — there is no auto-population // by IQueryExecutor (confirmed by reading QueryExecutor.ExecuteAsync directly; // the old comment here claiming otherwise was wrong and caused every write in // this file to fail live with "Must declare the scalar variable @TenantId"). public async Task GetSourceBindingByCode( int ClientId, string SourceCode, string Environment, LoginDTO LoginDTO) { try { return await _QueryExecutor.QuerySingleAsync( LoginDTO, GopFlowQB.GET_SOURCE_BINDING_BY_CODE, new { TenantId = ClientId, SourceCode, Environment }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } public async Task GetSourceBindingById( int sourceBindingId, LoginDTO loginDTO, CancellationToken ct = default) { return await _QueryExecutor.QuerySingleAsync( loginDTO, GopFlowQB.GET_SOURCE_BINDING_BY_ID, new { TenantId = loginDTO.ClientId, SourceBindingId = sourceBindingId }, cancellationToken: ct); } public async Task> GetAllSourceBindings( int ClientId, LoginDTO LoginDTO) { try { return await _QueryExecutor.QueryAsync( LoginDTO, GopFlowQB.GET_ALL_SOURCE_BINDINGS, new { TenantId = ClientId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } public async Task SaveSourceBinding( GopSourceBindingDTO dto, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct = default) { var now = DateTime.UtcNow; await _QueryExecutor.ExecuteAsync(loginDTO, GopFlowQB.INSERT_SOURCE_BINDING, new { dto.SourceBindingId, TenantId = loginDTO.ClientId, dto.SourceType, dto.SourceCode, dto.SourceName, dto.Description, dto.FlowId, dto.SnapshotVersion, dto.Environment, IsActive = true, SortOrder = 0, Status = 1, CreatedById = loginDTO.UserId, CreatedOn = now, ModifiedById = loginDTO.UserId, ModifiedOn = now }, tx, ct); } public async Task UpdateSourceBinding( GopSourceBindingDTO dto, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct = default) { await _QueryExecutor.ExecuteAsync(loginDTO, GopFlowQB.UPDATE_SOURCE_BINDING, new { dto.SourceBindingId, TenantId = loginDTO.ClientId, dto.SourceName, dto.Description, dto.FlowId, dto.SnapshotVersion, dto.Environment, dto.IsActive, ModifiedById = loginDTO.UserId, ModifiedOn = DateTime.UtcNow }, tx, ct); } public async Task DeleteSourceBinding( int SourceBindingId, LoginDTO LoginDTO, DbTransaction tx, CancellationToken ct = default) { await _QueryExecutor.ExecuteAsync(LoginDTO, GopFlowQB.DELETE_SOURCE_BINDING, new { SourceBindingId, TenantId = LoginDTO.ClientId, ModifiedById = LoginDTO.UserId, ModifiedOn = DateTime.UtcNow }, tx, ct); } // ── Flow ────────────────────────────────────────────────────────────────── public async Task GetFlowById(int ClientId, int FlowId, LoginDTO LoginDTO) { try { return await _QueryExecutor.QuerySingleAsync( LoginDTO, GopFlowQB.GET_FLOW_BY_ID, new { TenantId = ClientId, FlowId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } public async Task> GetFlowList( LoginDTO loginDTO, CancellationToken ct = default) { return await _QueryExecutor.QueryAsync( loginDTO, GopFlowQB.GET_FLOW_LIST, new { TenantId = loginDTO.ClientId }, cancellationToken: ct); } public async Task SaveFlow( GopFlowDTO dto, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct = default) { var now = DateTime.UtcNow; await _QueryExecutor.ExecuteAsync(loginDTO, GopFlowQB.INSERT_FLOW, new { dto.FlowId, TenantId = loginDTO.ClientId, dto.FlowCode, dto.FlowName, dto.Description, Status = 1, // INT: Active Version = Math.Max(1, dto.Version), SortOrder = 0, CreatedById = loginDTO.UserId, CreatedOn = now, ModifiedById = loginDTO.UserId, ModifiedOn = now }, tx, ct); } public async Task UpdateFlow( GopFlowDTO dto, LoginDTO loginDTO, DbTransaction tx, CancellationToken ct = default) { await _QueryExecutor.ExecuteAsync(loginDTO, GopFlowQB.UPDATE_FLOW, new { dto.FlowId, TenantId = loginDTO.ClientId, dto.FlowName, dto.Description, Status = 1, dto.Version, SortOrder = 0, ModifiedById = loginDTO.UserId, ModifiedOn = DateTime.UtcNow }, tx, ct); } public async Task DeleteFlow( int FlowId, LoginDTO LoginDTO, DbTransaction tx, CancellationToken ct = default) { await _QueryExecutor.ExecuteAsync(LoginDTO, GopFlowQB.DELETE_FLOW, new { FlowId, TenantId = LoginDTO.ClientId, ModifiedById = LoginDTO.UserId, ModifiedOn = DateTime.UtcNow }, tx, ct); } // ── Snapshot ────────────────────────────────────────────────────────────── public async Task GetSnapshotByFlowVersion( int ClientId, int FlowId, string SnapshotVersion, LoginDTO LoginDTO) { try { return await _QueryExecutor.QuerySingleAsync( LoginDTO, GopFlowQB.GET_SNAPSHOT_BY_FLOW_VERSION, new { TenantId = ClientId, FlowId, SnapshotVersion }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } public async Task> GetSnapshotsByFlow( int ClientId, int FlowId, LoginDTO LoginDTO) { try { return await _QueryExecutor.QueryAsync( LoginDTO, GopFlowQB.GET_SNAPSHOTS_BY_FLOW, new { TenantId = ClientId, FlowId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } // ── Snapshot Steps ──────────────────────────────────────────────────────── public async Task> GetSnapshotStepsByVersion( int ClientId, string SnapshotVersion, LoginDTO LoginDTO) { try { return await _QueryExecutor.QueryAsync( LoginDTO, GopFlowQB.GET_SNAPSHOT_STEPS_BY_VERSION, new { TenantId = ClientId, SnapshotVersion }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } // TenantId is the execution's OWN tenant (the explicit ClientId param), never // LoginDTO.ClientId — GopWorkerService's background poller connects using a // per-connection LoginDTO built from MSERVERCONFIG (the DB's real business // ClientId, e.g. GB5DEMO's -1399999868), which is NOT the same tenant an // individual queued execution belongs to (e.g. Promotion's framework-level // ClientId=-1 executions). Using LoginDTO.ClientId here silently returned zero // steps for any such execution — the pipeline then took its "no steps found" // branch and marked it Success without running anything. Found live 2026-08-21 // verifying Metadata Promotion's first-ever real GOP Ingest execution. public async Task> GetSnapshotSteps( int ClientId, int FlowId, string SnapshotVersion, LoginDTO LoginDTO) { return await _QueryExecutor.QueryAsync( LoginDTO, GopFlowQB.GET_SNAPSHOT_STEPS, new { TenantId = ClientId, FlowId, SnapshotVersion }); } // See GetSnapshotSteps' header comment — same fix, same reason. public async Task> GetSnapshotStepEdges( int ClientId, int FlowId, string SnapshotVersion, LoginDTO LoginDTO, CancellationToken ct = default) { return await _QueryExecutor.QueryAsync( LoginDTO, GopFlowQB.GET_SNAPSHOT_STEP_EDGES_BY_VERSION, new { TenantId = ClientId, FlowId, SnapshotVersion }, cancellationToken: ct); } // ── Flow Step (design-time graph node) ────────────────────────────────────── public async Task> GetFlowSteps( int FlowId, LoginDTO LoginDTO, CancellationToken ct = default) { return await _QueryExecutor.QueryAsync( LoginDTO, GopFlowQB.GET_FLOW_STEPS_BY_FLOW, new { TenantId = LoginDTO.ClientId, FlowId }, cancellationToken: ct); } public async Task GetFlowStepById( int FlowStepId, LoginDTO LoginDTO, CancellationToken ct = default) { return await _QueryExecutor.QuerySingleAsync( LoginDTO, GopFlowQB.GET_FLOW_STEP_BY_ID, new { TenantId = LoginDTO.ClientId, FlowStepId }, cancellationToken: ct); } public async Task InsertFlowStep( GopFlowStepDTO dto, LoginDTO LoginDTO, DbTransaction tx, CancellationToken ct = default) { var now = DateTime.UtcNow; await _QueryExecutor.ExecuteAsync(LoginDTO, GopFlowQB.INSERT_FLOW_STEP, new { dto.FlowStepId, TenantId = LoginDTO.ClientId, dto.FlowId, dto.NodeCode, dto.NodeName, dto.NodeType, dto.LogicalServiceName, dto.TargetOperationId, dto.TimeoutSeconds, dto.MaxRetries, dto.SupportsRemediation, dto.AssignedRole, dto.AssignedToUserId, dto.SortOrder, CreatedById = LoginDTO.UserId, CreatedOn = now, ModifiedById = LoginDTO.UserId, ModifiedOn = now }, tx, ct); } public async Task UpdateFlowStep( GopFlowStepDTO dto, LoginDTO LoginDTO, DbTransaction tx, CancellationToken ct = default) { await _QueryExecutor.ExecuteAsync(LoginDTO, GopFlowQB.UPDATE_FLOW_STEP, new { dto.FlowStepId, TenantId = LoginDTO.ClientId, dto.NodeName, dto.NodeType, dto.LogicalServiceName, dto.TargetOperationId, dto.TimeoutSeconds, dto.MaxRetries, dto.SupportsRemediation, dto.AssignedRole, dto.AssignedToUserId, dto.SortOrder, ModifiedById = LoginDTO.UserId, ModifiedOn = DateTime.UtcNow }, tx, ct); } public async Task DeleteFlowStep( int FlowStepId, LoginDTO LoginDTO, DbTransaction tx, CancellationToken ct = default) { await _QueryExecutor.ExecuteAsync(LoginDTO, GopFlowQB.DELETE_FLOW_STEP, new { FlowStepId, TenantId = LoginDTO.ClientId, ModifiedById = LoginDTO.UserId, ModifiedOn = DateTime.UtcNow }, tx, ct); } // ── Flow Step Edge (design-time DAG edge) ─────────────────────────────────── public async Task> GetFlowStepEdges( int FlowId, LoginDTO LoginDTO, CancellationToken ct = default) { return await _QueryExecutor.QueryAsync( LoginDTO, GopFlowQB.GET_FLOW_STEP_EDGES_BY_FLOW, new { TenantId = LoginDTO.ClientId, FlowId }, cancellationToken: ct); } public async Task GetFlowStepEdgeById( int FlowStepEdgeId, LoginDTO LoginDTO, CancellationToken ct = default) { return await _QueryExecutor.QuerySingleAsync( LoginDTO, GopFlowQB.GET_FLOW_STEP_EDGE_BY_ID, new { TenantId = LoginDTO.ClientId, FlowStepEdgeId }, cancellationToken: ct); } public async Task InsertFlowStepEdge( GopFlowStepEdgeDTO dto, LoginDTO LoginDTO, DbTransaction tx, CancellationToken ct = default) { var now = DateTime.UtcNow; await _QueryExecutor.ExecuteAsync(LoginDTO, GopFlowQB.INSERT_FLOW_STEP_EDGE, new { dto.FlowStepEdgeId, TenantId = LoginDTO.ClientId, dto.FlowId, dto.FromStepId, dto.ToStepId, dto.EdgeLabel, dto.EdgeCondition, CreatedById = LoginDTO.UserId, CreatedOn = now, ModifiedById = LoginDTO.UserId, ModifiedOn = now }, tx, ct); } public async Task UpdateFlowStepEdge( GopFlowStepEdgeDTO dto, LoginDTO LoginDTO, DbTransaction tx, CancellationToken ct = default) { await _QueryExecutor.ExecuteAsync(LoginDTO, GopFlowQB.UPDATE_FLOW_STEP_EDGE, new { dto.FlowStepEdgeId, TenantId = LoginDTO.ClientId, dto.EdgeLabel, dto.EdgeCondition, ModifiedById = LoginDTO.UserId, ModifiedOn = DateTime.UtcNow }, tx, ct); } public async Task DeleteFlowStepEdge( int FlowStepEdgeId, LoginDTO LoginDTO, DbTransaction tx, CancellationToken ct = default) { await _QueryExecutor.ExecuteAsync(LoginDTO, GopFlowQB.DELETE_FLOW_STEP_EDGE, new { FlowStepEdgeId, TenantId = LoginDTO.ClientId, ModifiedById = LoginDTO.UserId, ModifiedOn = DateTime.UtcNow }, tx, ct); } // ── Publish ─────────────────────────────────────────────────────────────── // GopFlowSnapshotDTO/GopFlowSnapshotStepDTO/GopFlowSnapshotStepEdgeDTO all use // a ClientId property (not TenantId) — INSERT_SNAPSHOT and INSERT_SNAPSHOT_STEP's // SQL bind @TenantId, which Dapper cannot satisfy from a DTO instance that has no // TenantId property at all (passing the DTO directly, as this file used to do, // fails live with "Must declare the scalar variable @TenantId"). Built explicit // parameter objects mapping ClientId -> TenantId instead. INSERT_SNAPSHOT_STEP_EDGE // is the one query in this file that already binds @ClientId directly, so passing // its DTO through unchanged is correct — left as-is. public async Task InsertSnapshot( GopFlowSnapshotDTO SnapshotDTO, LoginDTO LoginDTO, DbTransaction tx) { return await _QueryExecutor.ExecuteAsync(LoginDTO, GopFlowQB.INSERT_SNAPSHOT, new { SnapshotDTO.SnapshotId, TenantId = SnapshotDTO.ClientId, SnapshotDTO.FlowId, SnapshotDTO.SnapshotVersion, SnapshotDTO.Environment, SnapshotDTO.SpecAggregateHash, SnapshotDTO.SnapshotJson, SnapshotDTO.IsActive, SnapshotDTO.PublishedAt, SnapshotDTO.PublishedBy }, tx); } public async Task InsertSnapshotStep( GopFlowSnapshotStepDTO StepDTO, LoginDTO LoginDTO, DbTransaction tx) { await _QueryExecutor.ExecuteAsync(LoginDTO, GopFlowQB.INSERT_SNAPSHOT_STEP, new { StepDTO.SnapshotStepId, TenantId = StepDTO.ClientId, StepDTO.SnapshotId, StepDTO.SnapshotVersion, StepDTO.StepOrder, StepDTO.NodeId, StepDTO.NodeCode, StepDTO.NodeType, StepDTO.LogicalServiceName, StepDTO.ResolvedEndpoint, StepDTO.HttpMethod, StepDTO.RelativePath, StepDTO.RequestSchemaJson, StepDTO.ResponseSchemaJson, StepDTO.MappingJson, StepDTO.SpecHash, StepDTO.TimeoutSeconds, StepDTO.MaxRetries, StepDTO.SupportsRemediation, StepDTO.AssignedRole, StepDTO.AssignedToUserId, CreatedById = LoginDTO.UserId }, tx); } public async Task InsertSnapshotStepEdge( GopFlowSnapshotStepEdgeDTO EdgeDTO, LoginDTO LoginDTO, DbTransaction tx) { await _QueryExecutor.ExecuteAsync( LoginDTO, GopFlowQB.INSERT_SNAPSHOT_STEP_EDGE, EdgeDTO, tx); } public async Task DeactivateSnapshotsByFlow( int FlowId, LoginDTO LoginDTO, DbTransaction tx) { await _QueryExecutor.ExecuteAsync( LoginDTO, GopFlowQB.DEACTIVATE_SNAPSHOTS_BY_FLOW, new { FlowId, TenantId = LoginDTO.ClientId, ModifiedById = LoginDTO.UserId }, tx); } // ── Qualifier Rules (legacy) ────────────────────────────────────────────── public async Task> GetQualifierRulesByCode( int ClientId, string QualifierCode, LoginDTO LoginDTO) { try { return await _QueryExecutor.QueryAsync( LoginDTO, GopFlowQB.GET_QUALIFIER_RULES_BY_CODE, new { TenantId = ClientId, QualifierCode }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } } }