using Dapper; using GB5Shared.Connection; using GB5Shared.DTO.Framework.CommonConfig; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ServerConfig; using GB5Shared.PubSub.OutBox; 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 JobEngineSL.Services { // Polls TOUTBOX every 10 seconds and publishes unpublished events via Dapr. public sealed class OutBoxPollerService : BackgroundService { private readonly IServiceScopeFactory _scopeFactory; private readonly IOptionsMonitor _databaseOptions; private readonly ILogger _logger; private const string PubSubName = "pubsub"; 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'"; public OutBoxPollerService( IServiceScopeFactory scopeFactory, IOptionsMonitor databaseOptions, ILogger logger) { _scopeFactory = scopeFactory; _databaseOptions = databaseOptions; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { using var timer = new PeriodicTimer(TimeSpan.FromSeconds(10)); while (await timer.WaitForNextTickAsync(stoppingToken).ConfigureAwait(false)) { await PollAllTenantsAsync(stoppingToken).ConfigureAwait(false); } } private async Task PollAllTenantsAsync(CancellationToken ct) { using var scope = _scopeFactory.CreateScope(); var appConnection = scope.ServiceProvider.GetRequiredService(); var outBox = scope.ServiceProvider.GetRequiredService(); IReadOnlyList tenants; try { tenants = await LoadTenantsAsync(appConnection).ConfigureAwait(false); } catch (Exception ex) { _logger.LogError(ex, "OutBoxPoller: failed to load tenant list from MSERVERCONFIG"); return; } foreach (var tenant in tenants) { ct.ThrowIfCancellationRequested(); try { await appConnection.DBConnectionStringCached(tenant.ConnectionName).ConfigureAwait(false); } catch { _logger.LogWarning( "OutBoxPoller: skipping tenant {ClientId} ({ConnectionName}) — connection not configured", tenant.ClientId, tenant.ConnectionName); continue; } var tenantLogin = BuildLogin(tenant); try { await outBox.PublishPendingEventsAsync(PubSubName, tenantLogin).ConfigureAwait(false); } catch (OperationCanceledException) { throw; } catch (Exception ex) { _logger.LogError(ex, "OutBoxPoller failed for tenant {TenantId}", tenantLogin.ClientId); } } } private async Task> LoadTenantsAsync(IApplicationConnection appConnection) { var systemConn = await appConnection.Gb5SystemConnectionString().ConfigureAwait(false); 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, TenantQuery).ConfigureAwait(false); return rows.ToList(); } private static LoginDTO BuildLogin(ServerConfigDTO tenant) => new() { UserId = -1, ClientId = tenant.ClientId, ConnectionDatabaseName = tenant.ConnectionName, DatabaseName = tenant.DatabaseName, DatabaseType = tenant.DbType }; } }