using System; using System.Collections.Generic; using System.Data; using System.Linq; using System.Threading; using System.Threading.Tasks; using Dapper; using FrameworkDAL.CustomCode.GOP; using GB5Shared.Connection; using GB5Shared.DTO.Framework.CommonConfig; 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 Microsoft.Extensions.Options; using Npgsql; using static GB5Shared.GB5Constant.Constant; namespace FrameworkBLL.GOP.Worker { // ============================================================ // GopWorkerService — BackgroundService for GOP queue polling. // // Responsibility: // For every tenant database explicitly allowlisted via // GOP:WorkerAllowedConnectionNames, poll TGopExecutionHeader for // Created/Queued items, acquire an optimistic lock, dispatch to // GopExecutionPipeline, and release the lock in a finally block. // // Multi-tenant design (replaces an earlier single-hardcoded-database // stopgap): as a BackgroundService, this has no per-request LoginDTO to // inherit a DB connection from, unlike every FastEndpoint/BLL call in // this codebase. The real tenant registry lives in MSERVERCONFIG/MSERVER // (queried via the Gb5System database) — this mirrors // JobEngineSL.Services.OutBoxPollerService, a complete, already-deployed // reference implementation of exactly this "enumerate active tenants, // build a LoginDTO per tenant, loop with per-tenant error isolation" // pattern (an earlier version of this comment incorrectly called that // service "also a TODO stub" — verified directly, it is not). // // Two distinct levels of "multi-tenant" apply here, not one: // Level 1 — which PHYSICAL DATABASE to connect to. MSERVERCONFIG can // list many active databases (18 on the shared GB5DEMO box at the // time this was written — other teams' real client/test databases, // not just GB5DEMO). Discovering them is this fix's job. // Level 2 — GET_PENDING_EXECUTIONS/RESET_EXPIRED_LOCKS do NOT filter by // the TENANTID column at all (verified by reading the SQL) — one // connection to one physical database already surfaces every // ClientId that database happens to host. GopExecutionPipeline // already reflects this: it threads a separate `clientId` parameter // (sourced from each execution row's own TENANTID), never // loginDTO.ClientId. So this worker only ever needs one LoginDTO per // PHYSICAL DATABASE, never per ClientId. // // Safety gate: unlike OutBoxPollerService (safe to run for every active // tenant — it only relays already-decided outbound events), a GOP poll // cycle actually EXECUTES flow steps, including arbitrary ApiCall steps // with real external side effects. Discovery (which databases exist) and // permission (which ones this worker may act on) are kept as separate // concerns: GOP:WorkerAllowedConnectionNames is an explicit allowlist of // MSERVERCONFIG.CONNECTIONNAME values, empty/unset by default — the // worker safely no-ops every cycle until an operator deliberately opts a // specific tenant in, rather than defaulting to "poll everyone." // // WorkerId is unique per process-lifetime: // "{MachineName}_{ProcessId}_{Guid}" — prevents dead-worker // lock collision when a crashed worker restarts on the same host. // // Recovery sweep: TGOPEXECUTIONLOCK is also not TENANTID-filtered, so // expired-lock cleanup runs per PHYSICAL DATABASE (keyed by // ConnectionName), not once globally for the whole process — otherwise // only the first-polled database's locks would ever get reset. // // Uses IServiceScopeFactory because BackgroundService is Singleton // but all DAL/BLL/Pipeline dependencies are Scoped. One scope covers the // whole poll tick (all allowlisted tenants) — safe because IQueryExecutor/ // Dapper calls in this codebase are per-call-LoginDTO-driven, not // scope-bound (same pattern OutBoxPollerService itself relies on). // ============================================================ public class GopWorkerService : BackgroundService { // Unique identifier for this worker instance — stable for the process lifetime. // Format prevents two workers on the same host from sharing an ID after a restart. private static readonly string WorkerId = $"{Environment.MachineName}_{Environment.ProcessId}_{Guid.NewGuid():N}"; private readonly IServiceScopeFactory _ScopeFactory; private readonly IConfiguration _Configuration; private readonly IOptionsMonitor _DatabaseOptions; private readonly ILogger _Logger; // Configuration knobs — move to IOptions if tuning is needed per environment. private const int PollIntervalMs = 5_000; // 5 s between queue polls private const int RecoveryIntervalMs = 60_000; // 1 min between recovery sweeps private const int LockDurationSec = 300; // 5-min lock TTL per execution private const int BatchSize = 10; // max executions per poll cycle per tenant // Recovery-sweep timing, keyed per physical database (ConnectionName) — not global, // since RESET_EXPIRED_LOCKS is itself not TENANTID-filtered (see class comment). private readonly Dictionary _nextRecoveryAtByConnection = new(); // Same MSERVERCONFIG/MSERVER query as JobEngineSL.Services.OutBoxPollerService.TenantQuery — // duplicated locally rather than shared cross-project, matching this codebase's existing // convention of each poller keeping its own copy (no shared "active tenants" helper exists // today, and introducing one is out of scope for this fix). private const string ActiveTenantsQuery = @" 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'"; public GopWorkerService( IServiceScopeFactory scopeFactory, IConfiguration configuration, IOptionsMonitor databaseOptions, ILogger logger) { _ScopeFactory = scopeFactory; _Configuration = configuration; _DatabaseOptions = databaseOptions; _Logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _Logger.LogInformation( "GopWorkerService starting. WorkerId={WorkerId}", WorkerId); while (!stoppingToken.IsCancellationRequested) { try { await PollAsync(stoppingToken); } catch (OperationCanceledException) { break; } catch (Exception ex) { _Logger.LogError(ex, "GopWorkerService: unhandled exception in poll loop (WorkerId={WorkerId})", WorkerId); } try { await Task.Delay(PollIntervalMs, stoppingToken); } catch (OperationCanceledException) { break; } } _Logger.LogInformation( "GopWorkerService stopping. WorkerId={WorkerId}", WorkerId); } private async Task PollAsync(CancellationToken ct) { var allowedConnectionNames = GetAllowedConnectionNames(); if (allowedConnectionNames.Count == 0) { _Logger.LogWarning( "GopWorkerService: GOP:WorkerAllowedConnectionNames is empty — skipping this poll " + "cycle. Add one or more MSERVERCONFIG ConnectionName values to this list to let the " + "worker process GOP executions for those tenants."); return; } await using var scope = _ScopeFactory.CreateAsyncScope(); var appConnection = scope.ServiceProvider.GetRequiredService(); IReadOnlyList tenants; try { tenants = await LoadActiveTenantsAsync(appConnection); } catch (Exception ex) { _Logger.LogError(ex, "GopWorkerService: failed to load tenant list from MSERVERCONFIG — skipping this poll cycle"); return; } var allowedTenants = tenants .Where(t => allowedConnectionNames.Contains(t.ConnectionName, StringComparer.OrdinalIgnoreCase)) .ToList(); // Sequential by design, not a limitation to fix later: GOP steps can have real external // side effects (ApiCall), and the allowlist is meant to stay short — unlike // OutBoxPollerService's unconditional all-tenant loop (a fine trade-off for its own // job of relaying already-decided outbound events), one-tenant-at-a-time is the safer // default here, not something to parallelize away. foreach (var tenant in allowedTenants) { if (ct.IsCancellationRequested) break; LoginDTO workerLogin; try { await appConnection.DBConnectionStringCached(tenant.ConnectionName); workerLogin = BuildWorkerLogin(tenant); } catch (Exception ex) { _Logger.LogWarning(ex, "GopWorkerService: skipping tenant {ClientId} ({ConnectionName}) — connection not configured", tenant.ClientId, tenant.ConnectionName); continue; } try { await PollTenantAsync(scope, workerLogin, ct); } catch (OperationCanceledException) { throw; // propagate cancellation } catch (Exception ex) { _Logger.LogError(ex, "GopWorkerService: unhandled exception polling tenant {ClientId} ({ConnectionName}) (WorkerId={WorkerId})", tenant.ClientId, tenant.ConnectionName, WorkerId); } } } private async Task PollTenantAsync(AsyncServiceScope scope, LoginDTO workerLogin, CancellationToken ct) { var queueDal = scope.ServiceProvider.GetRequiredService(); var pipeline = scope.ServiceProvider.GetRequiredService(); var lockSvc = scope.ServiceProvider.GetRequiredService(); // ── Recovery sweep (per physical database, not global) ──────────── var recoveryKey = workerLogin.ConnectionDatabaseName ?? workerLogin.DatabaseName; if (!_nextRecoveryAtByConnection.TryGetValue(recoveryKey, out var nextAt) || DateTime.UtcNow >= nextAt) { await RunRecoverySweepAsync(queueDal, workerLogin); _nextRecoveryAtByConnection[recoveryKey] = DateTime.UtcNow.AddMilliseconds(RecoveryIntervalMs); } // ── Fetch pending executions ─────────────────────────────────────── var pending = await queueDal.GetPendingExecutions(BatchSize, workerLogin); foreach (var execution in pending) { if (ct.IsCancellationRequested) break; // Each execution is processed in its own try block so one failure // does not prevent other executions in the same batch from running. try { await ProcessExecutionAsync( execution.ExecutionId, execution.ClientId, queueDal, lockSvc, pipeline, workerLogin, ct); } catch (OperationCanceledException) { throw; // propagate cancellation } catch (Exception ex) { _Logger.LogError(ex, "GopWorkerService: failed to process execution {ExecutionId} (WorkerId={WorkerId})", execution.ExecutionId, WorkerId); } } } private async Task ProcessExecutionAsync( int executionId, int clientId, IGopQueueDAL queueDal, GopExecutionLockService lockSvc, GopExecutionPipeline pipeline, LoginDTO workerLogin, CancellationToken ct) { // Acquire lock — returns false if another worker already holds it. bool acquired = await lockSvc.TryAcquireLockAsync( clientId, executionId, WorkerId, LockDurationSec, workerLogin, ct); if (!acquired) { _Logger.LogDebug( "GopWorkerService: execution {ExecutionId} lock already held — skipping", executionId); return; } _Logger.LogInformation( "GopWorkerService: acquired lock for execution {ExecutionId} (WorkerId={WorkerId})", executionId, WorkerId); try { await pipeline.ExecuteAsync(executionId, clientId, workerLogin, ct); } finally { // Always release the lock — even if the pipeline threw. await lockSvc.ReleaseLockAsync(clientId, executionId, workerLogin, ct); _Logger.LogInformation( "GopWorkerService: released lock for execution {ExecutionId} (WorkerId={WorkerId})", executionId, WorkerId); } } private async Task RunRecoverySweepAsync(IGopQueueDAL queueDal, LoginDTO workerLogin) { _Logger.LogDebug( "GopWorkerService: running recovery sweep for {ConnectionName} (WorkerId={WorkerId})", workerLogin.ConnectionDatabaseName, WorkerId); await queueDal.ResetExpiredLocks(workerLogin); } // Reads the safety allowlist — see class comment. Empty/unset means the worker does // nothing, by design; this is not the same as "unconfigured = broken," it's the safe // default until an operator deliberately opts a specific tenant's ConnectionName in. private IReadOnlyList GetAllowedConnectionNames() { var values = _Configuration.GetSection("GOP:WorkerAllowedConnectionNames").Get(); return values is { Length: > 0 } ? values : Array.Empty(); } private async Task> LoadActiveTenantsAsync(IApplicationConnection appConnection) { var systemConn = await appConnection.Gb5SystemConnectionString(); var dbType = _DatabaseOptions.CurrentValue.DataBaseType; using IDbConnection conn = dbType switch { DBTYPE.SQL => new SqlConnection(systemConn), DBTYPE.POSTGRESQL => new NpgsqlConnection(systemConn), _ => throw new NotSupportedException($"Unsupported DB type: {dbType}") }; var rows = await SqlMapper.QueryAsync(conn, ActiveTenantsQuery); return rows.ToList(); } // System LoginDTO used by the background worker for one specific, allowlisted tenant // database. ClientId is carried through for completeness/logging but doesn't affect // correctness of the queries this worker runs (see class comment, Level 2) — only // DatabaseName/ConnectionDatabaseName (which physical database to connect to) does. private static LoginDTO BuildWorkerLogin(ServerConfigDTO tenant) => new() { UserId = -1, ClientId = tenant.ClientId, DatabaseName = tenant.DatabaseName, ConnectionDatabaseName = tenant.ConnectionName, DatabaseType = tenant.DbType, UserCode = WorkerId }; } }