using Dapper; using FrameworkDAL.CustomCode.ChangeTracking; using FrameworkDAL.CustomCode.EventLog; using FrameworkDAL.DTO.EventLog; using GB5Shared.ChangeTracking; 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 System; using System.Collections.Generic; using System.Data; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using static GB5Shared.GB5Constant.Constant; namespace FrameworkBLL.EventSub { public interface IEventLogSubBLL { /// /// Receives an OutboxEventMessage from Dapr and writes an audit entry to TEVENTLOG. /// Called by EventLogSubscribeController for every EVENTTYPEID:{id} topic. /// Task RecordEventAsync(OutboxEventMessage message, CancellationToken ct = default); } public sealed class EventLogSubBLL : IEventLogSubBLL { private readonly IEventLogDAL _dal; private readonly IChangeTrackingDAL _changeTrackingDal; private readonly IApplicationConnection _appConnection; private readonly IOptionsMonitor _systemDto; private readonly ILogger _logger; public EventLogSubBLL( IEventLogDAL dal, IChangeTrackingDAL changeTrackingDal, IApplicationConnection appConnection, IOptionsMonitor systemDto, ILogger logger) { _dal = dal; _changeTrackingDal = changeTrackingDal; _appConnection = appConnection; _systemDto = systemDto; _logger = logger; } public async Task RecordEventAsync(OutboxEventMessage message, CancellationToken ct = default) { var login = await BuildLoginAsync(message.TenantId, message.UserId).ConfigureAwait(false); if (login is null) { _logger.LogWarning( "EventLogSub: tenant {TenantId} not found — skipping EventTypeId={EventTypeId}", message.TenantId, message.EventTypeId); return; } var dto = new EventLogDTO { EventTypeId = message.EventTypeId, UserId = login.UserId, EventLogEventText = $"EventTypeId={message.EventTypeId} ObjectTypeId={message.ObjectTypeId} ObjectId={message.ObjectId}", EventLogData = message.Payload, EventLogTimeStamp = DateTime.UtcNow, Login = login }; await _dal.SaveEventLog(dto, login).ConfigureAwait(false); _logger.LogInformation( "EventLogSub: recorded | EventTypeId={EventTypeId} ObjectId={ObjectId} TenantId={TenantId}", message.EventTypeId, message.ObjectId, message.TenantId); // ── Persist field-level changes (UPDATE path only) ───────────────── // Changes is null for INSERT operations — only write when present. if (!string.IsNullOrWhiteSpace(message.Changes)) { try { var changes = JsonSerializer.Deserialize>(message.Changes); if (changes is { Count: > 0 }) { // EventLogId is set by SaveEventLog (INSERT returns identity via SAVE_LOGIN_EVENTLOG_SQL). // For async subscribers the EventLogId is 0 because SaveEventLog is fire-and-forget. // We store EventLogId from dto — DAL must populate it after INSERT. // If EventLogId is 0 (not yet supported), skip silently. if (dto.EventLogId > 0) { await _changeTrackingDal .SaveChangesAsync(dto.EventLogId, changes, login, ct) .ConfigureAwait(false); _logger.LogInformation( "EventLogSub: saved {Count} change record(s) for EventLogId={EventLogId}", changes.Count, dto.EventLogId); } else { _logger.LogDebug( "EventLogSub: Changes present but EventLogId=0 — change records skipped for correlation {Key}", message.CorrelationKey); } } } catch (JsonException jex) { // Malformed Changes JSON must not NACK the event — log and continue _logger.LogWarning(jex, "EventLogSub: failed to deserialize Changes for correlation {Key}", message.CorrelationKey); } } } // ── Helpers ─────────────────────────────────────────────────────────── private async Task BuildLoginAsync(int tenantId, int userId = -1) { var systemConn = await _appConnection.Gb5SystemConnectionString().ConfigureAwait(false); int dbType = _systemDto.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.CLIENTID = @TenantId"; using IDbConnection conn = dbType switch { DBTYPE.SQL => new SqlConnection(systemConn), DBTYPE.POSTGRESQL => new NpgsqlConnection(systemConn), _ => throw new NotSupportedException($"Unsupported DB type: {dbType}") }; var tenant = await conn.QueryFirstOrDefaultAsync( sql, new { TenantId = tenantId }).ConfigureAwait(false); if (tenant is null) return null; return new LoginDTO { UserId = userId, ClientId = tenant.ClientId, DatabaseName = tenant.DatabaseName, DatabaseType = tenant.DbType, ConnectionDatabaseName = tenant.ConnectionName }; } } }