using Dapper; using FrameworkBLL.SchedulerTaskGenerator; using FrameworkDAL.DTO.SchedulerTaskGenerator; using GB5Shared.Connection; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ServerConfig; using Microsoft.Data.SqlClient; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using System; using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; namespace FrameworkSL.Controllers.SchedulerTaskGenerator { /// /// Polls TJOBEXECUTION directly for due (STATUS=0, NEXTATTEMPTON elapsed) scheduled-job /// executions and delivers each via ISchedulerExecutionDeliveryService. Registered as /// AddHostedService in Program.cs, replacing SchedulerJobQueueProcessor. /// /// Replaces the previous TJOBQUEUE-based design, where SchedulerJobExecutorQuartzJob /// enqueued a "SCHEDULED_WEBSERVICE_CALL" TJOBQUEUE row for a BaseJobQueueProcessor /// subclass to pick up. That design kept two independent, redundant retry counters for /// the same failure — TJOBQUEUE.RETRYATTEMPT (governing whether the queue re-delivered) /// and TJOBEXECUTION.RETRYNUMBER (SchedulerTaskServiceBLL.ApplyRetryOrDeadLetterAsync, /// governing DLQ) — both fed from the same MJOBDEFINE.MaxRetry ceiling and both updated /// on every single failed attempt. It also meant a job's TJOBEXECUTION row sat at /// STATUS=0 (Pending), not 2 (InProgress), for the whole time it waited out a TJOBQUEUE /// retry backoff — so the ISCONCURRENT=0 guard (which only ever blocked on STATUS=2) /// failed to stop Quartz firing a brand-new, overlapping execution during that window. /// /// This version claims and retries directly against TJOBEXECUTION — one table, one /// retry counter (RETRYNUMBER), one DLQ flag (STATUS=5). The concurrency guard in /// SchedulerTaskGeneratorQB.CreatePendingExecution now blocks on STATUS IN (0,2), /// closing that gap. This also revives SchedulerTaskServiceBLL.ReapTimedOutExecutionsAsync, /// which had no caller left once the old SchedulerTaskManager poll loop went dead. /// /// TJOBQUEUE itself is untouched by this change — other modules still use it /// independently via GB5Shared/QueueReader's BaseJobQueueProcessor/IJobQueueEnqueuer. /// public sealed class SchedulerExecutionPoller : BackgroundService { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; // 0 = all tenants (single-server / dev default); mirrors QuartzSyncService's // "Scheduler:ServerId" config key — must match on every host for consistent coverage. private readonly int _schedulerServerId; private static readonly TimeSpan PollInterval = TimeSpan.FromSeconds(5); private const int BatchSize = 10; public SchedulerExecutionPoller( IServiceScopeFactory scopeFactory, ILogger logger, IConfiguration config) { _scopeFactory = scopeFactory; _logger = logger; _schedulerServerId = (config ?? throw new ArgumentNullException(nameof(config))) .GetValue("Scheduler:ServerId", 0); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { using var timer = new PeriodicTimer(PollInterval); while (await timer.WaitForNextTickAsync(stoppingToken).ConfigureAwait(false)) { await ProcessAllTenantsAsync(stoppingToken).ConfigureAwait(false); } } private async Task ProcessAllTenantsAsync(CancellationToken ct) { using var scope = _scopeFactory.CreateScope(); var appConnection = scope.ServiceProvider.GetRequiredService(); List tenants; try { tenants = await LoadActiveTenantsAsync(appConnection, _schedulerServerId).ConfigureAwait(false); } catch (Exception ex) when (ex is not OperationCanceledException) { _logger.LogError(ex, "SchedulerExecutionPoller: failed to load tenant logins"); return; } var service = scope.ServiceProvider.GetRequiredService(); var delivery = scope.ServiceProvider.GetRequiredService(); foreach (var tenant in tenants) { var login = BuildLogin(tenant); try { // Restores timeout reaping — dead since the old SchedulerTaskManager // poll loop (the only caller of ReapTimedOutExecutionsAsync) stopped // being registered as a hosted service. await service.ReapTimedOutExecutionsAsync(login, ct).ConfigureAwait(false); var claimed = await service.ClaimDueExecutionsAsync(BatchSize, login, ct) .ConfigureAwait(false); var tasks = claimed.Select(item => DeliverOneAsync(item, login, delivery, ct)); await Task.WhenAll(tasks).ConfigureAwait(false); } catch (Exception ex) when (ex is not OperationCanceledException) { _logger.LogError(ex, "SchedulerExecutionPoller: poll failed for tenant {TenantId}", tenant.ClientId); } } } private async Task DeliverOneAsync( SchedulerTaskDTO item, LoginDTO login, ISchedulerExecutionDeliveryService delivery, CancellationToken ct) { try { await delivery.DeliverAsync(item, login, ct).ConfigureAwait(false); } catch (OperationCanceledException) { throw; } catch (Exception ex) { // DeliverAsync already writes failure to TJOBEXECUTION internally via // UpdateExecutionResultAsync before it can throw — this only catches a // truly unexpected failure in the delivery plumbing itself (e.g. DI // resolution), so one bad item doesn't take down the whole poll tick. _logger.LogError(ex, "SchedulerExecutionPoller: unhandled delivery failure | JobExecutionId={Id}", item.JobExecutionId); } } // Tenant enumeration — same MSERVERCONFIG/MSERVER query used by // QuartzSyncService/SysJobExecutorQuartzJob/BaseJobQueueProcessor. No shared // abstraction exists for this in the codebase; duplicated here following the // same established convention (see QuartzSyncService.LoadActiveTenantsAsync). private static async Task> LoadActiveTenantsAsync( IApplicationConnection appConnection, int serverId) { var systemConn = await appConnection.Gb5SystemConnectionString().ConfigureAwait(false); const string tenantSql = @" SELECT SERVERCONFIG1.CLIENTID AS ClientId, SERVERCONFIG1.DATABASENAME AS DatabaseName, SERVERCONFIG1.DATABASETYPE AS DbType, SERVERCONFIG1.CONNECTIONNAME AS ConnectionName FROM MSERVERCONFIG SERVERCONFIG1 JOIN MSERVER SERVER1 ON SERVERCONFIG1.SERVERID = SERVER1.SERVERID WHERE SERVERCONFIG1.STATUS = 1 AND SERVERCONFIG1.CONNECTIONNAME <> 'ACTIVITI' AND (@ServerId = 0 OR SERVERCONFIG1.SERVERID = @ServerId) -- Tracker §47/§48 — per-tenant opt-in, default Disabled; see -- QuartzSyncService.LoadActiveTenantsAsync's own identical comment. AND SERVERCONFIG1.SCHEDULERENABLED = 1"; using var connection = new SqlConnection(systemConn); return (await SqlMapper.QueryAsync( connection, tenantSql, new { ServerId = serverId }).ConfigureAwait(false)).ToList(); } // ApplicationConnection.DBConnectionStringCached resolves connection strings by // treating LoginDTO.DatabaseName as a lookup key against MSERVERCONFIG.CONNECTIONNAME // (confusingly named, but confirmed elsewhere — see QuartzSyncService.BuildLogin / // ActionProcessorWorker.BuildLogin's own comment on the same gotcha). Using the // literal tenant.DatabaseName here would fail resolution for every tenant whose // CONNECTIONNAME differs from its raw DATABASENAME. private static LoginDTO BuildLogin(ServerConfigDTO tenant) => new() { UserId = -1, UserName = "SchedulerService", ClientId = tenant.ClientId, DatabaseName = tenant.ConnectionName, ConnectionDatabaseName = tenant.ConnectionName }; } }