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