using Dapr.Client; using GB5Shared.DTO.Framework.Login; using GB5Shared.Telemetry; using JobEngineDAL.Interfaces; using JobEngineSL.Hubs; using Microsoft.AspNetCore.SignalR; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Newtonsoft.Json; using Quartz; namespace JobEngineSL.Services { // Quartz IJob implementation. One instance per job execution. // Reads MJOBDEFINE from DB, calls HTTP endpoint (MWEBSERVICE) or Dapr topic, // records TJOBEXECUTION, and pushes SignalR events per-tenant. // ISCONCURRENT flag is applied per-job in QuartzSyncService — no class-level attribute. public class QuartzJobExecutor : IJob { private readonly IServiceScopeFactory _scopeFactory; private readonly IHubContext _hub; private readonly IHttpClientFactory _httpFactory; private readonly DaprClient _dapr; private readonly ILogger _logger; public QuartzJobExecutor( IServiceScopeFactory scopeFactory, IHubContext hub, IHttpClientFactory httpFactory, DaprClient dapr, ILogger logger) { _scopeFactory = scopeFactory; _hub = hub; _httpFactory = httpFactory; _dapr = dapr; _logger = logger; } public async Task Execute(IJobExecutionContext ctx) { var jobId = ctx.JobDetail.JobDataMap.GetInt("JobId"); var tenantId = ctx.JobDetail.JobDataMap.GetInt("TenantId"); var correlationId = Guid.NewGuid().ToString("N"); long executionId = 0; using var scope = _scopeFactory.CreateScope(); var jobDefineDAL = scope.ServiceProvider.GetRequiredService(); var executionDAL = scope.ServiceProvider.GetRequiredService(); // Build a minimal login for system-level DAL calls var sysLogin = new LoginDTO { ClientId = tenantId, UserId = 0 }; try { var job = await jobDefineDAL.GetJobDetailAsync(jobId, sysLogin, ctx.CancellationToken) .ConfigureAwait(false); if (job == null) { _logger.LogWarning("QuartzJobExecutor: JobId {JobId} not found — skipping", jobId); return; } GB5Trace.Step("quartz-execute-start", new { jobId, job.JobName, correlationId }); executionId = await executionDAL.InsertExecutionAsync(new JobEngineDAL.DTOs.JobExecutionDTO { JobId = jobId, JobName = job.JobName, TriggerType = "SCHEDULED", ScheduledFor = ctx.ScheduledFireTimeUtc?.UtcDateTime, CorrelationId = correlationId }, sysLogin, ctx.CancellationToken).ConfigureAwait(false); var hubGroup = $"jobengine:client:{tenantId}"; await _hub.Clients.Group(hubGroup) .ReceiveJobStarted(jobId, job.JobName, correlationId).ConfigureAwait(false); GB5Trace.Step("quartz-execute-call", new { job.DaprPubSubTopic, job.UriTemplate }); string? result = null; if (!string.IsNullOrWhiteSpace(job.DaprPubSubTopic)) { var payload = new { JobId = jobId, TenantId = tenantId, CorrelationId = correlationId }; await _dapr.PublishEventAsync("pubsub", job.DaprPubSubTopic, payload, ctx.CancellationToken).ConfigureAwait(false); result = $"Published to {job.DaprPubSubTopic}"; } else if (!string.IsNullOrWhiteSpace(job.UriTemplate)) { var http = _httpFactory.CreateClient("jobengine"); var uri = BuildUri(job.UriTemplate, job.UriParameterValue); var method = string.Equals(job.MethodType, "POST", StringComparison.OrdinalIgnoreCase) ? HttpMethod.Post : HttpMethod.Get; using var req = new HttpRequestMessage(method, uri); using var resp = await http.SendAsync(req, ctx.CancellationToken).ConfigureAwait(false); resp.EnsureSuccessStatusCode(); result = await resp.Content.ReadAsStringAsync(ctx.CancellationToken).ConfigureAwait(false); } var startedAt = ctx.FireTimeUtc.UtcDateTime; var durationMs = (long)(DateTime.UtcNow - startedAt).TotalMilliseconds; await executionDAL.CompleteExecutionAsync(executionId, "SUCCESS", result, null, sysLogin, ctx.CancellationToken) .ConfigureAwait(false); await jobDefineDAL.UpdateLastRunAsync(jobId, ctx.NextFireTimeUtc?.UtcDateTime, sysLogin, ctx.CancellationToken) .ConfigureAwait(false); await _hub.Clients.Group(hubGroup) .ReceiveJobCompleted(jobId, job.JobName, durationMs, null).ConfigureAwait(false); GB5Trace.Step("quartz-execute-done", new { jobId, durationMs }); } catch (Exception ex) { GB5Trace.MarkFailed("quartz-execute-failed", ex); _logger.LogError(ex, "QuartzJobExecutor failed for JobId {JobId}", jobId); // Use CancellationToken.None — ctx.CancellationToken may already be cancelled if (executionId > 0) await executionDAL.CompleteExecutionAsync( executionId, ex is OperationCanceledException ? "CANCELLED" : "FAILED", null, ex.Message, sysLogin, CancellationToken.None).ConfigureAwait(false); await _hub.Clients.Group($"jobengine:client:{tenantId}") .ReceiveJobFailed(jobId, $"Job-{jobId}", ex.Message, 0, false).ConfigureAwait(false); if (ex is OperationCanceledException) throw; throw new JobExecutionException(ex, false); } } private static string BuildUri(string uriTemplate, string? paramValue) { if (string.IsNullOrWhiteSpace(paramValue)) return uriTemplate; return uriTemplate.Replace("{Param}", Uri.EscapeDataString(paramValue)); } } }