using System.Data; using Microsoft.Data.SqlClient; 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 GB5Shared.QueryExecutor; using Microsoft.AspNetCore.Mvc; using Microsoft.Extensions.Options; using Npgsql; using Quartz; using Serilog; using static GB5Shared.GB5Constant.Constant; namespace FrameworkSL.Controllers.PubSub { public class OutboxPublishQuartzJob : IJob { private readonly IOutBox _outBox; private readonly IApplicationConnection _appConnection; private readonly IOptionsSnapshot _databaseDTO; private readonly IQueryExecutor _queryExecutor; public OutboxPublishQuartzJob( IOutBox outBox, IApplicationConnection appConnection, IOptionsSnapshot databaseDTO, IQueryExecutor queryExecutor) { _outBox = outBox; _appConnection = appConnection; _databaseDTO = databaseDTO; _queryExecutor = queryExecutor; } public async Task Execute(IJobExecutionContext context) { try { string systemConn = await _appConnection.Gb5SystemConnectionString(); int dbType = _databaseDTO.Value.DataBaseType; string sql = @" 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'"; using IDbConnection connection = dbType switch { DBTYPE.SQL => new SqlConnection(systemConn), DBTYPE.POSTGRESQL => new NpgsqlConnection(systemConn), _ => throw new NotSupportedException($"Unsupported DB Type: {dbType}") }; var tenants = (await connection.QueryAsync(sql)).ToList(); Log.Information("Fetched {Count} tenants", tenants.Count); Log.Information("Outbox publish: {Count} tenant(s) found", tenants.Count); foreach (var tenant in tenants) { Log.Information( "Outbox publish: processing tenant ClientId={ClientId} ConnectionName={Conn} DB={DB}", tenant.ClientId, tenant.ConnectionName, tenant.DatabaseName); var login = new LoginDTO { UserId = -1, ClientId = tenant.ClientId, ConnectionDatabaseName = tenant.ConnectionName, DatabaseName = tenant.DatabaseName, DatabaseType = tenant.DbType }; // Pre-validate using ConnectionName (same lookup path as PublishPendingEventsAsync) try { await _appConnection.DBConnectionStringCached(tenant.ConnectionName) .ConfigureAwait(false); } catch { Log.Warning( "Outbox publish: skipping tenant {ClientId} ({ConnectionName}) — connection not configured", tenant.ClientId, tenant.ConnectionName); continue; } try { // Reset FAILED rows older than 5 minutes back to PENDING so they are retried try { await _queryExecutor.ExecuteAsync(login, GB5Shared.Query.OutBox.OutBoxQB.ResetFailedToPending, new { Pending = GB5Shared.GB5Constant.Constant.OUTBOXSTATUS.PENDING, Failed = GB5Shared.GB5Constant.Constant.OUTBOXSTATUS.FAILED, MinutesOld = 5 }); } catch (Exception exReset) { Log.Warning(exReset, "Outbox publish: ResetFailedToPending skipped for ClientId={ClientId} ConnectionName={Conn}", tenant.ClientId, tenant.ConnectionName); } await _outBox.PublishPendingEventsAsync( pubsubName: PUBLISHTYPE.PUBSUB, LoginDTO: login, batchSize: 100 ); Log.Information( "✅ Outbox publish done for ClientId={ClientId} ConnectionName={Conn}", tenant.ClientId, tenant.ConnectionName); } catch (Exception exTenant) { Log.Error(exTenant, "❌ Outbox publish failed for ClientId={ClientId} ConnectionName={Conn}", tenant.ClientId, tenant.ConnectionName); } } Log.Information("📤 Outbox publish job completed"); } catch (Exception ex) { Log.Error(ex, "❌ Outbox publish job failed completely"); } } } }