using FrameworkBLL.SchedulerTaskGeneratorPublisher; using GB5Shared.ActionProcessor; using FrameworkDAL.CustomCode.ActionProcessor; using GB5Shared.Connection; using GB5Shared.DTO.Framework.CommonConfig; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ServerConfig; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using Npgsql; using Quartz; using System; using System.Data; using Microsoft.Data.SqlClient; using System.Linq; using System.Threading.Tasks; using static GB5Shared.GB5Constant.Constant; namespace FrameworkSL.Controllers.ActionProcessor { /// /// Quartz job: polls TACTIONOUTBOX for SENDSTATUS=0 rows and publishes them /// to RabbitMQ via IRabbitMqPublisher. Runs on the same Quartz scheduler as /// OutboxPublishQuartzJob — register in Program.cs with a trigger interval. /// [DisallowConcurrentExecution] public class ActionOutboxDispatcherQuartzJob : IJob { private const int BatchSize = 50; private readonly IRabbitMqPublisher _publisher; private readonly IActionOutboxDAL _outboxDal; private readonly IApplicationConnection _appConnection; private readonly IOptionsMonitor _databaseDTO; private readonly ILogger _logger; public ActionOutboxDispatcherQuartzJob( IRabbitMqPublisher publisher, IActionOutboxDAL outboxDal, IApplicationConnection appConnection, IOptionsMonitor databaseDTO, ILogger logger) { _publisher = publisher; _outboxDal = outboxDal; _appConnection = appConnection; _databaseDTO = databaseDTO; _logger = logger; } public async Task Execute(IJobExecutionContext context) { var ct = context.CancellationToken; try { _logger.LogInformation("ActionOutboxDispatcherQuartzJob started at {StartTime}", DateTimeOffset.UtcNow); var systemConn = await _appConnection.Gb5SystemConnectionString(); int dbType = _databaseDTO.CurrentValue.DataBaseType; const 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 Dapper.SqlMapper.QueryAsync(connection, sql)).ToList(); foreach (var tenant in tenants) { // ✅ DatabaseName must hold the CONNECTION NAME, not the literal DB name — // ApplicationConnection resolves connection strings by treating this field // as a lookup key against MSERVERCONFIG.CONNECTIONNAME. The line below it // (DBConnectionStringCached(tenant.ConnectionName)) already used the right // value, which is why the pre-check passed while the real query failed with // "No server configuration found for connection name: ". var login = new LoginDTO { UserId = -1, ClientId = tenant.ClientId, ConnectionDatabaseName = tenant.ConnectionName, DatabaseName = tenant.ConnectionName, DatabaseType = tenant.DbType }; // Pre-validate connection — skip tenants with no registered connection string try { await _appConnection.DBConnectionStringCached(tenant.ConnectionName).ConfigureAwait(false); } catch { _logger.LogWarning( "Action outbox dispatcher: skipping tenant {ClientId} ({ConnectionName}) — connection not configured", tenant.ClientId, tenant.ConnectionName); continue; } try { var batch = await _outboxDal.GetPendingBatchAsync(BatchSize, login, ct) .ConfigureAwait(false); _logger.LogInformation( "Action outbox dispatcher: tenant {ClientId} — {Count} pending row(s)", tenant.ClientId, batch.Count); foreach (var item in batch) { try { await _publisher.PublishAsync(item.DestinationTopic, item.Payload) .ConfigureAwait(false); await _outboxDal.MarkSentAsync(item.OutboxId, login, ct) .ConfigureAwait(false); _logger.LogInformation( "Action outbox dispatched | OutboxId={Id} Topic={Topic}", item.OutboxId, item.DestinationTopic); } catch (Exception exItem) { _logger.LogError(exItem, "Action outbox dispatch failed | OutboxId={Id}", item.OutboxId); await _outboxDal.MarkFailedAsync(item.OutboxId, login, ct) .ConfigureAwait(false); } } } catch (Exception exTenant) when ( exTenant.Message.Contains("No server configuration found") || exTenant.InnerException?.Message.Contains("No server configuration found") == true) { _logger.LogWarning( "Action outbox dispatcher: skipping tenant {ClientId} ({ConnectionName}) — connection not configured", tenant.ClientId, tenant.ConnectionName); } catch (Exception exTenant) { _logger.LogError(exTenant, "Action outbox dispatcher failed for ClientId={ClientId}", tenant.ClientId); } } } catch (Exception ex) { _logger.LogError(ex, "Action outbox dispatcher job failed"); } } } }