using FrameworkDAL.DTO.SchedulerTaskGenerator; using FrameworkDAL.Query.SchedulerTaskGenerator; using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; using System; using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; namespace FrameworkDAL.CustomCode.SchedulerTaskGenerator { public interface ISchedulerTaskDAL { Task> LoadReadyJobsAsync( LoginDTO login, CancellationToken ct = default); Task LoadJobDetailsAsync( int jobExecutionId, LoginDTO login, CancellationToken ct = default); /// /// Minimal, single-column fallback lookup used when LoadJobDetailsAsync's joined /// result comes back with UriParameterValue empty despite MJOBDEFINE holding a /// stable, non-null value (observed intermittently — root cause not yet isolated; /// this guarantees correctness at the point of use regardless of that cause). /// Task GetUriParameterValueAsync( int jobId, int tenantId, LoginDTO login, CancellationToken ct = default); Task> LoadActionsByJobIdsAsync( IEnumerable jobIds, LoginDTO login, CancellationToken ct = default); /// /// True if this job has at least one active MACTION row linked via the scheduler's /// fixed EVENTTYPEID (SchedulerActionEventConstants.SchedulerJobEventTypeId). /// Task HasActiveSchedulerActionAsync( int jobId, LoginDTO login, CancellationToken ct = default); /// /// Inserts a new TJOBEXECUTION row with STATUS=Pending, due immediately, for the /// given JobId, and returns the newly generated JOBEXECUTIONID (via the SQL's /// OUTPUT clause). Actual delivery happens later when SchedulerExecutionPoller /// claims the row via ClaimDueExecutionsAsync — this method only records that /// Quartz fired the job. /// ✅ Returns 0 if no row was inserted — e.g. the ISCONCURRENT guard blocked it /// because a prior execution for this job is still Pending or InProgress. /// ✅ nowIst is passed so NEXTRUNON (the informational next-scheduled-run /// estimate) is computed and saved in the DB at insert time. /// ✅ triggerType/hostInstance populate TJOBEXECUTION.TRIGGERTYPE/HOSTINSTANCE; /// CORRELATIONID is generated inside the SQL itself (NEWID()). /// Task CreatePendingExecutionAsync( int jobId, DateTime nowIst, string triggerType, string hostInstance, string correlationId, string? traceParent, string? traceState, LoginDTO login, CancellationToken ct = default); /// /// Atomically claims up to due TJOBEXECUTION rows /// (STATUS='PENDING', NEXTATTEMPTON due) for one tenant, flips them to 'PROCESSING' /// (InProgress), and returns everything SchedulerExecutionDeliveryService needs /// to make the call — no separate LoadJobDetailsAsync round trip needed. Safe to /// call concurrently from multiple poller instances (ROWLOCK, READPAST). /// Task> ClaimDueExecutionsAsync( int batchSize, LoginDTO login, CancellationToken ct = default); /// /// Returns the current RETRYNUMBER for a JobExecutionId and the MAXRETRIES /// configured on the owning MJOBDEFINE row. Used to decide retry vs DLQ. /// Task<(int RetryNumber, int MaxRetries, int JobId)> GetRetryStateAsync( int jobExecutionId, LoginDTO login, CancellationToken ct = default); /// /// Applies a retry-or-DLQ decision to an existing TJOBEXECUTION row. /// newStatus: SchedulerExecutionStatus.Pending ("PENDING", requeue after backoff) /// or .DLQ ("DLQ", retries exhausted). /// Task ApplyRetryDecisionAsync( int jobExecutionId, int newRetryNumber, string newStatus, DateTime? nextAttemptOn, LoginDTO login, CancellationToken ct = default); /// /// Finds InProgress executions whose STARTRUNON exceeds the owning job's /// TIMEOUTSECONDS, marks them FAILED, and returns their JobExecutionIds so the /// caller can apply the same retry-or-DLQ decision used for ordinary failures. /// Task> ReapTimedOutExecutionsAsync( LoginDTO login, CancellationToken ct = default); /// /// Saves a snapshot of the exact outbound request (URL, method, correlation id, etc.) /// to TJOBEXECUTION.REQUESTPAYLOAD — called right before the HTTP call. /// Task SaveRequestPayloadAsync( int jobExecutionId, string payload, LoginDTO login, CancellationToken ct = default); /// /// Updates an existing TJOBEXECUTION row with SUCCESS(3) or FAILED(4) result. /// Task UpdateExecutionResultAsync( int jobExecutionId, bool isSuccess, string? message, LoginDTO login, CancellationToken ct = default); /// /// Returns all active jobs with their latest execution state, /// next run time, action types, retry count, and last error. /// Used by GET /Scheduler/Monitor. /// Task> LoadSchedulerMonitorAsync( LoginDTO login, CancellationToken ct = default); /// /// Loads a single job by JobId regardless of active/paused status — used by /// QuartzSyncService to (re)build or tear down that job's native Quartz trigger. /// Returns null if no MJOBDEFINE/TSCHEDULER row is found for this tenant. /// Task LoadJobForSyncAsync( int jobId, LoginDTO login, CancellationToken ct = default); /// /// Returns every JobId pointing at the given SchedulerId — used when a /// TSCHEDULER row is edited so every job using it gets resynced too. /// Task> GetJobIdsBySchedulerIdAsync( int schedulerId, LoginDTO login, CancellationToken ct = default); /// /// Toggles MJOBDEFINE.STATUS (1=Active, 4=Inactive per GB5's common Status enum). /// Caller is responsible for syncing/removing the Quartz trigger afterward. /// Task SetJobStatusAsync( int jobId, int status, LoginDTO login, CancellationToken ct = default); /// /// Bulk-deactivates all active MJOBDEFINE rows whose RUNASUSEID matches the given /// userId. Called when MUSER.STATUS goes inactive so the deactivated user's /// scheduled jobs stop running. Returns rows affected. /// Task DeactivateJobsByUserIdAsync( int userId, LoginDTO login, CancellationToken ct = default); /// /// Loads the flat criteria attribute rows for one TCRITERIACONFIG record. /// Used by the BLL to build the POST body (ReportCallingDTO / CriteriaDTO) /// for report endpoints that are invoked by the scheduler. /// Returns an empty list if criteriaConfigId <= 0 or no rows exist. /// Task> LoadCriteriaAttributesAsync( int criteriaConfigId, LoginDTO login, CancellationToken ct = default); } /// /// Flat row returned by GetCriteriaAttributesForExecution — one row per /// TCRITERIACONFIGATTRIBUTE entry. The BLL groups these by SectionId to /// build a CriteriaDTO with nested SectionCriteriaDTO / AttributesCriteriaDTO. /// public sealed class SchedulerCriteriaAttributeRow { public int SectionId { get; set; } public int SectionSlNo { get; set; } public string? SectionJoin { get; set; } public int AttributeId { get; set; } public string? FieldName { get; set; } public string? AttributeType { get; set; } public string? FilterType { get; set; } public string? OperationType { get; set; } public string? FieldValue { get; set; } public string? FieldValueIn { get; set; } public string? JoinType { get; set; } public string? FieldDisplayValue { get; set; } public string? VariableField { get; set; } public int PeriodFilter { get; set; } public int MenuId { get; set; } public int ReportFormatId { get; set; } public int ReportViewId { get; set; } } /// /// Scheduler Task DAL — Dapper-based implementation. /// All DB access goes through IQueryExecutor. /// EF contexts (Gb4DbContext, CentralDbContext) are separate concerns and not used here. /// public sealed class SchedulerTaskServiceDAL : ISchedulerTaskDAL { private readonly IQueryExecutor _queryExecutor; public SchedulerTaskServiceDAL(IQueryExecutor queryExecutor) { _queryExecutor = queryExecutor ?? throw new ArgumentNullException(nameof(queryExecutor)); } // ───────────────────────────────────────────────────────────────── // LoadReadyJobsAsync // ───────────────────────────────────────────────────────────────── public async Task> LoadReadyJobsAsync( LoginDTO login, CancellationToken ct = default) { return (await _queryExecutor.QueryAsync( login, SchedulerTaskGeneratorQB.LoadReadyJobs, new { TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false)).ToList(); } // ───────────────────────────────────────────────────────────────── // LoadJobDetailsAsync // ───────────────────────────────────────────────────────────────── public async Task LoadJobDetailsAsync( int jobExecutionId, LoginDTO login, CancellationToken ct = default) { var job = (await _queryExecutor.QueryAsync( login, SchedulerTaskGeneratorQB.LoadJobDetails, new { JobExecutionId = jobExecutionId, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false)).FirstOrDefault(); if (job == null) return null; // Load actions inline for single-record detail fetch job.Actions = await LoadActionsByJobIdsAsync( new[] { job.JobId }, login, ct).ConfigureAwait(false); return job; } // ───────────────────────────────────────────────────────────────── // GetUriParameterValueAsync — single-column fallback (see interface doc) // ───────────────────────────────────────────────────────────────── public async Task GetUriParameterValueAsync( int jobId, int tenantId, LoginDTO login, CancellationToken ct = default) { return await _queryExecutor.QuerySingleAsync( login, "SELECT URIPARAMETERVALUE FROM MJOBDEFINE WHERE JOBID = @JobId AND TENANTID = @TenantId", new { JobId = jobId, TenantId = tenantId }, cancellationToken: ct ).ConfigureAwait(false); } // ───────────────────────────────────────────────────────────────── // LoadActionsByJobIdsAsync // ───────────────────────────────────────────────────────────────── public async Task> LoadActionsByJobIdsAsync( IEnumerable jobIds, LoginDTO login, CancellationToken ct = default) { var ids = jobIds as int[] ?? jobIds.ToArray(); if (ids.Length == 0) return new List(); return (await _queryExecutor.QueryAsync( login, SchedulerTaskGeneratorQB.LoadActionsByJobIds, new { JobIds = ids, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false)).ToList(); } // ───────────────────────────────────────────────────────────────── // HasActiveSchedulerActionAsync // ───────────────────────────────────────────────────────────────── public async Task HasActiveSchedulerActionAsync( int jobId, LoginDTO login, CancellationToken ct = default) { var found = await _queryExecutor.QuerySingleAsync( login, SchedulerTaskGeneratorQB.HasActiveSchedulerAction, new { EventTypeId = SchedulerActionEventConstants.SchedulerJobEventTypeId, JobId = jobId, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false); return found.HasValue; } // ───────────────────────────────────────────────────────────────── // CreatePendingExecutionAsync // ───────────────────────────────────────────────────────────────── /// /// Inserts a new TJOBEXECUTION row with STATUS=Pending, due immediately. /// ✅ nowIst is converted to UTC and passed as @NextRunOn so /// TJOBEXECUTION.NEXTRUNON (the informational next-scheduled-run /// estimate) is populated in the DB at insert time. /// public async Task CreatePendingExecutionAsync( int jobId, DateTime nowIst, string triggerType, string hostInstance, string correlationId, string? traceParent, string? traceState, LoginDTO login, CancellationToken ct = default) { var now = DateTime.UtcNow; // ✅ QuerySingleAsync (not ExecuteAsync) — the SQL's OUTPUT INSERTED.JOBEXECUTIONID // returns exactly one row (the new id) when the insert happens, and zero rows when // the ISCONCURRENT guard blocks it. QuerySingleAsync maps "no rows" to default(long) // = 0, which callers treat as "nothing was inserted." return await _queryExecutor.QuerySingleAsync( login, SchedulerTaskGeneratorQB.CreatePendingExecution, new { JobId = jobId, RunOn = now, NowIst = nowIst, // ✅ Passed to SQL for NEXTRUNON computation TriggerType = triggerType, HostInstance = hostInstance, CorrelationId = correlationId, TraceParent = traceParent, TraceState = traceState, UserId = login.UserId, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false); } // ───────────────────────────────────────────────────────────────── // ClaimDueExecutionsAsync // ───────────────────────────────────────────────────────────────── public async Task> ClaimDueExecutionsAsync( int batchSize, LoginDTO login, CancellationToken ct = default) { return (await _queryExecutor.QueryAsync( login, SchedulerTaskGeneratorQB.ClaimDueExecutions, new { BatchSize = batchSize, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false)).ToList(); } // ───────────────────────────────────────────────────────────────── // GetRetryStateAsync // ───────────────────────────────────────────────────────────────── public async Task<(int RetryNumber, int MaxRetries, int JobId)> GetRetryStateAsync( int jobExecutionId, LoginDTO login, CancellationToken ct = default) { var row = await _queryExecutor.QuerySingleAsync( login, SchedulerTaskGeneratorQB.GetRetryState, new { JobExecutionId = jobExecutionId, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false); if (row is null) throw new InvalidOperationException( $"GetRetryStateAsync: no TJOBEXECUTION/MJOBDEFINE row found for JobExecutionId={jobExecutionId}"); return (row.RetryNumber, row.MaxRetries, row.JobId); } private sealed class RetryStateRow { public int RetryNumber { get; set; } public int MaxRetries { get; set; } public int JobId { get; set; } } // ───────────────────────────────────────────────────────────────── // ApplyRetryDecisionAsync // ───────────────────────────────────────────────────────────────── public async Task ApplyRetryDecisionAsync( int jobExecutionId, int newRetryNumber, string newStatus, DateTime? nextAttemptOn, LoginDTO login, CancellationToken ct = default) { await _queryExecutor.ExecuteAsync( login, SchedulerTaskGeneratorQB.ApplyRetryDecision, new { JobExecutionId = jobExecutionId, NewRetryNumber = newRetryNumber, NewStatus = newStatus, NextAttemptOn = nextAttemptOn, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false); } // ───────────────────────────────────────────────────────────────── // ReapTimedOutExecutionsAsync // ───────────────────────────────────────────────────────────────── public async Task> ReapTimedOutExecutionsAsync( LoginDTO login, CancellationToken ct = default) { return (await _queryExecutor.QueryAsync( login, SchedulerTaskGeneratorQB.ReapTimedOutExecutions, new { TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false)).ToList(); } // ───────────────────────────────────────────────────────────────── // SaveRequestPayloadAsync // ───────────────────────────────────────────────────────────────── public async Task SaveRequestPayloadAsync( int jobExecutionId, string payload, LoginDTO login, CancellationToken ct = default) { await _queryExecutor.ExecuteAsync( login, SchedulerTaskGeneratorQB.SaveRequestPayload, new { JobExecutionId = jobExecutionId, Payload = payload, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false); } // ───────────────────────────────────────────────────────────────── // UpdateExecutionResultAsync // ───────────────────────────────────────────────────────────────── /// /// Writes COMPLETED or FAILED result to TJOBEXECUTION. /// public async Task UpdateExecutionResultAsync( int jobExecutionId, bool isSuccess, string? message, LoginDTO login, CancellationToken ct = default) { var now = DateTime.UtcNow; await _queryExecutor.ExecuteAsync( login, SchedulerTaskGeneratorQB.UpdateExecutionResult, new { JobExecutionId = jobExecutionId, Status = isSuccess ? FrameworkDAL.DTO.SchedulerTaskGenerator.SchedulerExecutionStatus.Success : FrameworkDAL.DTO.SchedulerTaskGenerator.SchedulerExecutionStatus.Failed, EndRunOn = now, Message = message, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false); } // ───────────────────────────────────────────────────────────────── // Add this method to SchedulerTaskServiceDAL implementation // ───────────────────────────────────────────────────────────────── /// /// Executes LoadSchedulerMonitor query and returns flat monitor rows. /// NextRunOn for jobs with no execution row yet is computed in BLL. /// public async Task> LoadSchedulerMonitorAsync( LoginDTO login, CancellationToken ct = default) { return (await _queryExecutor.QueryAsync( login, SchedulerTaskGeneratorQB.LoadSchedulerMonitor, new { TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false)).ToList(); } // ───────────────────────────────────────────────────────────────── // LoadJobForSyncAsync // ───────────────────────────────────────────────────────────────── public async Task LoadJobForSyncAsync( int jobId, LoginDTO login, CancellationToken ct = default) { return (await _queryExecutor.QueryAsync( login, SchedulerTaskGeneratorQB.LoadJobForSync, new { JobId = jobId, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false)).FirstOrDefault(); } // ───────────────────────────────────────────────────────────────── // GetJobIdsBySchedulerIdAsync // ───────────────────────────────────────────────────────────────── public async Task> GetJobIdsBySchedulerIdAsync( int schedulerId, LoginDTO login, CancellationToken ct = default) { return (await _queryExecutor.QueryAsync( login, SchedulerTaskGeneratorQB.GetJobIdsBySchedulerId, new { SchedulerId = schedulerId, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false)).ToList(); } // ───────────────────────────────────────────────────────────────── // SetJobStatusAsync // ───────────────────────────────────────────────────────────────── public async Task SetJobStatusAsync( int jobId, int status, LoginDTO login, CancellationToken ct = default) { await _queryExecutor.ExecuteAsync( login, SchedulerTaskGeneratorQB.SetJobStatus, new { JobId = jobId, Status = status, UserId = login.UserId, TenantId = login.ClientId }, cancellationToken: ct ).ConfigureAwait(false); } // ───────────────────────────────────────────────────────────────── // DeactivateJobsByUserIdAsync // ───────────────────────────────────────────────────────────────── public async Task DeactivateJobsByUserIdAsync( int userId, LoginDTO login, CancellationToken ct = default) { return await _queryExecutor.ExecuteAsync( login, SchedulerTaskGeneratorQB.DeactivateJobsByUserId, new { UserId = userId, TenantId = login.ClientId, RequestingUserId = login.UserId }, cancellationToken: ct ).ConfigureAwait(false); } // ───────────────────────────────────────────────────────────────── // LoadCriteriaAttributesAsync // ───────────────────────────────────────────────────────────────── public async Task> LoadCriteriaAttributesAsync( int criteriaConfigId, LoginDTO login, CancellationToken ct = default) { if (criteriaConfigId <= 0) return new List(); return (await _queryExecutor.QueryAsync( login, SchedulerTaskGeneratorQB.GetCriteriaAttributesForExecution, new { CriteriaConfigId = criteriaConfigId }, cancellationToken: ct ).ConfigureAwait(false)).ToList(); } } }