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");
}
}
}
}