using Dapper; using GB5Shared.Connection; using GB5Shared.DTO.Framework.CommonConfig; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ServerConfig; using GB5Shared.QueryExecutor; using GB5Shared.Telemetry; using Microsoft.Data.SqlClient; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using Npgsql; using System.Data; using static GB5Shared.GB5Constant.Constant; namespace GB5Shared.QueueReader { // Abstract base for all module queue processors. // Subclass registers handlers via the DI container; the base polls TJOBQUEUE atomically, // claims a batch, and dispatches each item to the matching IJobQueueHandler. // Exponential backoff: 30s → 2m → 8m → 30m → 2h → DLQ after MaxRetries. public abstract class BaseJobQueueProcessor : BackgroundService { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; private static readonly TimeSpan[] BackoffSchedule = [ TimeSpan.FromSeconds(30), TimeSpan.FromMinutes(2), TimeSpan.FromMinutes(8), TimeSpan.FromMinutes(30), TimeSpan.FromHours(2) ]; private const string ClaimSql = @" UPDATE TOP(@BatchSize) Q SET STATUS = 'PROCESSING', PROCESSINGSTARTED = GETUTCDATE(), UPDATEDDATE = GETUTCDATE() OUTPUT INSERTED.QUEUEID, INSERTED.TENANTID, INSERTED.JOBEXECUTIONID, INSERTED.MESSAGETYPE, INSERTED.PAYLOAD, INSERTED.PRIORITY, INSERTED.STATUS, INSERTED.SCHEDULEDFOR, INSERTED.MAXRETRIES, INSERTED.RETRYATTEMPT, INSERTED.MESSAGEID, INSERTED.CORRELATIONID, INSERTED.SOURCESERVICE, INSERTED.ERRORDETAILS, INSERTED.CREATEDDATE FROM TJOBQUEUE Q WITH (ROWLOCK, READPAST) WHERE STATUS IN ('PENDING', 'FAILED') AND TENANTID = @TenantId AND MESSAGETYPE IN @MessageTypes AND SCHEDULEDFOR <= GETUTCDATE() AND (NEXTATTEMPTAFTER IS NULL OR NEXTATTEMPTAFTER <= GETUTCDATE())"; private const string CompleteSql = "UPDATE TJOBQUEUE SET STATUS='COMPLETED', UPDATEDDATE=GETUTCDATE() WHERE QUEUEID=@QueueId AND TENANTID=@TenantId"; private const string RetrySql = "UPDATE TJOBQUEUE SET STATUS='FAILED', RETRYATTEMPT=RETRYATTEMPT+1, NEXTATTEMPTAFTER=@Next, ERRORDETAILS=@Error, UPDATEDDATE=GETUTCDATE() WHERE QUEUEID=@QueueId AND TENANTID=@TenantId"; private const string MoveToDLQSql = "UPDATE TJOBQUEUE SET STATUS='DLQ', ERRORDETAILS=@Error, UPDATEDDATE=GETUTCDATE() WHERE QUEUEID=@QueueId AND TENANTID=@TenantId"; protected BaseJobQueueProcessor(IServiceScopeFactory scopeFactory, ILogger logger) { _scopeFactory = scopeFactory; _logger = logger; } // Message types this processor handles — subclass provides the list protected abstract IReadOnlyList MessageTypes { get; } // Poll interval — override in subclass to change frequency protected virtual TimeSpan PollInterval => TimeSpan.FromSeconds(5); // Batch size per poll per tenant protected virtual int BatchSize => 10; // Override to restrict which tenants this processor polls. // Default loads all active tenants from MSERVERCONFIG — same pattern as OutBoxPollerService. protected virtual async Task> GetTenantLoginsAsync( IServiceScope scope, CancellationToken ct) { var appConnection = scope.ServiceProvider.GetRequiredService(); var dbOptions = scope.ServiceProvider.GetRequiredService>(); var systemConn = await appConnection.Gb5SystemConnectionString().ConfigureAwait(false); var dbType = dbOptions.CurrentValue.DataBaseType; using IDbConnection conn = dbType switch { DBTYPE.SQL => new SqlConnection(systemConn), DBTYPE.POSTGRESQL => new NpgsqlConnection(systemConn), _ => throw new NotSupportedException($"Unsupported DB type: {dbType}") }; // Tracker §47/§48 — an extra opt-in predicate a subclass can layer on top of the // shared TenantQuery below, WITHOUT changing TenantQuery itself (which is also used // by PAYJobQueueProcessor/PAYQueueProcessor/MMJobQueueProcessor — none of which are // "the scheduler" and must keep enumerating every STATUS=1 tenant exactly as before). // Only SchedulerActionEventQueueProcessor overrides this, requiring // MSERVERCONFIG.SCHEDULERENABLED=1 (default 0/Disabled) per tenant. var sql = string.IsNullOrWhiteSpace(AdditionalTenantFilterSql) ? TenantQuery : $"{TenantQuery}\n AND {AdditionalTenantFilterSql}"; var rows = await SqlMapper.QueryAsync(conn, sql).ConfigureAwait(false); return rows.Select(BuildLogin); } /// Extra SQL predicate (no leading AND) appended to the default /// tenant query — null/empty by default (every /// existing subclass's behavior is unchanged). Override only in a subclass whose polling /// genuinely represents "the scheduler" and should respect the per-tenant /// MSERVERCONFIG.SCHEDULERENABLED opt-in (tracker §47/§48) — never in a business-domain /// queue processor like PAYJobQueueProcessor/PAYQueueProcessor/MMJobQueueProcessor. protected virtual string? AdditionalTenantFilterSql => null; // BUG FIX (2026-08-17): ApplicationConnection.DBConnectionStringCached resolves a // connection by treating LoginDTO.DatabaseName as the MSERVERCONFIG.CONNECTIONNAME // lookup key (see FormConnectionString -> DatabaseConnectionObjectConnectionName) — // it is NOT the raw database name. Every other tenant-resolving path in this codebase // (e.g. SchedulerJobExecutorQuartzJob, SysJobExecutorQuartzJob) sets // LoginDTO.DatabaseName = tenant.ConnectionName accordingly. This method previously // set it to tenant.DatabaseName instead, so every BaseJobQueueProcessor subclass using // the default tenant enumeration (PAYJobQueueProcessor, PAYQueueProcessor, and now // SchedulerJobQueueProcessor) failed to resolve a connection for any tenant whose // CONNECTIONNAME differs from its raw DATABASENAME (e.g. "UNISOFTGB5" vs "unisoftgb4"). private static LoginDTO BuildLogin(ServerConfigDTO tenant) => new() { UserId = -1, ClientId = tenant.ClientId, ConnectionDatabaseName = tenant.ConnectionName, DatabaseName = tenant.ConnectionName, DatabaseType = tenant.DbType }; private const string TenantQuery = @" 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'"; 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(); IEnumerable tenants; try { tenants = await GetTenantLoginsAsync(scope, ct).ConfigureAwait(false); } catch (Exception ex) when (ex is not OperationCanceledException) { _logger.LogError(ex, "{Processor} failed to load tenant logins", GetType().Name); return; } var qe = scope.ServiceProvider.GetRequiredService(); var handlers = scope.ServiceProvider.GetServices().ToList(); foreach (var tenantLogin in tenants) { try { var claimed = (await qe.QueryAsync(tenantLogin, ClaimSql, new { BatchSize = BatchSize, TenantId = tenantLogin.ClientId, MessageTypes = MessageTypes }, cancellationToken: ct).ConfigureAwait(false)).ToList(); // Process items in parallel — each item independent var tasks = claimed.Select(item => ProcessItemAsync(item, tenantLogin, handlers, qe, ct)); await Task.WhenAll(tasks).ConfigureAwait(false); } catch (Exception ex) when (ex is not OperationCanceledException) { _logger.LogError(ex, "{Processor} poll failed for tenant {TenantId}", GetType().Name, tenantLogin.ClientId); } } } private async Task ProcessItemAsync( JobQueueItemDTO item, LoginDTO login, IEnumerable handlers, IQueryExecutor qe, CancellationToken ct) { var handler = handlers.FirstOrDefault(h => string.Equals(h.MessageType, item.MessageType, StringComparison.OrdinalIgnoreCase)); try { if (handler == null) throw new InvalidOperationException($"No handler registered for MessageType '{item.MessageType}'"); GB5Trace.Step("queue-process", new { item.QueueId, item.MessageType }); await handler.HandleAsync(item.Payload, login, item.QueueId, ct).ConfigureAwait(false); await qe.ExecuteAsync(login, CompleteSql, new { item.QueueId, TenantId = item.TenantId }, cancellationToken: ct).ConfigureAwait(false); _logger.LogInformation("Queue item {QueueId} ({MessageType}) completed", item.QueueId, item.MessageType); } catch (OperationCanceledException) { throw; } catch (Exception ex) { GB5Trace.MarkFailed($"queue-process-failed:{item.MessageType}", ex); _logger.LogError(ex, "Queue item {QueueId} ({MessageType}) failed attempt {Attempt}", item.QueueId, item.MessageType, item.RetryAttempt); if (item.RetryAttempt < item.MaxRetries) { var backoffIndex = Math.Min(item.RetryAttempt, BackoffSchedule.Length - 1); var next = DateTime.UtcNow.Add(BackoffSchedule[backoffIndex]); await qe.ExecuteAsync(login, RetrySql, new { item.QueueId, TenantId = item.TenantId, Next = next, Error = ex.Message }, cancellationToken: ct) .ConfigureAwait(false); } else { await qe.ExecuteAsync(login, MoveToDLQSql, new { item.QueueId, TenantId = item.TenantId, Error = ex.Message }, cancellationToken: ct) .ConfigureAwait(false); _logger.LogWarning("Queue item {QueueId} ({MessageType}) moved to DLQ after {Attempts} retries", item.QueueId, item.MessageType, item.RetryAttempt); } } } } }