using FrameworkBLL.SchedulerTaskGeneratorPublisher; using FrameworkDAL.DTO.SchedulerTaskGenerator; using GB5Shared.DTO.Framework.Login; using GB5Shared.Telemetry; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using System; using System.Collections.Generic; using System.Diagnostics; using System.Linq; using System.Threading; using System.Threading.Tasks; namespace FrameworkBLL.SchedulerTaskGenerator { // GB5ActivitySources.Scheduler is the canonical source for all scheduler spans. // Keeping this shim so call-sites that use SchedulerTelemetry.ActivitySource still compile. internal static class SchedulerTelemetry { public static ActivitySource ActivitySource => GB5ActivitySources.Scheduler; } public interface ISchedulerTaskManager { /// /// Loads all jobs due for execution at the current IST time. /// Task> LoadReadyTasksAsync( LoginDTO login, CancellationToken ct = default); /// /// Loads due jobs, marks each as in-progress, publishes to queue. /// Returns a summary of what was dispatched. /// Task ExecuteReadyTasksAsync( LoginDTO login, CancellationToken ct = default); } /// /// Orchestrates loading and publishing of ready scheduler tasks. /// Hosted by SchedulerBackgroundService or SchedulerQuartzJob — NOT a BackgroundService itself. /// /// Responsibilities: /// 1. Load jobs due now from BLL /// 2. Mark each as InProgress in DB /// 3. Publish JobExecutionId to queue /// 4. Return dispatch summary /// /// NOT responsible for: /// - Retry logic (belongs in the consumer after job execution) /// - DLQ routing (belongs in the consumer after MaxRetry exhausted) /// public sealed class SchedulerTaskManager : ISchedulerTaskManager { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; private static readonly TimeZoneInfo IST = TimeZoneInfo.FindSystemTimeZoneById("India Standard Time"); public SchedulerTaskManager( IServiceScopeFactory scopeFactory, ILogger logger) { _scopeFactory = scopeFactory ?? throw new ArgumentNullException(nameof(scopeFactory)); _logger = logger ?? throw new ArgumentNullException(nameof(logger)); } // ───────────────────────────────────────────────────────────────── // LoadReadyTasksAsync — read-only, no side effects // ───────────────────────────────────────────────────────────────── /// /// Returns all jobs due now. Read-only — no DB writes, no queue publish. /// Useful for monitoring / health-check endpoints. /// public async Task> LoadReadyTasksAsync( LoginDTO login, CancellationToken ct = default) { using var activity = SchedulerTelemetry.ActivitySource .StartActivity("Scheduler.LoadReadyTasks"); activity?.SetTag("db.name", login.DatabaseName); activity?.SetTag("scheduler.stage", "load"); using var scope = _scopeFactory.CreateScope(); var service = scope.ServiceProvider.GetRequiredService(); var nowIst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, IST); var tasks = await service.LoadExecutableTasksAsync(login, nowIst, ct) .ConfigureAwait(false); activity?.SetTag("scheduler.task.count", tasks.Count); return tasks; } // ───────────────────────────────────────────────────────────────── // ExecuteReadyTasksAsync — marks + publishes // ───────────────────────────────────────────────────────────────── /// /// Single DB load → mark each job in-progress → publish to queue. /// ✅ LoadExecutableTasksAsync is called ONCE per cycle (was called twice before). /// ✅ Retry/DLQ logic removed — belongs in the consumer, not the publisher. /// ✅ Publish failures are logged and skipped — job stays pending for next poll. /// public async Task ExecuteReadyTasksAsync( LoginDTO login, CancellationToken ct = default) { using var activity = SchedulerTelemetry.ActivitySource .StartActivity("Scheduler.ExecuteReadyTasks"); activity?.SetTag("db.name", login.DatabaseName); activity?.SetTag("scheduler.stage", "execute"); using var scope = _scopeFactory.CreateScope(); var service = scope.ServiceProvider.GetRequiredService(); var publisher = scope.ServiceProvider.GetRequiredService(); var nowIst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, IST); // ✅ Reap timed-out InProgress executions before loading ready jobs, so a job // stuck past its TIMEOUTSECONDS is retried/DLQ'd instead of blocking forever. try { var reaped = await service.ReapTimedOutExecutionsAsync(login, ct).ConfigureAwait(false); if (reaped > 0) { activity?.SetTag("scheduler.timedout.reaped", reaped); _logger.LogWarning("Scheduler reaped {Count} timed-out execution(s)", reaped); } } catch (Exception ex) { // Non-fatal — a reap failure must not block the rest of the poll cycle. _logger.LogError(ex, "ReapTimedOutExecutionsAsync failed — continuing poll cycle"); } // ✅ Single DB call — LoadExecutableTasksAsync already filters NextRunOn <= nowIst var tasks = await service.LoadExecutableTasksAsync(login, nowIst, ct) .ConfigureAwait(false); var summary = new ExecutionSummaryDTO { Total = tasks.Count }; if (tasks.Count == 0) { _logger.LogInformation( "Scheduler cycle — no jobs due | Time={NowIst:HH:mm:ss}", nowIst); return summary; } foreach (var task in tasks) { if (ct.IsCancellationRequested) break; using var jobActivity = SchedulerTelemetry.ActivitySource .StartActivity("Scheduler.PublishJob"); jobActivity?.SetTag("job.id", task.JobId); jobActivity?.SetTag("job.name", task.JobName); try { // 1. Mark execution as InProgress in DB (also saves NextRunOn). // ✅ Capture the id MarkExecutionStartedAsync actually just inserted — // task.JobExecutionId (loaded before this call) was stale: 0 for a // first-ever run, or a pre-existing row's id on retry. Publishing that // stale value meant the consumer's UpdateExecutionToInProgressAsync/ // GenerateAsync always operated on the wrong (or a nonexistent) row. var newJobExecutionId = await service.MarkExecutionStartedAsync(task.JobId, login, ct) .ConfigureAwait(false); task.JobExecutionId = newJobExecutionId; jobActivity?.SetTag("job.execution.id", task.JobExecutionId); if (newJobExecutionId == 0) { // Blocked by the ISCONCURRENT guard (a prior execution is still // running) — nothing was inserted, so there is nothing to publish. summary.Skipped++; _logger.LogInformation( "Scheduler skipped — already in progress | JobId={JobId}", task.JobId); jobActivity?.SetStatus(ActivityStatusCode.Ok); continue; } // 2. Publish full task to queue — consumer uses JobExecutionId to route // the correct TJOBEXECUTION row, and Actions to know what to run. // ✅ Previously only ActionCount was published, never the Actions list // itself — the consumer's `foreach (var action in task.Actions)` // always iterated zero times, so no action was ever generated. await publisher.PublishAsync("Scheduler.Ready", new { task.JobExecutionId, task.JobId, task.JobName, task.TenantId, task.NextRunOn, task.Actions, ActionCount = task.Actions?.Count ?? 0 }).ConfigureAwait(false); summary.Published++; _logger.LogInformation( "Scheduler published | JobId={JobId} ExecutionId={ExecutionId} NextRunOn={NextRunOn:HH:mm:ss}", task.JobId, task.JobExecutionId, task.NextRunOn); jobActivity?.SetStatus(ActivityStatusCode.Ok); } catch (Exception ex) { // ✅ Publish failure — log and skip, do NOT retry here. // The job remains STATUS=Pending in DB so the next poll cycle picks it up. // Retry (execution failure) is handled by the consumer after the job runs. summary.Skipped++; jobActivity?.SetStatus(ActivityStatusCode.Error, ex.Message); _logger.LogError(ex, "Scheduler publish failed — skipping | JobId={JobId} ExecutionId={ExecutionId}", task.JobId, task.JobExecutionId); } } activity?.SetTag("scheduler.published", summary.Published); activity?.SetTag("scheduler.skipped", summary.Skipped); _logger.LogInformation( "Scheduler cycle complete | Total={Total} Published={Published} Skipped={Skipped}", summary.Total, summary.Published, summary.Skipped); return summary; } } /// /// Summary of a single scheduler execution cycle. /// public sealed class ExecutionSummaryDTO { /// Total jobs due for execution this cycle. public int Total { get; set; } /// Jobs successfully marked in-progress and published to queue. public int Published { get; set; } /// Jobs skipped due to publish error (will retry on next poll). public int Skipped { get; set; } /// Jobs loaded but not yet due — always 0 after BLL filter fix. public int Waiting { get; set; } } }