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; }
}
}