using System.Data;
using System.Text.Json;
using Dapper;
using FrameworkDAL.Query.DataAlert;
using GB5Shared.Connection;
using GB5Shared.DTO.Framework.CommonConfig;
using GB5Shared.DTO.Framework.Login;
using GB5Shared.DTO.Framework.ServerConfig;
using GB5Shared.QueryExecutor;
using Microsoft.Data.SqlClient;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using Npgsql;
using static GB5Shared.GB5Constant.Constant;
namespace FrameworkBLL.DataAlert;
///
/// Background job that evaluates MDATAALERT rules with ALERTSOURCE=2 (observability alerts)
/// against TOBSERVABILITY metrics every 5 minutes.
/// When a threshold is breached, recipients are resolved from MSUBSCRIPTION and a
/// TNOTIFICATION + TNOTIFICATIONTO row is written per-tenant.
///
public sealed class ObservabilityAlertEvaluatorJob : BackgroundService
{
private readonly IServiceProvider _services;
private readonly IApplicationConnection _appConnection;
private readonly IOptionsMonitor _systemDto;
private readonly ILogger _logger;
private static readonly TimeSpan EvalInterval = TimeSpan.FromMinutes(5);
public ObservabilityAlertEvaluatorJob(
IServiceProvider services,
IApplicationConnection appConnection,
IOptionsMonitor systemDto,
ILogger logger)
{
_services = services;
_appConnection = appConnection;
_systemDto = systemDto;
_logger = logger;
}
protected override async Task ExecuteAsync(CancellationToken ct)
{
_logger.LogInformation("ObservabilityAlertEvaluatorJob started. Interval=5m.");
while (!ct.IsCancellationRequested)
{
try
{
await EvaluateAllTenantsAsync(ct).ConfigureAwait(false);
}
catch (OperationCanceledException) { break; }
catch (Exception ex)
{
_logger.LogError(ex, "ObservabilityAlertEvaluatorJob cycle failed at {Time:HH:mm:ss}", DateTime.UtcNow);
}
try { await Task.Delay(EvalInterval, ct).ConfigureAwait(false); }
catch (TaskCanceledException) { break; }
}
_logger.LogInformation("ObservabilityAlertEvaluatorJob stopped.");
}
private async Task EvaluateAllTenantsAsync(CancellationToken ct)
{
var tenants = await GetTenantsAsync(ct).ConfigureAwait(false);
foreach (var tenant in tenants)
{
try
{
await EvaluateTenantAsync(tenant, ct).ConfigureAwait(false);
}
catch (OperationCanceledException) { throw; }
catch (Exception ex)
{
_logger.LogError(ex,
"ObservabilityAlertEvaluatorJob failed for tenant {DatabaseName}", tenant.DatabaseName);
}
}
}
private async Task EvaluateTenantAsync(ServerConfigDTO tenant, CancellationToken ct)
{
using var scope = _services.CreateScope();
var qe = scope.ServiceProvider.GetRequiredService();
var login = BuildSystemLogin(tenant);
// 1. Load active observability alert rules for this tenant
var alerts = (await qe.QueryAsync(
login, ObservabilityAlertQB.GET_ACTIVE_OBS_ALERTS, null, cancellationToken: ct)
.ConfigureAwait(false))?.ToList() ?? [];
if (alerts.Count == 0) return;
// 2. Load TOBSERVABILITY metrics from system DB (shared across tenants)
using var sysConn = await OpenSystemConnectionAsync(ct).ConfigureAwait(false);
foreach (var alert in alerts)
{
double.TryParse(alert.Filter1, out var errorRateThreshold); // FILTER1 = error rate %
int.TryParse(alert.Filter2, out var slowCountThreshold); // FILTER2 = slow tx count
int.TryParse(alert.Filter3, out var hoursBack); // FILTER3 = lookback hours
if (hoursBack <= 0) hoursBack = 1;
var metrics = await sysConn.QuerySingleOrDefaultAsync(
new CommandDefinition(
ObservabilityAlertQB.EVAL_OBSERVABILITY_METRICS,
new { TenantId = tenant.ClientId, HoursBack = hoursBack },
cancellationToken: ct)).ConfigureAwait(false);
if (metrics is null) continue;
bool errorRateBreached = errorRateThreshold > 0 && metrics.ErrorRatePct >= errorRateThreshold;
bool slowCountBreached = slowCountThreshold > 0 && metrics.SlowCount >= slowCountThreshold;
if (!errorRateBreached && !slowCountBreached) continue;
_logger.LogWarning(
"ObservabilityAlert {AlertName} breached for tenant {TenantId}: ErrorRate={Rate:F1}% SlowCount={Slow}",
alert.DataAlertName, tenant.ClientId, metrics.ErrorRatePct, metrics.SlowCount);
await FireAlertAsync(alert, metrics, login, qe, ct).ConfigureAwait(false);
}
}
private async Task FireAlertAsync(
AlertRuleRow alert, MetricsRow metrics,
LoginDTO login, IQueryExecutor qe, CancellationToken ct)
{
// Resolve recipient UserIds from MSUBSCRIPTION (JSON array stored in USERIDS column)
var userIdJsonRows = (await qe.QueryAsync(
login,
"SELECT S.USERIDS FROM MSUBSCRIPTION S WHERE S.OBJECTTYPEID = 1 AND S.OBJECTID = @DataAlertId AND S.STATUS = 1",
new { DataAlertId = alert.DataAlertId },
cancellationToken: ct).ConfigureAwait(false))?.ToList() ?? [];
var userIds = new List();
foreach (var json in userIdJsonRows)
{
try
{
var ids = JsonSerializer.Deserialize>(json ?? "[]");
if (ids is not null) userIds.AddRange(ids);
}
catch { /* malformed JSON in USERIDS — skip silently */ }
}
if (userIds.Count == 0) return;
var message = $"{alert.DataAlertName}: ErrorRate={metrics.ErrorRatePct:F1}%, SlowCount={metrics.SlowCount}. {alert.Message}";
var insertSql = login.DatabaseType == DBTYPE.SQL
? ObservabilityAlertQB.INSERT_NOTIFICATION
: ObservabilityAlertQB.INSERT_NOTIFICATION_POSTGRESQL;
var notificationId = await qe.ExecuteIdentityAsync(
login, insertSql,
new { Message = message, DataAlertId = alert.DataAlertId }).ConfigureAwait(false);
foreach (var userId in userIds.Distinct())
{
await qe.ExecuteAsync(
login,
ObservabilityAlertQB.INSERT_NOTIFICATION_TO,
new { NotificationId = notificationId, ToUserId = userId },
cancellationToken: ct).ConfigureAwait(false);
}
await qe.ExecuteAsync(
login,
ObservabilityAlertQB.UPSERT_DATAALERT_SUMMARY,
new { DataAlertId = alert.DataAlertId, NumberOfRecords = userIds.Count },
cancellationToken: ct).ConfigureAwait(false);
}
// ── Helpers ───────────────────────────────────────────────────────────────
private async Task> GetTenantsAsync(CancellationToken ct)
{
using var conn = await OpenSystemConnectionAsync(ct).ConfigureAwait(false);
var rows = await conn.QueryAsync(
new CommandDefinition(
"SELECT CLIENTID AS ClientId, DATABASENAME AS DatabaseName, DATABASETYPE AS DatabaseType FROM MSERVERCONFIG WHERE STATUS = 1",
cancellationToken: ct)).ConfigureAwait(false);
return rows?.ToList() ?? [];
}
private async Task OpenSystemConnectionAsync(CancellationToken ct)
{
var connStr = await _appConnection.Gb5SystemConnectionString().ConfigureAwait(false);
var dbType = _systemDto.CurrentValue.DataBaseType;
IDbConnection conn = dbType switch
{
DBTYPE.SQL => new SqlConnection(connStr),
DBTYPE.POSTGRESQL => new NpgsqlConnection(connStr),
_ => throw new NotSupportedException($"Unsupported system DB type: {dbType}")
};
conn.Open();
return conn;
}
private static LoginDTO BuildSystemLogin(ServerConfigDTO tenant) => new()
{
ClientId = tenant.ClientId,
DatabaseName = tenant.DatabaseName,
DatabaseType = tenant.DbType,
UserId = -1
};
// ── Private DTOs (projection only — not exposed outside job) ─────────────
private sealed record AlertRuleRow(
int DataAlertId,
string DataAlertName,
int AlertLevel,
string Message,
string? Filter1,
string? Filter2,
string? Filter3);
private sealed record MetricsRow(
int TotalRequests,
int ErrorCount,
double ErrorRatePct,
int SlowCount);
}