using System; using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; using FrameworkDAL.CustomCode.GOP; using GB5Shared.DaprCache; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.GOP; using GB5Shared.EventLogPublish; using GB5Shared.GenerateAutoNumber; using GB5Shared.QueryExecutor; using GB5Shared.Resource.Response; using GB5Shared.Telemetry; using GB5Shared.Validation; using Microsoft.Extensions.Logging; using static GB5Shared.GB5Constant.Constant; namespace FrameworkBLL.GOP { public class GopFlowBLL : IGopFlowBLL { private readonly IGopFlowDAL _GopFlowDAL; private readonly IGopPublishService _PublishService; private readonly IQueryExecutor _QueryExecutor; private readonly AutoNumber _AutoNumber; private readonly IValidation _Validation; private readonly EventLogPublish _EventLog; private readonly KeyInvalidate _KeyInvalidate; private readonly ILogger _Logger; public GopFlowBLL( IGopFlowDAL gopFlowDAL, IGopPublishService publishService, IQueryExecutor queryExecutor, AutoNumber autoNumber, IValidation validation, EventLogPublish eventLog, KeyInvalidate keyInvalidate, ILogger logger) { _GopFlowDAL = gopFlowDAL; _PublishService = publishService; _QueryExecutor = queryExecutor; _AutoNumber = autoNumber; _Validation = validation; _EventLog = eventLog; _KeyInvalidate = keyInvalidate; _Logger = logger; } // ── Flow management ─────────────────────────────────────────────────────── public async Task SaveGopFlow( GopFlowDTO dto, LoginDTO loginDTO, CancellationToken ct = default) { await _Validation.NotEmpty(dto.FlowCode, nameof(dto.FlowCode)); await _Validation.NotEmpty(dto.FlowName, nameof(dto.FlowName)); bool isNew = dto.FlowId == 0; GB5Trace.Step("save-gop-flow", new { dto.FlowCode, isNew }); var trans = await _QueryExecutor.BeginTransactionAsync(loginDTO); try { if (isNew) { var auto = await _AutoNumber.GetNumberAsync(1, AUTONUMBERCONSTANT.GOPFLOW, loginDTO); dto.FlowId = auto.StartNumber; dto.Version = 1; dto.Status = "1"; await _GopFlowDAL.SaveFlow(dto, loginDTO, trans, ct); } else { dto.Version++; await _GopFlowDAL.UpdateFlow(dto, loginDTO, trans, ct); } await _QueryExecutor.CommitAsync(trans); await _KeyInvalidate.InvalidateByKey( dto.FlowId, EntityConstant.GOPFLOW, CacheKeyLevel.CLIENT_LEVEL, loginDTO); await _KeyInvalidate.InvalidateByKey( loginDTO.ClientId, EntityConstant.GOPFLOW, CacheKeyLevel.CLIENT_LEVEL, loginDTO); _Logger.LogInformation( "GopFlowBLL: flow {FlowCode} (Id={FlowId}) {Action} by user {UserId}", dto.FlowCode, dto.FlowId, isNew ? "created" : "updated", loginDTO.UserId); return isNew ? $"{SuccessResponse.SaveSuccessMessage} {dto.FlowId}" : $"{SuccessResponse.UpdateSuccessMessage} {dto.FlowId}"; } catch (Exception ex) { await _QueryExecutor.RollbackAsync(trans); GB5Trace.MarkFailed("save-gop-flow-failed", ex); throw; } } public async Task DeleteGopFlow( int flowId, LoginDTO loginDTO, CancellationToken ct = default) { GB5Trace.Step("delete-gop-flow", new { flowId }); var trans = await _QueryExecutor.BeginTransactionAsync(loginDTO); try { await _GopFlowDAL.DeleteFlow(flowId, loginDTO, trans, ct); await _QueryExecutor.CommitAsync(trans); await _KeyInvalidate.InvalidateByKey( flowId, EntityConstant.GOPFLOW, CacheKeyLevel.CLIENT_LEVEL, loginDTO); await _KeyInvalidate.InvalidateByKey( loginDTO.ClientId, EntityConstant.GOPFLOW, CacheKeyLevel.CLIENT_LEVEL, loginDTO); _Logger.LogInformation( "GopFlowBLL: flow {FlowId} deleted by user {UserId}", flowId, loginDTO.UserId); return SuccessResponse.DeleteSuccessMessage; } catch (Exception ex) { await _QueryExecutor.RollbackAsync(trans); GB5Trace.MarkFailed("delete-gop-flow-failed", ex); throw; } } // ── Flow Step / Edge (design-time graph) management ─────────────────────── public async Task GetGopFlowGraph( int flowId, LoginDTO loginDTO, CancellationToken ct = default) { var flow = await _GopFlowDAL.GetFlowById(loginDTO.ClientId, flowId, loginDTO); var steps = await _GopFlowDAL.GetFlowSteps(flowId, loginDTO, ct); var edges = await _GopFlowDAL.GetFlowStepEdges(flowId, loginDTO, ct); return new GopFlowGraphDTO { Flow = flow ?? new GopFlowDTO { FlowId = flowId }, Steps = steps.ToList(), Edges = edges.ToList() }; } public async Task SaveGopFlowStep( GopFlowStepDTO dto, LoginDTO loginDTO, CancellationToken ct = default) { await _Validation.NotEmpty(dto.NodeCode, nameof(dto.NodeCode)); await _Validation.NotEmpty(dto.NodeName, nameof(dto.NodeName)); await _Validation.NotEmpty(dto.NodeType, nameof(dto.NodeType)); if (dto.FlowId == 0) throw new ArgumentException("FlowId is required."); bool isNew = dto.FlowStepId == 0; GB5Trace.Step("save-gop-flow-step", new { dto.FlowId, dto.NodeCode, dto.NodeType, isNew }); var trans = await _QueryExecutor.BeginTransactionAsync(loginDTO); try { if (isNew) { var auto = await _AutoNumber.GetNumberAsync(1, AUTONUMBERCONSTANT.GOPFLOWSTEP, loginDTO); dto.FlowStepId = auto.StartNumber; await _GopFlowDAL.InsertFlowStep(dto, loginDTO, trans, ct); } else { await _GopFlowDAL.UpdateFlowStep(dto, loginDTO, trans, ct); } await _QueryExecutor.CommitAsync(trans); await _KeyInvalidate.InvalidateByKey( dto.FlowId, EntityConstant.GOPFLOWSTEP, CacheKeyLevel.CLIENT_LEVEL, loginDTO); await _EventLog.PublishEventLogAsync( isNew ? "GOP Flow Step Created" : "GOP Flow Step Updated", dto, EventTypeConstant.GOPFLOWSTEPSAVEDEVENTTYPEID, dto.FlowStepId, loginDTO, ct: ct); return isNew ? $"{SuccessResponse.SaveSuccessMessage} {dto.FlowStepId}" : $"{SuccessResponse.UpdateSuccessMessage} {dto.FlowStepId}"; } catch (Exception ex) { await _QueryExecutor.RollbackAsync(trans); GB5Trace.MarkFailed("save-gop-flow-step-failed", ex); throw; } } public async Task GetGopFlowStep( int flowStepId, LoginDTO loginDTO, CancellationToken ct = default) { return await _GopFlowDAL.GetFlowStepById(flowStepId, loginDTO, ct); } public async Task> GetGopFlowStepList( int flowId, LoginDTO loginDTO, CancellationToken ct = default) { return await _GopFlowDAL.GetFlowSteps(flowId, loginDTO, ct); } public async Task DeleteGopFlowStep( int flowStepId, LoginDTO loginDTO, CancellationToken ct = default) { GB5Trace.Step("delete-gop-flow-step", new { flowStepId }); // Resolved before delete — needed to invalidate the Graph cache (keyed by FlowId), not // available from the delete call's own parameters. var existing = await _GopFlowDAL.GetFlowStepById(flowStepId, loginDTO, ct); var trans = await _QueryExecutor.BeginTransactionAsync(loginDTO); try { await _GopFlowDAL.DeleteFlowStep(flowStepId, loginDTO, trans, ct); await _QueryExecutor.CommitAsync(trans); if (existing != null) { await _KeyInvalidate.InvalidateByKey( existing.FlowId, EntityConstant.GOPFLOWSTEP, CacheKeyLevel.CLIENT_LEVEL, loginDTO); } return SuccessResponse.DeleteSuccessMessage; } catch (Exception ex) { await _QueryExecutor.RollbackAsync(trans); GB5Trace.MarkFailed("delete-gop-flow-step-failed", ex); throw; } } public async Task SaveGopFlowStepEdge( GopFlowStepEdgeDTO dto, LoginDTO loginDTO, CancellationToken ct = default) { if (dto.FlowId == 0) throw new ArgumentException("FlowId is required."); if (dto.FromStepId == 0) throw new ArgumentException("FromStepId is required."); if (dto.ToStepId == 0) throw new ArgumentException("ToStepId is required."); if (dto.FromStepId == dto.ToStepId) throw new ArgumentException("A step cannot have an edge to itself."); bool isNew = dto.FlowStepEdgeId == 0; GB5Trace.Step("save-gop-flow-step-edge", new { dto.FlowId, dto.FromStepId, dto.ToStepId, dto.EdgeCondition, isNew }); var trans = await _QueryExecutor.BeginTransactionAsync(loginDTO); try { if (isNew) { var auto = await _AutoNumber.GetNumberAsync(1, AUTONUMBERCONSTANT.GOPFLOWSTEPEDGE, loginDTO); dto.FlowStepEdgeId = auto.StartNumber; await _GopFlowDAL.InsertFlowStepEdge(dto, loginDTO, trans, ct); } else { await _GopFlowDAL.UpdateFlowStepEdge(dto, loginDTO, trans, ct); } await _QueryExecutor.CommitAsync(trans); await _KeyInvalidate.InvalidateByKey( dto.FlowId, EntityConstant.GOPFLOWSTEP, CacheKeyLevel.CLIENT_LEVEL, loginDTO); return isNew ? $"{SuccessResponse.SaveSuccessMessage} {dto.FlowStepEdgeId}" : $"{SuccessResponse.UpdateSuccessMessage} {dto.FlowStepEdgeId}"; } catch (Exception ex) { await _QueryExecutor.RollbackAsync(trans); GB5Trace.MarkFailed("save-gop-flow-step-edge-failed", ex); throw; } } public async Task GetGopFlowStepEdge( int flowStepEdgeId, LoginDTO loginDTO, CancellationToken ct = default) { return await _GopFlowDAL.GetFlowStepEdgeById(flowStepEdgeId, loginDTO, ct); } public async Task> GetGopFlowStepEdgeList( int flowId, LoginDTO loginDTO, CancellationToken ct = default) { return await _GopFlowDAL.GetFlowStepEdges(flowId, loginDTO, ct); } public async Task DeleteGopFlowStepEdge( int flowStepEdgeId, LoginDTO loginDTO, CancellationToken ct = default) { GB5Trace.Step("delete-gop-flow-step-edge", new { flowStepEdgeId }); // Resolved before delete — needed to invalidate the Graph cache (keyed by FlowId), not // available from the delete call's own parameters. var existing = await _GopFlowDAL.GetFlowStepEdgeById(flowStepEdgeId, loginDTO, ct); var trans = await _QueryExecutor.BeginTransactionAsync(loginDTO); try { await _GopFlowDAL.DeleteFlowStepEdge(flowStepEdgeId, loginDTO, trans, ct); await _QueryExecutor.CommitAsync(trans); if (existing != null) { await _KeyInvalidate.InvalidateByKey( existing.FlowId, EntityConstant.GOPFLOWSTEP, CacheKeyLevel.CLIENT_LEVEL, loginDTO); } return SuccessResponse.DeleteSuccessMessage; } catch (Exception ex) { await _QueryExecutor.RollbackAsync(trans); GB5Trace.MarkFailed("delete-gop-flow-step-edge-failed", ex); throw; } } public async Task GetGopFlow( int flowId, LoginDTO loginDTO, CancellationToken ct = default) { return await _GopFlowDAL.GetFlowById(loginDTO.ClientId, flowId, loginDTO); } public async Task> GetGopFlowList( LoginDTO loginDTO, CancellationToken ct = default) { return await _GopFlowDAL.GetFlowList(loginDTO, ct); } public async Task PublishGopFlow( 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 result = await _PublishService.PublishAsync(flowId, snapshotVersionLabel, loginDTO, ct); if (result.IsSuccess) { await _KeyInvalidate.InvalidateByKey( flowId, EntityConstant.GOPFLOWSNAPSHOT, CacheKeyLevel.CLIENT_LEVEL, loginDTO); } return result; } public async Task> GetGopFlowVersionHistory( int flowId, LoginDTO loginDTO, CancellationToken ct = default) { return await _GopFlowDAL.GetSnapshotsByFlow(loginDTO.ClientId, flowId, loginDTO); } // ── Source Binding management ───────────────────────────────────────────── public async Task SaveGopSourceBinding( GopSourceBindingDTO dto, LoginDTO loginDTO, CancellationToken ct = default) { await _Validation.NotEmpty(dto.SourceCode, nameof(dto.SourceCode)); await _Validation.NotEmpty(dto.SourceName, nameof(dto.SourceName)); await _Validation.NotEmpty(dto.SourceType, nameof(dto.SourceType)); bool isNew = dto.SourceBindingId == 0; var trans = await _QueryExecutor.BeginTransactionAsync(loginDTO); try { if (isNew) { var auto = await _AutoNumber.GetNumberAsync(1, AUTONUMBERCONSTANT.GOPSOURCEBINDING, loginDTO); dto.SourceBindingId = auto.StartNumber; dto.IsActive = true; await _GopFlowDAL.SaveSourceBinding(dto, loginDTO, trans, ct); } else { await _GopFlowDAL.UpdateSourceBinding(dto, loginDTO, trans, ct); } await _QueryExecutor.CommitAsync(trans); await _KeyInvalidate.InvalidateByKey( dto.SourceBindingId, EntityConstant.GOPSOURCEBINDING, CacheKeyLevel.CLIENT_LEVEL, loginDTO); await _KeyInvalidate.InvalidateByKey( loginDTO.ClientId, EntityConstant.GOPSOURCEBINDING, CacheKeyLevel.CLIENT_LEVEL, loginDTO); return isNew ? $"{SuccessResponse.SaveSuccessMessage} {dto.SourceBindingId}" : $"{SuccessResponse.UpdateSuccessMessage} {dto.SourceBindingId}"; } catch { await _QueryExecutor.RollbackAsync(trans); throw; } } public async Task DeleteGopSourceBinding( int sourceBindingId, LoginDTO loginDTO, CancellationToken ct = default) { var trans = await _QueryExecutor.BeginTransactionAsync(loginDTO); try { await _GopFlowDAL.DeleteSourceBinding(sourceBindingId, loginDTO, trans, ct); await _QueryExecutor.CommitAsync(trans); await _KeyInvalidate.InvalidateByKey( sourceBindingId, EntityConstant.GOPSOURCEBINDING, CacheKeyLevel.CLIENT_LEVEL, loginDTO); await _KeyInvalidate.InvalidateByKey( loginDTO.ClientId, EntityConstant.GOPSOURCEBINDING, CacheKeyLevel.CLIENT_LEVEL, loginDTO); return SuccessResponse.DeleteSuccessMessage; } catch { await _QueryExecutor.RollbackAsync(trans); throw; } } public async Task GetGopSourceBinding( int sourceBindingId, LoginDTO loginDTO, CancellationToken ct = default) { return await _GopFlowDAL.GetSourceBindingById(sourceBindingId, loginDTO, ct); } public async Task> GetGopSourceBindingList( LoginDTO loginDTO, CancellationToken ct = default) { return await _GopFlowDAL.GetAllSourceBindings(loginDTO.ClientId, loginDTO); } // ── Read-only (used internally / by execution) ──────────────────────────── public async Task> GetAllSourceBindings(LoginDTO LoginDTO) { return await _GopFlowDAL.GetAllSourceBindings(LoginDTO.ClientId, LoginDTO); } public async Task GetFlowSnapshot( int flowId, string snapshotVersion, LoginDTO loginDTO) { return await _GopFlowDAL.GetSnapshotByFlowVersion(loginDTO.ClientId, flowId, snapshotVersion, loginDTO); } public async Task> GetSnapshotSteps( string snapshotVersion, LoginDTO loginDTO) { return await _GopFlowDAL.GetSnapshotStepsByVersion(loginDTO.ClientId, snapshotVersion, loginDTO); } } }