using Dapper; using Npgsql; namespace KArtSell.Host.Features.Observability; /// /// Data queries for observability metrics. /// All queries use PIT (point-in-time) pattern: published_at <= cutoff. /// Schema-qualified, explicit columns, no SELECT *. /// public class MetricsSql { private readonly NpgsqlDataSource _dataSource; public MetricsSql(NpgsqlDataSource dataSource) { _dataSource = dataSource; } public async Task<(int Total, int OnTime, TimeSpan AvgTime)?> GetBatchSlaAsync(CancellationToken cancellationToken = default) { const string sql = """ SELECT COUNT(*) as total, COUNT(CASE WHEN status = 'success' THEN 1 END) as on_time, AVG(duration_seconds) as avg_seconds FROM observability.batch_sla_metrics WHERE published_at <= @now AND completed_at >= @sevenDaysAgo """; await using var connection = await _dataSource.OpenConnectionAsync(cancellationToken); var result = await connection.QueryFirstOrDefaultAsync<(int, int, double)?>( sql, new { now = DateTime.UtcNow, sevenDaysAgo = DateTime.UtcNow.AddDays(-7) }, commandTimeout: 5); if (result == null || result.Value.Item1 == 0) return null; var (total, onTime, avgSec) = result.Value; return (total, onTime, TimeSpan.FromSeconds(avgSec)); } public async Task<(int Quarantined, int Total, List Errors)?> GetDataQualityQuarantineAsync(CancellationToken cancellationToken = default) { const string sql = """ SELECT COUNT(*) as total, COUNT(CASE WHEN resolution_status IS NULL THEN 1 END) as quarantined, STRING_AGG(DISTINCT reason, ', ') as errors FROM observability.data_quality_quarantine WHERE published_at <= @now AND quarantined_at >= @sevenDaysAgo """; await using var connection = await _dataSource.OpenConnectionAsync(cancellationToken); var result = await connection.QueryFirstOrDefaultAsync<(int Total, int Quarantined, string? Errors)?>( sql, new { now = DateTime.UtcNow, sevenDaysAgo = DateTime.UtcNow.AddDays(-7) }, commandTimeout: 5); if (result == null || result.Value.Item1 == 0) return null; var (total, quarantined, errors) = result.Value; var errorList = string.IsNullOrEmpty(errors) ? new List() : errors.Split(',').Select(e => e.Trim()).Take(5).ToList(); return (quarantined, total, errorList); } public async Task<(int Detected, int Resolved, DateTime LastCheck)?> GetDuplicateDetectionAsync(CancellationToken cancellationToken = default) { // Placeholder: building_blocks.outbox_message table exists (0000_building_blocks.sql). // Duplicate detection logging (via operation_audit_trail or dedicated table) not yet implemented. // Returns null until OutboxPollerJob hooks duplicate tracking (see DEBT-014). await Task.CompletedTask; // Async compliance return null; } public async Task<(int Detected, int Resolved, List Pending)?> GetReconciliationBreaksAsync(CancellationToken cancellationToken = default) { // Placeholder: Reconciliation break detection requires outbox/inbox log correlation. // Requires audit trail showing Evidence version mismatches. Not yet implemented. // Returns null until operation_audit_trail is populated by job consumers (see DEBT-014). await Task.CompletedTask; // Async compliance return null; } public async Task<(decimal Baseline, decimal Current)?> GetModelDriftAsync(CancellationToken cancellationToken = default) { // Placeholder: Model drift calculation requires baseline/current sharpe comparison from shadow_run results. // Returns null until Gate 3 rehearsal populates model_operations.shadow_run with real metrics. // Once shadow_run results exist, baseline/current sharpe can be calculated and compared (see DEBT-009). await Task.CompletedTask; // Async compliance return null; } }