using Cronos; using FrameworkDAL.CustomCode.SchedulerTaskGenerator; using FrameworkDAL.DTO.SchedulerTaskGenerator; using GB5Shared.DTO.Framework.Login; using GB5Shared.Telemetry; using Microsoft.Extensions.Logging; using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; using System.Text.Json; using System.Threading; using System.Threading.Tasks; namespace FrameworkBLL.SchedulerTaskGenerator { // ═══════════════════════════════════════════════════════════════════ // ISchedulerTaskServiceBLL // Owns the scheduler execution lifecycle: // Load → Mark Started → Publish → Mark Published → Update Result // ═══════════════════════════════════════════════════════════════════ /// /// Interface for Scheduler Task Service /// public interface ISchedulerTaskServiceBLL { /// /// Returns all scheduler jobs which are due now (NextRunOn computed + <= nowIst, actions loaded). /// Task> LoadExecutableTasksAsync( LoginDTO login, DateTime nowIst, CancellationToken ct = default); /// /// Returns a single scheduler job with full details + computed NextRunOn. /// Task GetTaskDetailsAsync( int jobExecutionId, LoginDTO login, CancellationToken ct = default); /// /// Computes NextRunOn and inserts a new TJOBEXECUTION row with STATUS=InProgress. /// Called by the scheduler before publishing to queue. Returns the newly generated /// JOBEXECUTIONID (0 if no row was inserted — e.g. blocked by ISCONCURRENT) so the /// caller can publish/tag with the correct id instead of a stale pre-insert value. /// Task MarkExecutionStartedAsync( int jobId, LoginDTO login, CancellationToken ct = default, string triggerType = "SCHEDULED"); /// /// Reaps InProgress executions that exceeded their MJOBDEFINE.TIMEOUTSECONDS, /// marks them Failed, and applies the same retry-or-DLQ decision as an /// ordinary execution failure. Called once per poll cycle before loading ready jobs. /// Task ReapTimedOutExecutionsAsync( LoginDTO login, CancellationToken ct = default); /// /// Atomically claims up to due TJOBEXECUTION rows /// (STATUS='PENDING', NEXTATTEMPTON due) for delivery, flipping them to 'PROCESSING' /// (InProgress). Called by SchedulerExecutionPoller once per poll tick per tenant. /// Task> ClaimDueExecutionsAsync( int batchSize, LoginDTO login, CancellationToken ct = default); /// /// True if this job has at least one active MACTION row linked via the scheduler's /// fixed EVENTTYPEID — used to decide whether a SCHEDULER_ACTION_EVENT is worth /// enqueueing to TJOBQUEUE after the job completes successfully. /// Task HasActiveSchedulerActionAsync( int jobId, LoginDTO login, CancellationToken ct = default); /// /// Saves a snapshot of the exact outbound request 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/failure result. /// If successful — schedules next run by inserting a new TJOBEXECUTION row. /// Task UpdateExecutionResultAsync( int jobExecutionId, bool isSuccess, string? message, LoginDTO login, CancellationToken ct = default); /// /// Bulk-deactivates all MJOBDEFINE rows whose RUNASUSEID matches userId, /// then removes their Quartz triggers. Called when a user is deactivated in MUSER /// so the deactivated user's scheduled jobs stop firing automatically. /// Task DeactivateJobsByUserAsync( int userId, LoginDTO login, CancellationToken ct = default); /// /// Loads the TCRITERIACONFIG for criteriaConfigId, resolves variable field tokens /// against the given login (WorkOUId, WorkPartyBranchId, date macros, etc.), and /// returns a JSON-serialized ReportCallingDTO ready to send as the POST body for a /// report endpoint. Returns null if criteriaConfigId <= 0 or no attributes exist. /// Task LoadCriteriaForExecutionAsync( int criteriaConfigId, LoginDTO login, CancellationToken ct = default); } // ═══════════════════════════════════════════════════════════════════ // ISchedulerMonitorBLL // Separate interface for monitor concerns — keeps monitor logic // out of ISchedulerTaskServiceBLL which owns the execution flow. // SL layer only calls this interface — never touches BLL internals. // ═══════════════════════════════════════════════════════════════════ /// /// Interface for Scheduler Monitor Service. /// Returns a full monitor view of all jobs — status, next run, /// action types, retry count, duration, and last error. /// public interface ISchedulerMonitorBLL { /// /// Returns all active jobs with latest execution state. /// Jobs with no execution row get NextRunOn computed in-memory. /// Result includes a summary count block and the full job list. /// Task GetMonitorAsync( LoginDTO login, CancellationToken ct = default); } // ═══════════════════════════════════════════════════════════════════ // SchedulerCronHelper (internal static — shared by both BLL classes) // Quartz → Cronos normalizer. Keeping it here avoids duplicating // the NormalizeQuartz logic across SchedulerTaskServiceBLL and // SchedulerMonitorBLL. // ═══════════════════════════════════════════════════════════════════ internal static class SchedulerCronHelper { /// /// Converts a Quartz 7-part cron expression to Cronos-compatible 6-part format. /// Quartz: seconds minutes hours day-of-month month day-of-week [year] /// Cronos: seconds minutes hours day-of-month month day-of-week (no year) /// internal static string NormalizeQuartz(string quartz) { var parts = quartz .Split(' ', StringSplitOptions.RemoveEmptyEntries) .ToList(); // Remove Quartz year field (7th part) if (parts.Count == 7) parts.RemoveAt(6); return string.Join(" ", parts.Select(p => p .Replace("?", "*") // Quartz ? wildcard → Cronos * .Replace("1/1", "*") // ✅ "every 1 from day 1" → "every day" .Replace("0/1", "*") // ✅ "every 1 from day 0" → "every day" )); } } // ═══════════════════════════════════════════════════════════════════ // SchedulerTaskServiceBLL // ═══════════════════════════════════════════════════════════════════ /// /// Scheduler Task Service Implementation. /// Owns the full job execution lifecycle. /// public sealed class SchedulerTaskServiceBLL : ISchedulerTaskServiceBLL { private readonly ISchedulerTaskDAL _dal; private readonly ILogger _logger; private static readonly TimeZoneInfo IST = TimeZoneInfo.FindSystemTimeZoneById("India Standard Time"); // Track JobIds that already logged the multi-schedule warning — avoids log flood private static readonly ConcurrentDictionary _warnedJobIds = new(); public SchedulerTaskServiceBLL( ISchedulerTaskDAL dal, ILogger logger) { _dal = dal ?? throw new ArgumentNullException(nameof(dal)); _logger = logger ?? throw new ArgumentNullException(nameof(logger)); } // ───────────────────────────────────────────────────────────────── // LoadExecutableTasksAsync // ───────────────────────────────────────────────────────────────── /// /// Loads ready jobs, computes NextRunOn, filters to only those DUE (NextRunOn <= nowIst), /// batch-loads actions, and returns ordered by NextRunOn. /// public async Task> LoadExecutableTasksAsync( LoginDTO login, DateTime nowIst, CancellationToken ct = default) { GB5Trace.Step("load-executable-tasks", new { nowIst }); try { var jobs = await _dal.LoadReadyJobsAsync(login, ct).ConfigureAwait(false); if (jobs.Count == 0) return jobs; // ✅ Due-ness is derived from LastRunOn (or "never run" if null) via IsDue, // NOT from ComputeNextRun's result. ComputeNextRun always advances to a // strictly-future slot (see its while-loop) so it can never satisfy a // "<= nowIst" check — using it for due-filtering meant this method could // never return any job, regardless of schedule type or LastRunOn. var dueJobs = jobs.Where(j => IsDue(j, nowIst)).ToList(); if (dueJobs.Count == 0) return dueJobs; // ComputeNextRun here is informational only (logging / RabbitMQ payload) — // it does not affect which jobs were selected as due. foreach (var job in dueJobs) job.NextRunOn = ComputeNextRun(job, nowIst); dueJobs = dueJobs.OrderBy(j => j.NextRunOn).ToList(); // Batch-load actions for all due jobs in one query — no N+1 var actions = await _dal.LoadActionsByJobIdsAsync( dueJobs.Select(j => j.JobId), login, ct).ConfigureAwait(false); foreach (var job in dueJobs) job.Actions = actions.Where(a => a.JobId == job.JobId).ToList(); GB5Trace.Tag("gb5.scheduler.due_count", dueJobs.Count); return dueJobs; } catch (Exception ex) { GB5Trace.MarkFailed("load-executable-tasks-failed", ex); _logger.LogError(ex, "LoadExecutableTasksAsync failed"); throw; } } // ───────────────────────────────────────────────────────────────── // GetTaskDetailsAsync // ───────────────────────────────────────────────────────────────── /// /// Returns a single job with full details. /// Computes NextRunOn in-memory since NEXTRUNON is not always stored in DB. /// public async Task GetTaskDetailsAsync( int jobExecutionId, LoginDTO login, CancellationToken ct = default) { GB5Trace.Step("get-task-details", new { jobExecutionId }); try { var job = await _dal.LoadJobDetailsAsync(jobExecutionId, login, ct) .ConfigureAwait(false); if (job is null) return null; var nowIst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, IST); job.NextRunOn = ComputeNextRun(job, nowIst); return job; } catch (Exception ex) { GB5Trace.MarkFailed("get-task-details-failed", ex); _logger.LogError(ex, "GetTaskDetailsAsync failed | JobExecutionId={Id}", jobExecutionId); throw; } } // ───────────────────────────────────────────────────────────────── // MarkExecutionStartedAsync // ───────────────────────────────────────────────────────────────── /// /// Computes NextRunOn for the job, then inserts a TJOBEXECUTION row with /// STATUS=InProgress and NEXTRUNON saved — so the DB always has the next run time. /// public async Task MarkExecutionStartedAsync( int jobId, LoginDTO login, CancellationToken ct = default, string triggerType = "SCHEDULED") { GB5Trace.Step("mark-execution-started", new { jobId, triggerType }); try { var nowIst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, IST); // Pass nowIst to DAL so TJOBEXECUTION.NEXTRUNON is saved in DB. // HostInstance identifies which pod/machine dispatched this execution — // needed for cluster attribution once multiple SchedulerQuartzJob instances run. // Same convention as SchedulerJobExecutorQuartzJob: generate the CorrelationId // here, app-side, so TJOBEXECUTION.CORRELATIONID is set at insert time rather // than via a separate DB-side NEWID(). var correlationId = Guid.NewGuid().ToString().ToUpperInvariant(); var newJobExecutionId = await _dal.CreatePendingExecutionAsync( jobId, nowIst, triggerType, Environment.MachineName, correlationId, traceParent: null, traceState: null, login, ct) .ConfigureAwait(false); _logger.LogInformation( "Scheduler execution started | JobId={JobId} JobExecutionId={JobExecutionId} " + "NowIst={NowIst:HH:mm:ss} TriggerType={TriggerType}", jobId, newJobExecutionId, nowIst, triggerType); return newJobExecutionId; } catch (Exception ex) { GB5Trace.MarkFailed("mark-execution-started-failed", ex); _logger.LogError(ex, "MarkExecutionStartedAsync failed | JobId={JobId}", jobId); throw; } } // ───────────────────────────────────────────────────────────────── // ReapTimedOutExecutionsAsync // ───────────────────────────────────────────────────────────────── /// /// Finds InProgress executions whose STARTRUNON exceeds MJOBDEFINE.TIMEOUTSECONDS, /// marks them Failed, and routes each through the same retry-or-DLQ decision as an /// ordinary failure so MAXRETRIES is respected consistently for timeouts too. /// public async Task ReapTimedOutExecutionsAsync( LoginDTO login, CancellationToken ct = default) { GB5Trace.Step("reap-timedout-executions"); try { var timedOutIds = await _dal.ReapTimedOutExecutionsAsync(login, ct) .ConfigureAwait(false); foreach (var jobExecutionId in timedOutIds) { try { await ApplyRetryOrDeadLetterAsync(jobExecutionId, login, ct) .ConfigureAwait(false); } catch (Exception ex) { // Don't let one bad retry decision stop the rest of the reap batch. GB5Trace.RecordError(ex, context: "reap-retry-decision-failed"); _logger.LogError(ex, "ReapTimedOutExecutionsAsync: retry/DLQ decision failed | JobExecutionId={Id}", jobExecutionId); } } if (timedOutIds.Count > 0) _logger.LogWarning( "Scheduler reaped {Count} timed-out execution(s) | Ids=[{Ids}]", timedOutIds.Count, string.Join(",", timedOutIds)); return timedOutIds.Count; } catch (Exception ex) { GB5Trace.MarkFailed("reap-timedout-executions-failed", ex); _logger.LogError(ex, "ReapTimedOutExecutionsAsync failed"); throw; } } // ───────────────────────────────────────────────────────────────── // ClaimDueExecutionsAsync // ───────────────────────────────────────────────────────────────── public async Task> ClaimDueExecutionsAsync( int batchSize, LoginDTO login, CancellationToken ct = default) { GB5Trace.Step("claim-due-executions", new { batchSize }); try { var claimed = await _dal.ClaimDueExecutionsAsync(batchSize, login, ct) .ConfigureAwait(false); if (claimed.Count > 0) _logger.LogInformation("Scheduler claimed {Count} due execution(s)", claimed.Count); return claimed; } catch (Exception ex) { GB5Trace.MarkFailed("claim-due-executions-failed", ex); _logger.LogError(ex, "ClaimDueExecutionsAsync failed"); throw; } } // ───────────────────────────────────────────────────────────────── // HasActiveSchedulerActionAsync // ───────────────────────────────────────────────────────────────── public async Task HasActiveSchedulerActionAsync( int jobId, LoginDTO login, CancellationToken ct = default) { GB5Trace.Step("has-active-scheduler-action", new { jobId }); try { return await _dal.HasActiveSchedulerActionAsync(jobId, login, ct) .ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("has-active-scheduler-action-failed", ex); _logger.LogError(ex, "HasActiveSchedulerActionAsync failed | JobId={JobId}", jobId); throw; } } // ───────────────────────────────────────────────────────────────── // ApplyRetryOrDeadLetterAsync (internal) // Shared by UpdateExecutionResultAsync (ordinary failure) and // ReapTimedOutExecutionsAsync (timeout failure). Backoff schedule // matches the documented TJOBQUEUE schedule: 30s → 2m → 8m → 30m → 2h → DLQ. // ───────────────────────────────────────────────────────────────── private static readonly int[] RetryBackoffSeconds = { 30, 120, 480, 1800, 7200 }; internal async Task ApplyRetryOrDeadLetterAsync( int jobExecutionId, LoginDTO login, CancellationToken ct = default) { var (retryNumber, maxRetries, jobId) = await _dal .GetRetryStateAsync(jobExecutionId, login, ct) .ConfigureAwait(false); var newRetryNumber = retryNumber + 1; if (newRetryNumber >= maxRetries) { // Exhausted — DLQ, no further NEXTATTEMPTON change. await _dal.ApplyRetryDecisionAsync( jobExecutionId, newRetryNumber, newStatus: SchedulerExecutionStatus.DLQ, nextAttemptOn: null, login, ct) .ConfigureAwait(false); _logger.LogError( "Scheduler execution moved to DLQ | JobExecutionId={Id} JobId={JobId} RetryNumber={Retry}/{Max}", jobExecutionId, jobId, newRetryNumber, maxRetries); GB5Trace.Tag("gb5.scheduler.dlq", true); return; } var backoffIndex = Math.Min(newRetryNumber - 1, RetryBackoffSeconds.Length - 1); var nextAttemptOn = DateTime.UtcNow.AddSeconds(RetryBackoffSeconds[backoffIndex]); await _dal.ApplyRetryDecisionAsync( jobExecutionId, newRetryNumber, newStatus: SchedulerExecutionStatus.Pending, nextAttemptOn, login, ct) .ConfigureAwait(false); _logger.LogWarning( "Scheduler execution requeued for retry | JobExecutionId={Id} JobId={JobId} " + "RetryNumber={Retry}/{Max} NextAttemptOn={NextAttempt:HH:mm:ss}", jobExecutionId, jobId, newRetryNumber, maxRetries, nextAttemptOn); } // ───────────────────────────────────────────────────────────────── // SaveRequestPayloadAsync // ───────────────────────────────────────────────────────────────── public async Task SaveRequestPayloadAsync( int jobExecutionId, string payload, LoginDTO login, CancellationToken ct = default) { GB5Trace.Step("save-request-payload", new { jobExecutionId }); try { await _dal.SaveRequestPayloadAsync(jobExecutionId, payload, login, ct) .ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("save-request-payload-failed", ex); _logger.LogError(ex, "SaveRequestPayloadAsync failed | JobExecutionId={Id}", jobExecutionId); throw; } } // ───────────────────────────────────────────────────────────────── // UpdateExecutionResultAsync // ───────────────────────────────────────────────────────────────── /// /// Writes success/failure result to TJOBEXECUTION. /// On success — computes next run and inserts a fresh TJOBEXECUTION (STATUS=Pending) /// so the scheduler picks it up on the next poll cycle. /// public async Task UpdateExecutionResultAsync( int jobExecutionId, bool isSuccess, string? message, LoginDTO login, CancellationToken ct = default) { GB5Trace.Step("update-execution-result", new { jobExecutionId, isSuccess }); try { await _dal.UpdateExecutionResultAsync(jobExecutionId, isSuccess, message, login, ct) .ConfigureAwait(false); _logger.LogInformation( "Scheduler execution result saved | JobExecutionId={JobExecutionId} Success={IsSuccess}", jobExecutionId, isSuccess); if (!isSuccess) { // ✅ Failure: route through the shared retry-or-DLQ decision instead of // leaving RETRYNUMBER at 0 forever — respects MJOBDEFINE.MAXRETRIES. await ApplyRetryOrDeadLetterAsync(jobExecutionId, login, ct) .ConfigureAwait(false); return; } // ✅ No new row is inserted here. This completed row's own LASTRUNON is now // the historical anchor — the next poll's LoadReadyJobsAsync/IsDue derives // due-ness fresh from it (see LoadReadyJobs SQL, which keeps the most recent // row regardless of status). Previously this called MarkJobInProgressAsync to // eagerly seed a "next" row, but that inserted it as STATUS=2 (InProgress) — // a status LoadReadyJobs's OUTER APPLY excluded from future consideration — // so a job would silently stop recurring after its first success. var details = await _dal.LoadJobDetailsAsync(jobExecutionId, login, ct) .ConfigureAwait(false); if (details is null) { _logger.LogWarning( "Job details not found post-success — cannot log next run | JobExecutionId={JobExecutionId}", jobExecutionId); return; } var nowIst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, IST); var nextRun = ComputeNextRun(details, nowIst); _logger.LogInformation( "Scheduler execution succeeded | JobId={JobId} EstimatedNextRun={NextRunOn:yyyy-MM-dd HH:mm:ss}", details.JobId, nextRun); } catch (Exception ex) { GB5Trace.MarkFailed("update-execution-result-failed", ex); _logger.LogError(ex, "UpdateExecutionResultAsync failed | JobExecutionId={Id}", jobExecutionId); throw; } } // ───────────────────────────────────────────────────────────────── // DeactivateJobsByUserAsync // ───────────────────────────────────────────────────────────────── public async Task DeactivateJobsByUserAsync( int userId, LoginDTO login, CancellationToken ct = default) { GB5Trace.Step("deactivate-jobs-by-user", new { userId }); try { var count = await _dal.DeactivateJobsByUserIdAsync(userId, login, ct) .ConfigureAwait(false); if (count > 0) _logger.LogInformation( "DeactivateJobsByUserAsync: deactivated {Count} job(s) for UserId={UserId}", count, userId); return count; } catch (Exception ex) { GB5Trace.MarkFailed("deactivate-jobs-by-user-failed", ex); _logger.LogError(ex, "DeactivateJobsByUserAsync failed | UserId={UserId}", userId); throw; } } // ───────────────────────────────────────────────────────────────── // LoadCriteriaForExecutionAsync // ───────────────────────────────────────────────────────────────── public async Task LoadCriteriaForExecutionAsync( int criteriaConfigId, LoginDTO login, CancellationToken ct = default) { if (criteriaConfigId <= 0) return null; GB5Trace.Step("load-criteria-for-execution", new { criteriaConfigId }); try { var rows = await _dal.LoadCriteriaAttributesAsync(criteriaConfigId, login, ct) .ConfigureAwait(false); if (rows.Count == 0) return null; var nowIst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, IST); var today = nowIst.Date; // Variable-field token map — matches MenuDAL.cs lines 229-264 exactly. // Each {Token} is replaced with the login's org context or a computed date. var tokens = new System.Collections.Generic.Dictionary(StringComparer.OrdinalIgnoreCase) { ["{WorkPartyBranchID}"] = login.WorkPartyBranchId.ToString(), ["{WorkOUID}"] = login.WorkOUId.ToString(), ["{WorkStoreId}"] = login.WorkStoreId.ToString(), ["{WorkPeriodId}"] = login.WorkPeriodId.ToString(), ["{TodayFrom}"] = today.ToString("yyyy-MM-dd 00:00:00.000"), ["{TodayTo}"] = today.ToString("yyyy-MM-dd 23:59:59.000"), ["{TodayTimeStamp}"] = nowIst.ToString("yyyy-MM-dd HH:mm:ss.000"), ["{YesterdayFrom}"] = today.AddDays(-1).ToString("yyyy-MM-dd 00:00:00.000"), ["{YesterdayTo}"] = today.AddDays(-1).ToString("yyyy-MM-dd 23:59:59.000"), ["{YesterdayTimeStamp}"]= today.AddDays(-1).ToString("yyyy-MM-dd HH:mm:ss.000"), ["{WeekFrom}"] = today.AddDays(-7).ToString("yyyy-MM-dd 00:00:00.000"), ["{MonthFrom}"] = today.AddMonths(-1).ToString("yyyy-MM-dd 00:00:00.000"), ["{YearFrom}"] = today.AddYears(-1).ToString("yyyy-MM-dd 00:00:00.000"), ["{CurrentMonthFrom}"] = new DateTime(today.Year, today.Month, 1).ToString("yyyy-MM-dd 00:00:00.000"), ["{CurrentYearFrom}"] = new DateTime(today.Year, 1, 1).ToString("yyyy-MM-dd 00:00:00.000"), ["{CurrentYearTo}"] = new DateTime(today.Year, 12, 31).ToString("yyyy-MM-dd 23:59:59.000"), }; string ResolveVariable(string? variableField, string? fieldValue) { if (string.IsNullOrWhiteSpace(variableField)) return fieldValue ?? string.Empty; var resolved = variableField; foreach (var kv in tokens) resolved = resolved.Replace(kv.Key, kv.Value, StringComparison.OrdinalIgnoreCase); return resolved; } // Group by section, preserving section order var sections = rows .GroupBy(r => r.SectionId) .Select(g => { var first = g.First(); var attributes = g.Select(r => new { FieldName = r.FieldName ?? string.Empty, OperationType = r.OperationType ?? "Equal", FieldValue = (object)ResolveVariable(r.VariableField, r.FieldValue), InArray = string.IsNullOrWhiteSpace(r.FieldValueIn) ? System.Array.Empty() : (object[])r.FieldValueIn.Split(',').Cast().ToArray(), JoinType = r.JoinType ?? "And", CriteriaAttributeId = r.AttributeId, FilterType = r.FilterType ?? string.Empty, CriteriaFieldDisplayValue = r.FieldDisplayValue ?? string.Empty, TableAlias = string.Empty }).ToList(); return new { SectionId = first.SectionId, OperationType = first.SectionJoin ?? "And", AttributesCriteriaList = attributes }; }).ToList(); var first0 = rows.First(); var criteriaDto = new { SectionCriteriaList = sections }; var reportCallingDto = new { CriteriaDTO = criteriaDto, MenuId = first0.MenuId, ReportFormatId = first0.ReportFormatId, ReportViewId = first0.ReportViewId, PeriodFilter = first0.PeriodFilter }; return System.Text.Json.JsonSerializer.Serialize(reportCallingDto); } catch (Exception ex) { GB5Trace.MarkFailed("load-criteria-for-execution-failed", ex); _logger.LogError(ex, "LoadCriteriaForExecutionAsync failed | CriteriaConfigId={CriteriaConfigId}", criteriaConfigId); throw; } } // ───────────────────────────────────────────────────────────────── // IsDue (internal) // The real due-check. Deliberately separate from ComputeNextRun: // ComputeNextRun always advances to a future slot (for display / // messaging), so it can never be compared "<= nowIst" to detect // due-ness. This computes the occurrence that should follow // LastRunOn WITHOUT clamping it into the future, and checks whether // that occurrence has already arrived. // ───────────────────────────────────────────────────────────────── internal bool IsDue(SchedulerTaskDTO job, DateTime nowIst) { try { // TJOBEXECUTION.LASTRUNON is stored in UTC (DAL uses DateTime.UtcNow) — // convert to IST once, up front, so every branch below compares like-for-like // against nowIst (IST). Mixing UTC LastRunOn with IST nowIst directly was the // root of a ~5.5 hour skew bug. DateTime? lastRunIst = job.LastRunOn.HasValue ? TimeZoneInfo.ConvertTimeFromUtc(job.LastRunOn.Value, IST) : null; if (string.Equals(job.SchedulerType, "ONE_TIME", StringComparison.OrdinalIgnoreCase)) return lastRunIst is null && job.ScheduledDate.HasValue && job.ScheduledDate.Value <= nowIst; if (string.Equals(job.SchedulerType, "INTERVAL", StringComparison.OrdinalIgnoreCase) && job.IntervalSeconds is > 0) { if (lastRunIst is null) return true; // never run — due now return lastRunIst.Value.AddSeconds(job.IntervalSeconds.Value) <= nowIst; } if (!string.IsNullOrWhiteSpace(job.CronExpression)) { if (job.LastRunOn is null) return true; // never run — due now var normalized = SchedulerCronHelper.NormalizeQuartz(job.CronExpression); var cron = CronExpression.Parse(normalized, CronFormat.IncludeSeconds); // Cronos works natively in UTC — pass LastRunOn (already UTC) directly. var nextUtc = cron.GetNextOccurrence(job.LastRunOn.Value); if (!nextUtc.HasValue) return false; var nextIst = TimeZoneInfo.ConvertTimeFromUtc(nextUtc.Value, IST); return nextIst <= nowIst; } if (TimeSpan.TryParse(job.OccursAtTime, out var time)) { if (lastRunIst is null) return true; // never run — due now var todaysSlot = lastRunIst.Value.Date.Add(time); var nextSlot = todaysSlot > lastRunIst.Value ? todaysSlot : todaysSlot.AddDays(1); return nextSlot <= nowIst; } if (job.EveryNoOfMinutes > 0 || job.EveryNoOfHours > 0) { if (lastRunIst is null) return true; // never run — due now var next = lastRunIst.Value.AddHours(job.EveryNoOfHours).AddMinutes(job.EveryNoOfMinutes); return next <= nowIst; } } catch (Exception ex) { _logger.LogWarning(ex, "IsDue compute failed | JobId={JobId}", job.JobId); } return false; } // ───────────────────────────────────────────────────────────────── // ComputeNextRun (internal) // Informational only — computes the upcoming (always-future) slot for // logging / RabbitMQ payload / monitor display. Never used for due- // checking (see IsDue). Used only within the BLL assembly. // ───────────────────────────────────────────────────────────────── internal DateTime? ComputeNextRun(SchedulerTaskDTO job, DateTime nowIst) { try { // TJOBEXECUTION.LASTRUNON is stored in UTC — convert to IST before using // it alongside nowIst (IST). See IsDue for the same fix. DateTime? lastRunIst = job.LastRunOn.HasValue ? TimeZoneInfo.ConvertTimeFromUtc(job.LastRunOn.Value, IST) : null; // ── SchedulerType-driven schedule (added with 2026-08-13 schema) ── // Only engages when a job's TSCHEDULER row explicitly sets SCHEDULERTYPE — // existing jobs have this column NULL, so this never changes their behavior. if (string.Equals(job.SchedulerType, "ONE_TIME", StringComparison.OrdinalIgnoreCase)) { // A one-time job never reschedules once its ScheduledDate has passed. return job.ScheduledDate.HasValue && job.ScheduledDate.Value > nowIst ? job.ScheduledDate.Value : null; } if (string.Equals(job.SchedulerType, "INTERVAL", StringComparison.OrdinalIgnoreCase) && job.IntervalSeconds is > 0) { var baseTime = lastRunIst ?? nowIst; var next = baseTime.AddSeconds(job.IntervalSeconds.Value); while (next <= nowIst) next = next.AddSeconds(job.IntervalSeconds.Value); return next; } // ── CRON based schedule ─────────────────────────────────── if (!string.IsNullOrWhiteSpace(job.CronExpression)) { // Log multi-schedule warning only ONCE per JobId — prevents log flood if ((!string.IsNullOrWhiteSpace(job.OccursAtTime) || job.EveryNoOfMinutes > 0 || job.EveryNoOfHours > 0) && _warnedJobIds.TryAdd(job.JobId, true)) { _logger.LogWarning( "Job has multiple schedule fields set — CronExpression takes priority | JobId={JobId}", job.JobId); } var normalized = SchedulerCronHelper.NormalizeQuartz(job.CronExpression); var cron = CronExpression.Parse(normalized, CronFormat.IncludeSeconds); var baseUtc = TimeZoneInfo.ConvertTimeToUtc(nowIst, IST); var nextUtc = cron.GetNextOccurrence(baseUtc); return nextUtc.HasValue ? TimeZoneInfo.ConvertTimeFromUtc(nextUtc.Value, IST) : null; } // ── Fixed daily execution time ──────────────────────────── if (TimeSpan.TryParse(job.OccursAtTime, out var time)) { if ((job.EveryNoOfMinutes > 0 || job.EveryNoOfHours > 0) && _warnedJobIds.TryAdd(job.JobId, true)) { _logger.LogWarning( "Job has both OccursAtTime and interval fields set — OccursAtTime takes priority | JobId={JobId}", job.JobId); } var next = nowIst.Date.Add(time); return next <= nowIst ? next.AddDays(1) : next; } // ── Interval based execution ────────────────────────────── if (job.EveryNoOfMinutes > 0 || job.EveryNoOfHours > 0) { var baseTime = lastRunIst ?? nowIst.AddHours(-job.EveryNoOfHours).AddMinutes(-job.EveryNoOfMinutes); var next = baseTime .AddHours(job.EveryNoOfHours) .AddMinutes(job.EveryNoOfMinutes); // Advance until future — prevents overdue jobs firing on every poll while (next <= nowIst) next = next .AddHours(job.EveryNoOfHours) .AddMinutes(job.EveryNoOfMinutes); return next; } } catch (Exception ex) { // Do not crash — log and skip so other jobs still run _logger.LogWarning(ex, "Scheduler NextRun compute failed | JobId={JobId} CronExpression={Cron}", job.JobId, job.CronExpression); } return null; } } // ═══════════════════════════════════════════════════════════════════ // SchedulerMonitorBLL // Separate class for monitor concerns — reads DB state and builds // the monitor response. SL endpoint calls ISchedulerMonitorBLL only. // Uses SchedulerCronHelper.NormalizeQuartz shared with above class. // ═══════════════════════════════════════════════════════════════════ /// /// Scheduler Monitor Service Implementation. /// Loads all jobs from DAL and computes NextRunOn in-memory for /// jobs that have no execution row yet in TJOBEXECUTION. /// public sealed class SchedulerMonitorBLL : ISchedulerMonitorBLL { private readonly ISchedulerTaskDAL _dal; private readonly ILogger _logger; private static readonly TimeZoneInfo IST = TimeZoneInfo.FindSystemTimeZoneById("India Standard Time"); // Shared warned-job tracker — prevents monitor log flood private static readonly ConcurrentDictionary _warnedJobIds = new(); public SchedulerMonitorBLL( ISchedulerTaskDAL dal, ILogger logger) { _dal = dal ?? throw new ArgumentNullException(nameof(dal)); _logger = logger ?? throw new ArgumentNullException(nameof(logger)); } // ───────────────────────────────────────────────────────────────── // GetMonitorAsync // ───────────────────────────────────────────────────────────────── /// /// Returns all active jobs with latest execution state. /// NextRunOn is computed in BLL for jobs where the DB value is null. /// Builds summary counts and returns SchedulerMonitorResponseDTO. /// public async Task GetMonitorAsync( LoginDTO login, CancellationToken ct = default) { GB5Trace.Step("get-scheduler-monitor"); try { var jobs = await _dal.LoadSchedulerMonitorAsync(login, ct) .ConfigureAwait(false); var nowIst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, IST); // Compute NextRunOn in BLL for jobs with no DB value. // SL never calls ComputeNextRun directly — this is the correct layer. foreach (var job in jobs.Where(j => j.NextRunOn == null)) { try { job.NextRunOn = ComputeNextRun( job.CronExpression, job.OccursAtTime, job.EveryNoOfHours, job.EveryNoOfMinutes, job.LastRunOn, job.JobId, nowIst); } catch (Exception ex) { _logger.LogWarning(ex, "Monitor NextRun compute failed | JobId={JobId}", job.JobId); } } // Build summary counts using shared SchedulerExecutionStatus constants var summary = new SchedulerMonitorSummaryDTO { Total = jobs.Count, Pending = jobs.Count(j => j.ExecutionStatus == SchedulerExecutionStatus.Pending), Running = jobs.Count(j => j.ExecutionStatus == SchedulerExecutionStatus.InProgress), Success = jobs.Count(j => j.ExecutionStatus == SchedulerExecutionStatus.Success), Failed = jobs.Count(j => j.ExecutionStatus == SchedulerExecutionStatus.Failed), DLQ = jobs.Count(j => j.ExecutionStatus == SchedulerExecutionStatus.DLQ), NotStarted = jobs.Count(j => j.ExecutionStatus == null) }; _logger.LogInformation( "Scheduler monitor loaded | Total={Total} Running={Running} Failed={Failed} DLQ={DLQ}", summary.Total, summary.Running, summary.Failed, summary.DLQ); return new SchedulerMonitorResponseDTO { Summary = summary, Jobs = jobs }; } catch (Exception ex) { GB5Trace.MarkFailed("get-scheduler-monitor-failed", ex); _logger.LogError(ex, "GetMonitorAsync failed"); throw; } } // ───────────────────────────────────────────────────────────────── // ComputeNextRun (private — monitor-specific overload) // Takes individual fields from SchedulerMonitorDTO — does NOT // depend on SchedulerTaskDTO (wrong DTO type for a monitor read). // Uses SchedulerCronHelper.NormalizeQuartz shared with main BLL. // ───────────────────────────────────────────────────────────────── private DateTime? ComputeNextRun( string? cronExpression, string? occursAtTime, int everyNoOfHours, int everyNoOfMinutes, DateTime? lastRunOn, int jobId, DateTime nowIst) { // ── CRON based schedule ─────────────────────────────────────── if (!string.IsNullOrWhiteSpace(cronExpression)) { if (_warnedJobIds.TryAdd(jobId, true) && (!string.IsNullOrWhiteSpace(occursAtTime) || everyNoOfMinutes > 0 || everyNoOfHours > 0)) { _logger.LogWarning( "Monitor: multiple schedule fields set — CronExpression takes priority | JobId={JobId}", jobId); } var normalized = SchedulerCronHelper.NormalizeQuartz(cronExpression); var cron = CronExpression.Parse(normalized, CronFormat.IncludeSeconds); var baseUtc = TimeZoneInfo.ConvertTimeToUtc(nowIst, IST); var nextUtc = cron.GetNextOccurrence(baseUtc); return nextUtc.HasValue ? TimeZoneInfo.ConvertTimeFromUtc(nextUtc.Value, IST) : null; } // ── Fixed daily execution time ──────────────────────────────── if (TimeSpan.TryParse(occursAtTime, out var time)) { var next = nowIst.Date.Add(time); return next <= nowIst ? next.AddDays(1) : next; } // ── Interval based execution ────────────────────────────────── if (everyNoOfMinutes > 0 || everyNoOfHours > 0) { // ✅ LASTRUNON is stored in UTC — convert to IST before comparing/arithmetic // against nowIst (IST). Same fix as SchedulerTaskServiceBLL.ComputeNextRun. var lastRunIst = lastRunOn.HasValue ? TimeZoneInfo.ConvertTimeFromUtc(lastRunOn.Value, IST) : (DateTime?)null; var baseTime = lastRunIst ?? nowIst.AddHours(-everyNoOfHours).AddMinutes(-everyNoOfMinutes); var next = baseTime .AddHours(everyNoOfHours) .AddMinutes(everyNoOfMinutes); // Advance until future — prevents overdue jobs appearing as immediate while (next <= nowIst) next = next .AddHours(everyNoOfHours) .AddMinutes(everyNoOfMinutes); return next; } return null; } } }