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