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