using Dapper; using GB5Shared.Attachment; using GB5Shared.Connection; using GB5Shared.DTO.Framework.CommonConfig; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ServerConfig; using GB5Shared.DTO.PubSub; using Microsoft.Data.SqlClient; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using Npgsql; using GB5Shared.Telemetry; using System; using System.Collections.Generic; using System.Data; using System.Threading; using System.Threading.Tasks; using static GB5Shared.GB5Constant.Constant; namespace FrameworkBLL.Attachment { public interface IAttachmentResolveSubBLL { Task ResolveAsync(AttachmentResolveMessage message, CancellationToken ct = default); } public sealed class AttachmentResolveSubBLL : IAttachmentResolveSubBLL { private readonly IAttachmentResolveDAL _dal; private readonly IApplicationConnection _appConnection; private readonly IOptionsMonitor _systemDto; private readonly ILogger _logger; public AttachmentResolveSubBLL( IAttachmentResolveDAL dal, IApplicationConnection appConnection, IOptionsMonitor systemDto, ILogger logger) { _dal = dal; _appConnection = appConnection; _systemDto = systemDto; _logger = logger; } // ── Entry point called by the Dapr subscriber controller ────────────── public async Task ResolveAsync(AttachmentResolveMessage message, CancellationToken ct = default) { using var activity = GB5Shared.Telemetry.Dapr.DaprEventTracing.StartSubscribeActivity( "ATTACHMENT-RESOLVE", PUBLISHTYPE.PUBSUB, message.TraceParent, message.TraceState); activity?.SetTag("gb5.object.header_type_id", message.ObjectHeaderTypeId); activity?.SetTag("gb5.object.id", message.ObjectId); activity?.SetTag("gb5.tenant.id", message.TenantId); activity?.SetTag("gb5.header_row_guid", message.HeaderRowGuid.ToString()); // Load only the databases belonging to this specific client (TenantId = CLIENTID). // One client can have multiple databases (BASICTEST, TRANSACTIONCTRL, etc.). // We run the SP in all of them — only the one that holds the matching // TATTACHMENT rows will update anything; the rest update 0 rows safely. var tenants = await LoadClientDatabasesAsync(message.TenantId).ConfigureAwait(false); int resolved = 0, skipped = 0, errors = 0; foreach (var tenant in tenants) { // ── Validate the connection is reachable before attempting SP ── try { await _appConnection.DBConnectionStringCached(tenant.ConnectionName) .ConfigureAwait(false); } catch { _logger.LogWarning( "[attachment-resolve] Skipping tenant {ClientId} ({ConnectionName}) — connection not configured", tenant.ClientId, tenant.ConnectionName); skipped++; continue; } var login = new LoginDTO { ClientId = tenant.ClientId, DatabaseName = tenant.DatabaseName, DatabaseType = tenant.DbType, ConnectionDatabaseName = tenant.ConnectionName }; try { await _dal.ResolveAsync( message.ObjectHeaderTypeId, message.HeaderRowGuid, message.ObjectId, login, ct).ConfigureAwait(false); resolved++; } catch (SqlException sqlEx) when (sqlEx.Number == 2812) { // 2812 = stored procedure not found — SP not yet deployed to this DB. // Log a warning so the deployer knows which DB still needs the SP. _logger.LogWarning( "[attachment-resolve] SP_RESOLVE_ATTACHMENTS_BATCH not found | DB={DB} | Action=Deploy the SP to this database", tenant.DatabaseName); skipped++; } catch (Exception ex) { _logger.LogError(ex, "[attachment-resolve] ERROR | ClientId={ClientId} | DB={DB} | HeaderRowGuid={HeaderRowGuid} | ObjectId={ObjectId} | Error={Error}", tenant.ClientId, tenant.DatabaseName, message.HeaderRowGuid, message.ObjectId, ex.Message); errors++; } } _logger.LogInformation( "[attachment-resolve] All tenants processed | HeaderRowGuid={HeaderRowGuid} | ObjectId={ObjectId} | Resolved={Resolved} | Skipped={Skipped} | Errors={Errors}", message.HeaderRowGuid, message.ObjectId, resolved, skipped, errors); activity?.SetTag("gb5.attachment_resolve.resolved", resolved); activity?.SetTag("gb5.attachment_resolve.skipped", skipped); activity?.SetTag("gb5.attachment_resolve.errors", errors); // Re-throw to NACK Dapr only when every single tenant errored (not just skipped). if (errors > 0 && resolved == 0 && skipped == 0) { var ex = new InvalidOperationException( $"[attachment-resolve] SP failed in all {errors} tenant DB(s) — Dapr will retry."); GB5Trace.RecordError(ex, context: "attachment-resolve-all-tenants-failed"); activity?.SetStatus(System.Diagnostics.ActivityStatusCode.Error, ex.Message); throw ex; } activity?.SetStatus(System.Diagnostics.ActivityStatusCode.Ok); } // ── Load all databases that belong to a specific client (CLIENTID) ────── private async Task> LoadClientDatabasesAsync(int clientId) { var systemConn = await _appConnection.Gb5SystemConnectionString().ConfigureAwait(false); int dbType = _systemDto.CurrentValue.DataBaseType; // A single client (TenantId / CLIENTID) can have multiple database entries // in MSERVERCONFIG (e.g. BASICTEST, TRANSACTIONCTRL, MASTERTEST …). // We scope strictly to this client so we never touch other clients' data. 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.CLIENTID = @ClientId AND SERVERCONFIG1.CONNECTIONNAME <> 'ACTIVITI'"; using IDbConnection conn = dbType switch { DBTYPE.SQL => new SqlConnection(systemConn), DBTYPE.POSTGRESQL => new NpgsqlConnection(systemConn), _ => throw new NotSupportedException($"Unsupported DB type: {dbType}") }; return await conn.QueryAsync(sql, new { ClientId = clientId }) .ConfigureAwait(false); } } }