using System; using System.Collections.Generic; using System.Data.Common; 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 GopQueueDAL : IGopQueueDAL { private readonly IQueryExecutor _QueryExecutor; private readonly IValidation _Validation; public GopQueueDAL(IQueryExecutor QueryExecutor, IValidation Validation) { _QueryExecutor = QueryExecutor; _Validation = Validation; } // ── Status code translation ──────────────────────────────────────────────── // Live schema stores these as TINYINT codes (see GopQueueQB.cs header comment // for the full mapping and its source — GOP_DDL.sql's own column comments). // Business logic above this layer (GopExecutionPipeline, GopQueueBLL, // GopWorkerService) keeps using the string values below unchanged. private static byte MapExecutionStatus(string status) => status switch { "Created" => 0, "Queued" => 0, // no distinct DDL code — reuses Created "Running" => 1, "Success" => 2, "Failed" => 3, "PendingApproval" => 4, _ => throw new ArgumentException($"Unknown execution status '{status}'.") }; private static byte MapQueuedStatus(string status) => status switch { "Pending" => 0, "Processing" => 1, "Completed" => 2, "Failed" => 3, "Poison" => 3, // no distinct DDL code — reuses Failed _ => throw new ArgumentException($"Unknown queue status '{status}'.") }; private static byte MapNodeStatus(string status) => status switch { "Running" => 1, "Success" => 2, "Failed" => 3, _ => throw new ArgumentException($"Unknown node status '{status}'.") }; private static byte MapRemediationStatus(string status) => status switch { "Open" => 0, "InProgress" => 1, "Resolved" => 2, "Closed" => 3, _ => throw new ArgumentException($"Unknown remediation status '{status}'.") }; // ── Submit (Queue + Header atomic) ──────────────────────── public async Task<(int QueueId, int ExecutionId)> SubmitToQueue( GopExecutionQueueDTO QueueDTO, GopExecutionHeaderDTO HeaderDTO, LoginDTO LoginDTO) { DbTransaction? tx = null; try { tx = await _QueryExecutor.BeginTransactionAsync(LoginDTO); var now = DateTime.UtcNow; // TGOPEXECUTIONQUEUE/TGOPEXECUTIONHEADER have IDENTITY PKs on the live // schema — insert without the ID column, read it back via SCOPE_IDENTITY(). int queueId = await _QueryExecutor.ExecuteIdentityAsync( LoginDTO, GopQueueQB.INSERT_QUEUE, new { QueueDTO.ClientId, QueueDTO.FlowId, QueueDTO.FlowCode, QueueDTO.SnapshotVersion, QueueDTO.SourceCode, QueueDTO.PayloadJson, QueueDTO.IdempotencyKey, QueueDTO.Priority, QueueDTO.MaxRetry, CreatedById = LoginDTO.UserId, CreatedOn = now, ModifiedById = LoginDTO.UserId, ModifiedOn = now }, tx); QueueDTO.QueueId = queueId; HeaderDTO.QueueId = queueId; int execId = await _QueryExecutor.ExecuteIdentityAsync( LoginDTO, GopQueueQB.INSERT_EXECUTION_HEADER, new { HeaderDTO.ClientId, HeaderDTO.QueueId, ParentExecutionId = HeaderDTO.ParentExecutionId ?? -1, HeaderDTO.FlowId, HeaderDTO.FlowCode, HeaderDTO.SnapshotVersion, HeaderDTO.SourcePayloadJson, HeaderDTO.IsReplay, CreatedById = LoginDTO.UserId, CreatedOn = now, ModifiedById = LoginDTO.UserId, ModifiedOn = now }, tx); HeaderDTO.ExecutionId = execId; // Initial state log: null → Created await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.INSERT_STATE_LOG, new { HeaderDTO.ClientId, ExecutionId = execId, PreviousStatus = (byte?)null, NewStatus = MapExecutionStatus("Created"), ReasonCode = "EXECUTION_CREATED", ChangedBy = "gop-api" }, tx); await _QueryExecutor.CommitAsync(tx); return (queueId, execId); } catch (Exception ex) { if (tx != null) await _QueryExecutor.RollbackAsync(tx); throw new Exception(await _Validation.HandleException(ex, ErrorResponse.SaveErrorMessage)); } } // ── Queue Status ────────────────────────────────────────── public async Task UpdateQueueStatus(int ClientId, int QueueId, string Status, string? ErrorMessage, LoginDTO LoginDTO) { try { var now = DateTime.UtcNow; await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.UPDATE_QUEUE_STATUS, new { ClientId, QueueId, QueuedStatus = MapQueuedStatus(Status), ErrorMessage, ModifiedById = LoginDTO.UserId, ModifiedOn = now }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.UpdateErrorMessage)); } } // ── Execution Header ────────────────────────────────────── public async Task UpdateExecutionStatus(int ClientId, int ExecutionId, string Status, string? ErrorMessage, long? TotalDurationMs, LoginDTO LoginDTO, string? LastOutputPayloadJson = null) { try { var now = DateTime.UtcNow; await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.UPDATE_EXECUTION_STATUS, new { ClientId, ExecutionId, ExecutionStatus = MapExecutionStatus(Status), ErrorMessage, TotalDurationMs, LastOutputPayloadJson, ModifiedById = LoginDTO.UserId, ModifiedOn = now }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.UpdateErrorMessage)); } } public async Task UpdateExecutionCurrentNode(int ClientId, int ExecutionId, int NodeId, string NodeCode, LoginDTO LoginDTO) { try { await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.UPDATE_EXECUTION_CURRENT_NODE, new { ClientId, ExecutionId, CurrentNodeId = NodeId, CurrentNodeCode = NodeCode }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.UpdateErrorMessage)); } } public async Task UpdateExecutionApprovalPending( int ClientId, int ExecutionId, int NodeId, string NodeCode, string? LastOutputPayloadJson, string? AssignedRole, int? AssignedToUserId, LoginDTO LoginDTO) { try { await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.UPDATE_EXECUTION_APPROVAL_PENDING, new { ClientId, ExecutionId, CurrentNodeId = NodeId, CurrentNodeCode = NodeCode, LastOutputPayloadJson, AssignedRole, AssignedToUserId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.UpdateErrorMessage)); } } public async Task ResumeExecutionAfterApproval(int ClientId, int ExecutionId, LoginDTO LoginDTO) { try { int rows = await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.RESUME_EXECUTION_AFTER_APPROVAL, new { ClientId, ExecutionId }); return rows > 0; } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.UpdateErrorMessage)); } } public async Task GetExecutionById(int ClientId, int ExecutionId, LoginDTO LoginDTO) { try { return await _QueryExecutor.QuerySingleAsync( LoginDTO, GopQueueQB.GET_EXECUTION_BY_ID, new { ClientId, ExecutionId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } public async Task> GetExecutionsByStatus(int ClientId, string? Status, int Page, int PageSize, LoginDTO LoginDTO) { try { int offset = (Page - 1) * PageSize; byte? executionStatus = Status is null ? null : MapExecutionStatus(Status); return await _QueryExecutor.QueryAsync( LoginDTO, GopQueueQB.GET_EXECUTIONS_BY_STATUS, new { ClientId, ExecutionStatus = executionStatus, Offset = offset, PageSize }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } // ── Node Log ────────────────────────────────────────────── public async Task InsertNodeLog(GopExecutionNodeLogDTO NodeLogDTO, LoginDTO LoginDTO) { try { return await _QueryExecutor.ExecuteIdentityAsync( LoginDTO, GopQueueQB.INSERT_NODE_LOG, new { NodeLogDTO.ClientId, NodeLogDTO.ExecutionId, NodeLogDTO.NodeId, NodeLogDTO.NodeCode, NodeLogDTO.StepOrder, NodeLogDTO.AttemptNo, NodeLogDTO.RequestPayloadJson, NodeLogDTO.HttpMethod, NodeLogDTO.EndpointUrl }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.SaveErrorMessage)); } } public async Task UpdateNodeLog(int ClientId, int NodeLogId, string Status, string? ResponsePayload, int? HttpStatusCode, long DurationMs, string? ErrorCode, string? ErrorMessage, LoginDTO LoginDTO) { try { await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.UPDATE_NODE_LOG, new { ClientId, NodeLogId, NodeStatus = MapNodeStatus(Status), ResponsePayloadJson = ResponsePayload, HttpStatusCode, DurationMs, ErrorCode, ErrorMessage }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.UpdateErrorMessage)); } } public async Task> GetNodeLogsByExecution(int ClientId, int ExecutionId, LoginDTO LoginDTO) { try { return await _QueryExecutor.QueryAsync( LoginDTO, GopQueueQB.GET_NODE_LOGS_BY_EXECUTION, new { ClientId, ExecutionId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } // ── State Log ───────────────────────────────────────────── public async Task InsertStateLog(int ClientId, int ExecutionId, string? PreviousStatus, string NewStatus, string? ReasonCode, string ChangedBy, LoginDTO LoginDTO) { try { await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.INSERT_STATE_LOG, new { ClientId, ExecutionId, PreviousStatus = PreviousStatus is null ? (byte?)null : MapExecutionStatus(PreviousStatus), NewStatus = MapExecutionStatus(NewStatus), ReasonCode, ChangedBy }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.SaveErrorMessage)); } } public async Task> GetStateLogsByExecution(int ClientId, int ExecutionId, LoginDTO LoginDTO) { try { return await _QueryExecutor.QueryAsync( LoginDTO, GopQueueQB.GET_STATE_LOGS_BY_EXECUTION, new { ClientId, ExecutionId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } // ── Metrics ─────────────────────────────────────────────── public async Task InsertExecutionMetrics(GopExecutionMetricsDTO MetricsDTO, LoginDTO LoginDTO) { try { await _QueryExecutor.ExecuteAsync(LoginDTO, GopQueueQB.INSERT_EXECUTION_METRICS, MetricsDTO); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.SaveErrorMessage)); } } public async Task InsertNodeMetrics(GopExecutionNodeMetricsDTO NodeMetricsDTO, LoginDTO LoginDTO) { try { await _QueryExecutor.ExecuteAsync(LoginDTO, GopQueueQB.INSERT_NODE_METRICS, NodeMetricsDTO); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.SaveErrorMessage)); } } // ── Dead Letter ─────────────────────────────────────────── public async Task InsertDeadLetter(GopDeadLetterDTO DlqDTO, LoginDTO LoginDTO) { DbTransaction? tx = null; try { tx = await _QueryExecutor.BeginTransactionAsync(LoginDTO); var now = DateTime.UtcNow; // TGOPDEADLETTER has an IDENTITY PK on the live schema. int dlqId = await _QueryExecutor.ExecuteIdentityAsync(LoginDTO, GopQueueQB.INSERT_DEAD_LETTER, new { DlqDTO.ClientId, DlqDTO.ExecutionId, DlqDTO.QueueId, DlqDTO.FlowId, DlqDTO.FlowCode, DlqDTO.SnapshotVersion, DlqDTO.FailedNodeId, DlqDTO.FailedNodeCode, DlqDTO.PayloadJson, DlqDTO.ErrorMessage, DlqDTO.ErrorStack, DlqDTO.RetryCount, DlqDTO.FailureCategory, CreatedById = LoginDTO.UserId, CreatedOn = now, ModifiedById = LoginDTO.UserId, ModifiedOn = now }, tx); DlqDTO.DlqId = dlqId; await _QueryExecutor.CommitAsync(tx); return dlqId; } catch (Exception ex) { if (tx != null) await _QueryExecutor.RollbackAsync(tx); throw new Exception(await _Validation.HandleException(ex, ErrorResponse.SaveErrorMessage)); } } public async Task> GetDeadLetters(int ClientId, LoginDTO LoginDTO) { try { return await _QueryExecutor.QueryAsync( LoginDTO, GopQueueQB.GET_DEAD_LETTERS, new { ClientId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } public async Task ResolveDeadLetter(int ClientId, int DlqId, string ResolvedBy, LoginDTO LoginDTO) { try { // TGOPDEADLETTER.RESOLVEDBYID is an INT FK — attribution always reflects the // authenticated caller (LoginDTO.UserId), not the free-text ResolvedBy string // (kept on the interface for the caller's own audit trail / display purposes). await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.RESOLVE_DEAD_LETTER, new { ClientId, DlqId, ResolvedById = LoginDTO.UserId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.UpdateErrorMessage)); } } // ── Remediation ─────────────────────────────────────────── public async Task InsertRemediation(GopRemediationDTO RemediationDTO, LoginDTO LoginDTO) { DbTransaction? tx = null; try { tx = await _QueryExecutor.BeginTransactionAsync(LoginDTO); var now = DateTime.UtcNow; // TGOPREMEDIATION has an IDENTITY PK on the live schema. int remId = await _QueryExecutor.ExecuteIdentityAsync(LoginDTO, GopQueueQB.INSERT_REMEDIATION, new { RemediationDTO.ClientId, RemediationDTO.ExecutionId, RemediationDTO.QueueId, RemediationDTO.FlowId, RemediationDTO.FlowCode, RemediationDTO.SnapshotVersion, RemediationDTO.FailedNodeId, RemediationDTO.FailedNodeCode, RemediationDTO.OriginalPayloadJson, RemediationDTO.TargetResponseJson, RemediationDTO.IssueListJson, RemediationDTO.AssignedTo, CreatedById = LoginDTO.UserId, CreatedOn = now, ModifiedById = LoginDTO.UserId, ModifiedOn = now }, tx); RemediationDTO.RemediationId = remId; await _QueryExecutor.CommitAsync(tx); return remId; } catch (Exception ex) { if (tx != null) await _QueryExecutor.RollbackAsync(tx); throw new Exception(await _Validation.HandleException(ex, ErrorResponse.SaveErrorMessage)); } } public async Task UpdateRemediation(int ClientId, int RemediationId, string Status, string ResolutionNotes, string ResolvedBy, LoginDTO LoginDTO) { try { // See ResolveDeadLetter — RESOLVEDBYID is an INT FK, so attribution uses // the authenticated caller rather than the free-text ResolvedBy string. await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.UPDATE_REMEDIATION, new { ClientId, RemediationId, RemediationStatus = MapRemediationStatus(Status), ResolutionNotes, ResolvedById = LoginDTO.UserId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.UpdateErrorMessage)); } } public async Task> GetOpenRemediations(int ClientId, LoginDTO LoginDTO) { try { return await _QueryExecutor.QueryAsync( LoginDTO, GopQueueQB.GET_OPEN_REMEDIATIONS, new { ClientId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } // ── Lock ────────────────────────────────────────────────── public async Task AcquireLock(int ClientId, int ExecutionId, string LockedBy, int LockDurationSeconds, LoginDTO LoginDTO) { try { int result = await _QueryExecutor.ExecuteScalarAsync( LoginDTO, GopQueueQB.ACQUIRE_LOCK, new { ClientId, ExecutionId, LockedBy, LockDurationSeconds }); return result == 1; } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.SaveErrorMessage)); } } public async Task ReleaseLock(int ClientId, int ExecutionId, LoginDTO LoginDTO) { try { await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.RELEASE_LOCK, new { ClientId, ExecutionId }); } catch (Exception ex) { // Non-critical - log silently _ = ex; } } // ── Worker polling ──────────────────────────────────────── public async Task> GetPendingExecutions( int limit, LoginDTO loginDTO) { return await _QueryExecutor.QueryAsync( loginDTO, GopQueueQB.GET_PENDING_EXECUTIONS, new { Limit = limit }); } public async Task ResetExpiredLocks(LoginDTO loginDTO) { try { await _QueryExecutor.ExecuteAsync(loginDTO, GopQueueQB.RESET_EXPIRED_LOCKS, new { }); } catch (Exception ex) { // Recovery sweep failures are non-critical — log silently rather than crashing the worker. _ = ex; } } // ── Dashboard ───────────────────────────────────────────── public async Task> GetQueueDashboard(int ClientId, LoginDTO LoginDTO) { try { return await _QueryExecutor.QueryAsync( LoginDTO, GopQueueQB.GET_QUEUE_DASHBOARD, new { ClientId }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.GetFetchingMessge)); } } // ── Source File ─────────────────────────────────────────── public async Task InsertSourceFile(GopSourceFileDTO FileDTO, LoginDTO LoginDTO) { try { var now = DateTime.UtcNow; return await _QueryExecutor.ExecuteIdentityAsync(LoginDTO, GopQueueQB.INSERT_SOURCE_FILE, new { FileDTO.ClientId, FileDTO.SourceCode, FileDTO.OriginalFileName, FileDTO.StoredFilePath, FileDTO.FileSizeBytes, FileDTO.MimeType, FileDTO.UploadedBy, FileDTO.QueueId, CreatedById = LoginDTO.UserId, CreatedOn = now, ModifiedById = LoginDTO.UserId, ModifiedOn = now }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.SaveErrorMessage)); } } public async Task UpdateSourceFileStatus(int ClientId, int FileId, string Status, int? QueueId, LoginDTO LoginDTO) { try { byte uploadStatus = Status switch { "Uploaded" => 0, "Queued" => 1, "PartialFailed" => 2, "Failed" => 3, "ParseFailed" => 4, _ => throw new ArgumentException($"Unknown source file status '{Status}'.") }; await _QueryExecutor.ExecuteAsync( LoginDTO, GopQueueQB.UPDATE_SOURCE_FILE_STATUS, new { ClientId, FileId, UploadStatus = uploadStatus, QueueId, ModifiedById = LoginDTO.UserId, ModifiedOn = DateTime.UtcNow }); } catch (Exception ex) { throw new Exception(await _Validation.HandleException(ex, ErrorResponse.UpdateErrorMessage)); } } } }