using FrameworkDAL.CustomCode.SchedulerTaskGenerator; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.JobEngine; using GB5Shared.Telemetry; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Quartz; using System; using System.Diagnostics; using System.Threading.Tasks; namespace FrameworkBLL.SchedulerTaskGenerator { /// /// Per-job Quartz executor — replaces the old generic 30s-poll SchedulerQuartzJob. /// Quartz fires this directly for ONE job (via its own native, persistent, per-job /// trigger registered by QuartzSyncService) — there is no "which jobs are due" /// poll/scan here, Quartz already decided that by firing this instance at all. /// /// This executor's own responsibility ends at "recorded as Pending in TJOBEXECUTION" — /// actual delivery (the HTTP call, retry/backoff, DLQ) is driven straight off /// TJOBEXECUTION by SchedulerExecutionPoller/SchedulerExecutionDeliveryService /// (FrameworkSL/Controllers/SchedulerTaskGenerator). There is deliberately no /// TJOBQUEUE hop here any more — see SchedulerExecutionPoller's header comment for /// why (TJOBQUEUE.RETRYATTEMPT and TJOBEXECUTION.RETRYNUMBER used to be two /// independent, redundant retry counters for the same failure). /// [DisallowConcurrentExecution] public sealed class SchedulerJobExecutorQuartzJob : IJob { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; private static readonly string HostInstance = $"{Environment.MachineName}:{Environment.ProcessId}"; public SchedulerJobExecutorQuartzJob( IServiceScopeFactory scopeFactory, ILogger logger) { _scopeFactory = scopeFactory ?? throw new ArgumentNullException(nameof(scopeFactory)); _logger = logger ?? throw new ArgumentNullException(nameof(logger)); } public async Task Execute(IJobExecutionContext context) { var ct = context.CancellationToken; // Generated here — BEFORE the root span opens and BEFORE the TJOBEXECUTION insert — // so the same value can be tagged on the root span (gb5.correlation.id, visible and // searchable inside the trace) AND written as TJOBEXECUTION.CORRELATIONID, rather // than letting SQL generate it independently via NEWID() after the fact. // TJOBQUEUE.CORRELATIONID then matches automatically — SchedulerExecutionDeliveryService // already carries item.CorrelationId (read back from TJOBEXECUTION) into // IJobQueueEnqueuer.EnqueueAsync. Upper-invariant to match SQL Server's own NEWID() // rendering and the existing CORRELATIONID values already in the DB. var correlationId = Guid.NewGuid().ToString().ToUpperInvariant(); // Root span for the entire distributed pipeline this one Quartz fire kicks off — // poll-claim, HTTP delivery, TJOBQUEUE enqueue, Dapr publish/subscribe, and (if an // action fires) all the way through RabbitMQ to the actual send. Every later hop // reconstructs its parent from the traceparent/tracestate persisted below (first // to TJOBEXECUTION, then carried through each subsequent message payload) — NOT // from Activity.Current, since every later hop runs in a separate async // continuation (a different poll tick, a different process boundary) where no // ambient Activity survives. Disposing this activity immediately after Quartz's // synchronous firing work is correct: a root span's own duration only covers the // operation that started the trace, not every child span linked to it later. // Span name stays plain "jobscheduler" — the CorrelationId is carried as the // gb5.correlation.id TAG instead, visible inside the trace/span detail in Zipkin // and searchable there, without changing the span name itself. using var rootActivity = GB5ActivitySources.Scheduler.StartActivity( "jobscheduler", ActivityKind.Internal); rootActivity?.SetTag("gb5.correlation.id", correlationId); // JobId/TenantId/DatabaseName/ConnectionName were put into the job-level // JobDataMap by QuartzSyncService.SyncJobAsync when this job's trigger was // registered — Quartz merges job-level + trigger-level JobDataMap at fire time, // so a manual "Trigger Now" (Section 5 — trigger-level TriggerType="MANUAL") // is read the same way via context.MergedJobDataMap below. var map = context.MergedJobDataMap; var jobId = map.GetInt("JobId"); var tenantId = map.GetInt("TenantId"); var databaseName = map.GetString("DatabaseName") ?? string.Empty; var connectionName = map.GetString("ConnectionName") ?? databaseName; var triggerType = map.ContainsKey("TriggerType") ? map.GetString("TriggerType") : "SCHEDULED"; rootActivity?.SetTag("gb5.job.id", jobId); rootActivity?.SetTag("gb5.tenant.id", tenantId); rootActivity?.SetTag("gb5.trigger.type", triggerType ?? "SCHEDULED"); var login = new LoginDTO { UserId = -1, UserName = "SchedulerService", ClientId = tenantId, DatabaseName = connectionName, ConnectionDatabaseName = connectionName }; using var scope = _scopeFactory.CreateScope(); var dal = scope.ServiceProvider.GetRequiredService(); var hub = scope.ServiceProvider.GetRequiredService(); var nowIst = TimeZoneInfo.ConvertTimeFromUtc( DateTime.UtcNow, TimeZoneInfo.FindSystemTimeZoneById("India Standard Time")); long jobExecutionId; try { jobExecutionId = await dal.CreatePendingExecutionAsync( jobId, nowIst, triggerType ?? "SCHEDULED", HostInstance, correlationId, traceParent: rootActivity?.Id, traceState: rootActivity?.TraceStateString, login, ct) .ConfigureAwait(false); } catch (Exception ex) { _logger.LogError(ex, "SchedulerJobExecutorQuartzJob: CreatePendingExecution failed | JobId={JobId}", jobId); throw new JobExecutionException(ex, refireImmediately: false); } if (jobExecutionId == 0) { // ISCONCURRENT guard blocked it — a previous run for this job is still // Pending or InProgress and this job doesn't allow concurrent execution. // Normal, not an error. _logger.LogInformation( "SchedulerJobExecutorQuartzJob: skipped | JobId={JobId} (concurrency guard — previous run still pending/in progress)", jobId); return; } var job = await dal.LoadJobDetailsAsync((int)jobExecutionId, login, ct).ConfigureAwait(false); await hub.PushJobStarted(tenantId, new JobStartedDTO { JobId = jobId, JobExecutionId = jobExecutionId, JobName = job?.JobName ?? string.Empty, TriggerType = triggerType, CorrelationId = job?.CorrelationId, StartedAtUtc = DateTime.UtcNow }, ct).ConfigureAwait(false); // No enqueue, no URL resolution here — SchedulerExecutionPoller claims this // row (STATUS 0->2) on its next poll tick and SchedulerExecutionDeliveryService // resolves the webhook URL fresh at delivery time. No next-run bookkeeping // needed either — Quartz's own trigger repeat semantics (SimpleTrigger/ // CronTrigger) fire the next Execute() automatically. } } }