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