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