using Dapper;
using FrameworkDAL.Query.EventSub;
using GB5Shared.Connection;
using GB5Shared.DTO.Framework.CommonConfig;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using Microsoft.Data.SqlClient;
using Npgsql;
using System;
using System.Collections.Generic;
using System.Data;
using System.Threading;
using System.Threading.Tasks;
using static GB5Shared.GB5Constant.Constant;
namespace FrameworkDAL.CustomCode.EventSub
{
public interface IDaprSubscriptionDAL
{
///
/// Returns all distinct EventTypeIds found in MACTION across every active tenant.
/// Used to build dynamic Dapr topic subscriptions at runtime — no hardcoded EventTypeIds.
///
Task> GetActiveEventTypeIdsAsync(CancellationToken ct = default);
}
public sealed class DaprSubscriptionDAL : IDaprSubscriptionDAL
{
private readonly IApplicationConnection _appConnection;
private readonly IOptionsMonitor _systemDto;
private readonly ILogger _logger;
public DaprSubscriptionDAL(
IApplicationConnection appConnection,
IOptionsMonitor systemDto,
ILogger logger)
{
_appConnection = appConnection;
_systemDto = systemDto;
_logger = logger;
}
public async Task> GetActiveEventTypeIdsAsync(CancellationToken ct = default)
{
var systemConnStr = await _appConnection.Gb5SystemConnectionString().ConfigureAwait(false);
int dbType = _systemDto.CurrentValue.DataBaseType;
_logger.LogInformation("DaprSubscriptionDAL: querying tenant connections from system DB");
// Step 1 — get all tenant connection names from the system DB
IEnumerable connectionNames;
using (IDbConnection systemConn = CreateConnection(dbType, systemConnStr))
{
connectionNames = await systemConn.QueryAsync(
new CommandDefinition(
DaprSubscriptionQB.GET_ALL_TENANT_CONNECTIONS,
cancellationToken: ct))
.ConfigureAwait(false);
}
// Step 2 — query MACTION in each tenant DB, collect distinct EventTypeIds
var eventTypeIds = new HashSet();
foreach (var connectionName in connectionNames)
{
try
{
_logger.LogInformation("DaprSubscriptionDAL: querying MACTION for connection={ConnectionName}", connectionName);
var tenantConnStr = await _appConnection
.DBConnectionStringCached(connectionName)
.ConfigureAwait(false);
using IDbConnection tenantConn = CreateConnection(dbType, tenantConnStr);
var ids = await tenantConn.QueryAsync(
new CommandDefinition(
DaprSubscriptionQB.GET_DISTINCT_EVENT_TYPE_IDS_FROM_ACTION,
cancellationToken: ct))
.ConfigureAwait(false);
int count = 0;
foreach (var id in ids)
{
eventTypeIds.Add(id);
count++;
}
_logger.LogInformation("DaprSubscriptionDAL: connection={ConnectionName} → {Count} EventTypeId(s)", connectionName, count);
}
catch (Exception ex)
{
_logger.LogError(ex, "DaprSubscriptionDAL: FAILED for connection={ConnectionName} — skipping", connectionName);
}
}
_logger.LogInformation("DaprSubscriptionDAL: total distinct EventTypeIds={Total}", eventTypeIds.Count);
return eventTypeIds;
}
private static IDbConnection CreateConnection(int dbType, string connStr) =>
dbType switch
{
DBTYPE.SQL => new SqlConnection(connStr),
DBTYPE.POSTGRESQL => new NpgsqlConnection(connStr),
_ => throw new NotSupportedException($"Unsupported DB type: {dbType}")
};
}
}