diff --git a/src/KArtSell.Host/KArtSell.Host.csproj b/src/KArtSell.Host/KArtSell.Host.csproj index a861f0db..dd2d9071 100644 --- a/src/KArtSell.Host/KArtSell.Host.csproj +++ b/src/KArtSell.Host/KArtSell.Host.csproj @@ -1,4 +1,5 @@ - + bab7e095-067e-4797-b2ab-df4c1f8b447d + diff --git a/src/KArtSell.Host/Program.cs b/src/KArtSell.Host/Program.cs index 3b50bf94..991559fb 100644 --- a/src/KArtSell.Host/Program.cs +++ b/src/KArtSell.Host/Program.cs @@ -199,7 +199,7 @@ static string? ResolveSecret(string? configValue, string environmentVariable) // 2. Check if config has a placeholder (e.g., "${VAR_NAME}") if (!string.IsNullOrEmpty(configValue)) { - if (configValue.StartsWith("${") && configValue.EndsWith("}")) + if (configValue.StartsWith("${", StringComparison.Ordinal) && configValue.EndsWith('}')) { // This is a placeholder, try to resolve from environment return Environment.GetEnvironmentVariable(environmentVariable); diff --git a/src/KArtSell.Modules.ModelOperations/Features/GetObservabilityMetrics/Endpoint.cs b/src/KArtSell.Modules.ModelOperations/Features/GetObservabilityMetrics/Endpoint.cs index 5e60e783..d891c058 100644 --- a/src/KArtSell.Modules.ModelOperations/Features/GetObservabilityMetrics/Endpoint.cs +++ b/src/KArtSell.Modules.ModelOperations/Features/GetObservabilityMetrics/Endpoint.cs @@ -15,18 +15,14 @@ public sealed class Endpoint(IObservabilityService observability, IClock clock) public override async Task HandleAsync(CancellationToken ct) { - var batchSla = await observability.GetBatchSlaMetricsAsync(ct); - var dataQuality = await observability.GetDataQualityMetricsAsync(ct); - var duplicates = await observability.GetDuplicateDetectionMetricsAsync(ct); - var reconciliation = await observability.GetReconciliationMetricsAsync(ct); - var modelDrift = await observability.GetModelDriftMetricsAsync(ct); + var metrics = await observability.GetMetricsAsync(ct); var response = new Response( - batchSla, - dataQuality, - duplicates, - reconciliation, - modelDrift, + metrics.BatchSla ?? new(), + metrics.DataQuality ?? new(), + metrics.DuplicateDetection ?? new(), + metrics.Reconciliation ?? new(), + metrics.ModelDrift ?? new(), clock.UtcNow.DateTime); await Send.OkAsync(response, ct); diff --git a/src/KArtSell.Modules.ModelOperations/ModelOperationsModule.cs b/src/KArtSell.Modules.ModelOperations/ModelOperationsModule.cs index a3e686a4..d4593b10 100644 --- a/src/KArtSell.Modules.ModelOperations/ModelOperationsModule.cs +++ b/src/KArtSell.Modules.ModelOperations/ModelOperationsModule.cs @@ -1,6 +1,9 @@ using KArtSell.Modules.ModelOperations.Application; using KArtSell.Modules.ModelOperations.Domain; using KArtSell.Modules.ModelOperations.Infrastructure; +using KArtSell.Modules.ModelOperations.Observability; +using KArtSell.Modules.ModelOperations.ShadowRun; +using KArtSell.Modules.ModelOperations.ShadowRun.Services; using Microsoft.Extensions.DependencyInjection; namespace KArtSell.Modules.ModelOperations; @@ -13,6 +16,9 @@ public static class ModelOperationsModule services.AddScoped(); services.AddScoped(); services.AddScoped(); + services.AddSingleton(); + services.AddScoped(); + services.AddScoped(); return services; } } diff --git a/src/KArtSell.Modules.ModelOperations/Observability/IObservabilityService.cs b/src/KArtSell.Modules.ModelOperations/Observability/IObservabilityService.cs index 979d01df..5d496b1c 100644 --- a/src/KArtSell.Modules.ModelOperations/Observability/IObservabilityService.cs +++ b/src/KArtSell.Modules.ModelOperations/Observability/IObservabilityService.cs @@ -1,81 +1,48 @@ namespace KArtSell.Modules.ModelOperations.Observability; -/// -/// Observability service for production readiness metrics -/// Covers: Batch SLA, Data Quality, Duplicates, Reconciliation, Model Drift -/// public interface IObservabilityService { - /// - /// Batch SLA metrics: Job completion times and queue depths - /// - Task GetBatchSlaMetricsAsync(CancellationToken ct); - - /// - /// Data Quality metrics: Jobs in quarantine (marked as `dq`) - /// - Task GetDataQualityMetricsAsync(CancellationToken ct); - - /// - /// Duplicate detection: Outbox message dedup violations - /// - Task GetDuplicateDetectionMetricsAsync(CancellationToken ct); - - /// - /// Reconciliation: Evidence vs current state mismatches - /// - Task GetReconciliationMetricsAsync(CancellationToken ct); - - /// - /// Model Drift: Out-of-sample performance vs baseline - /// - Task GetModelDriftMetricsAsync(CancellationToken ct); + Task GetMetricsAsync(CancellationToken cancellationToken); } -/// -/// Batch SLA: Job completion times, queue depths, retry rates -/// -public sealed record BatchSlaMetrics( - int QueueDepth, - double AverageCompletionTimeMs, - int TotalJobsCompleted, - int RetryCount, - DateTime MeasuredAt); +public class ObservabilityMetricsDto +{ + public BatchSlaMetrics? BatchSla { get; set; } + public DataQualityMetrics? DataQuality { get; set; } + public DuplicateDetectionMetrics? DuplicateDetection { get; set; } + public ReconciliationMetrics? Reconciliation { get; set; } + public ModelDriftMetrics? ModelDrift { get; set; } +} -/// -/// Data Quality: Jobs in quarantine, reasons, age -/// -public sealed record DataQualityMetrics( - int QuarantinedJobCount, - string[] TopQuarantineReasons, - double AverageQuarantineAgeHours, - DateTime MeasuredAt); +public class BatchSlaMetrics +{ + public int QueueDepth { get; set; } + public double AverageCompletionTimeSeconds { get; set; } + public double RetryRate { get; set; } +} -/// -/// Duplicate Detection: Constraint violations, affected messages -/// -public sealed record DuplicateDetectionMetrics( - int DuplicateViolationCount, - int AffectedMessageCount, - DateTime LastViolationAt, - DateTime MeasuredAt); +public class DataQualityMetrics +{ + public int QuarantineCount { get; set; } + public int AgeMinutes { get; set; } + public List TopFailureReasons { get; set; } = new(); +} -/// -/// Reconciliation: Evidence vs state mismatches, audit trail completeness -/// -public sealed record ReconciliationMetrics( - int MismatchCount, - double AuditTrailCompleteness, - int OutboxMessageCount, - int InboxProcessedCount, - DateTime MeasuredAt); +public class DuplicateDetectionMetrics +{ + public int ConstraintViolationCount { get; set; } + public DateTime LastDetected { get; set; } +} -/// -/// Model Drift: OOS performance, baseline comparison, risk flags -/// -public sealed record ModelDriftMetrics( - int ModelsUnderMonitoring, - double AverageOosPerformance, - int PerformanceDegradedCount, - double BaselineSharpeRatio, - DateTime MeasuredAt); +public class ReconciliationMetrics +{ + public double CompletenessPercentage { get; set; } + public int AuditRecordsCount { get; set; } +} + +public class ModelDriftMetrics +{ + public double OosPerformanceValue { get; set; } + public double BaselineComparison { get; set; } + public bool DegradationFlag { get; set; } +} diff --git a/src/KArtSell.Modules.ModelOperations/Observability/ObservabilityService.cs b/src/KArtSell.Modules.ModelOperations/Observability/ObservabilityService.cs deleted file mode 100644 index 5f70e29f..00000000 --- a/src/KArtSell.Modules.ModelOperations/Observability/ObservabilityService.cs +++ /dev/null @@ -1,172 +0,0 @@ -using Dapper; -using KArtSell.BuildingBlocks.Data; -using KArtSell.BuildingBlocks.Time; - -namespace KArtSell.Modules.ModelOperations.Observability; - -/// -/// Observability service implementation for production readiness metrics -/// Following AGENTS.md v16.0: Evidence-based monitoring, constraint validation -/// -public sealed class ObservabilityService(IDbConnectionFactory connectionFactory, IClock clock) : IObservabilityService -{ - public async Task GetBatchSlaMetricsAsync(CancellationToken ct) - { - await using var connection = await connectionFactory.OpenAsync(ct); - - // Query Hangfire jobs for SLA metrics - // Counts pending jobs and calculates average completion time for completed jobs - var metrics = await connection.QuerySingleAsync<(int QueueDepth, double AvgTime, int Completed, int Retries)>(""" - SELECT - COUNT(CASE WHEN status = 'Enqueued' THEN 1 END) as QueueDepth, - COALESCE(AVG(EXTRACT(EPOCH FROM (completed_at - created_at)) * 1000), 0) as AvgTime, - COUNT(CASE WHEN status = 'Succeeded' THEN 1 END) as Completed, - COUNT(CASE WHEN attempts > 1 THEN 1 END) as Retries - FROM hangfire.job - WHERE created_at >= NOW() - INTERVAL '24 hours' - """); - - return new BatchSlaMetrics( - metrics.QueueDepth, - metrics.AvgTime, - metrics.Completed, - metrics.Retries, - clock.UtcNow.DateTime); - } - - public async Task GetDataQualityMetricsAsync(CancellationToken ct) - { - await using var connection = await connectionFactory.OpenAsync(ct); - - // Query for jobs in quarantine (status = Failed with specific error patterns) - var quarantined = await connection.QueryAsync<(string Error, int Count)>(""" - SELECT - COALESCE(SUBSTRING(exception_message FROM 1 FOR 100), 'Unknown') as Error, - COUNT(*) as Count - FROM hangfire.job - WHERE status = 'Failed' - AND created_at >= NOW() - INTERVAL '24 hours' - AND exception_message LIKE '%data quality%' OR exception_message LIKE '%dq%' OR exception_message LIKE '%validation%' - GROUP BY Error - ORDER BY Count DESC - LIMIT 5 - """); - - var quarantineList = quarantined.ToList(); - var ageQuery = await connection.QuerySingleAsync<(int Count, double AvgAgeHours)>(""" - SELECT - COUNT(*) as Count, - COALESCE(AVG(EXTRACT(EPOCH FROM (NOW() - created_at)) / 3600), 0) as AvgAgeHours - FROM hangfire.job - WHERE status = 'Failed' - AND created_at >= NOW() - INTERVAL '7 days' - """); - - return new DataQualityMetrics( - ageQuery.Count, - quarantineList.Select(q => q.Error).ToArray(), - ageQuery.AvgAgeHours, - clock.UtcNow.DateTime); - } - - public async Task GetDuplicateDetectionMetricsAsync(CancellationToken ct) - { - await using var connection = await connectionFactory.OpenAsync(ct); - - // Query for duplicate inbox records (same message_id, consumer_id pair) - var duplicates = await connection.QueryAsync<(Guid MessageId, string Consumer, int Count)>(""" - SELECT - outbox_id as MessageId, - consumer_id as Consumer, - COUNT(*) as Count - FROM outbox.inbox - WHERE created_at >= NOW() - INTERVAL '24 hours' - GROUP BY outbox_id, consumer_id - HAVING COUNT(*) > 1 - """); - - var duplicateList = duplicates.ToList(); - var totalAffected = duplicateList.Sum(d => d.Count); - - // Get last constraint violation timestamp (if any) - var lastViolation = await connection.QuerySingleOrDefaultAsync(""" - SELECT MAX(created_at) - FROM outbox.inbox - WHERE created_at >= NOW() - INTERVAL '7 days' - GROUP BY outbox_id, consumer_id - HAVING COUNT(*) > 1 - LIMIT 1 - """); - - return new DuplicateDetectionMetrics( - duplicateList.Count, - totalAffected, - lastViolation ?? clock.UtcNow.DateTime, - clock.UtcNow.DateTime); - } - - public async Task GetReconciliationMetricsAsync(CancellationToken ct) - { - await using var connection = await connectionFactory.OpenAsync(ct); - - // Query for reconciliation: Compare outbox vs inbox message counts - var reconciliation = await connection.QuerySingleAsync<(int OutboxCount, int InboxProcessed, int Mismatches)>(""" - SELECT - (SELECT COUNT(*) FROM building_blocks.outbox_message WHERE published_at IS NOT NULL) as OutboxCount, - (SELECT COUNT(*) FROM outbox.inbox WHERE status = 'Processed') as InboxProcessed, - (SELECT COUNT(*) - FROM building_blocks.outbox_message o - WHERE o.published_at IS NOT NULL - AND NOT EXISTS ( - SELECT 1 FROM outbox.inbox i WHERE i.outbox_id = o.message_id AND i.status = 'Processed' - )) as Mismatches - """); - - // Calculate audit trail completeness (% of outbox messages with corresponding inbox entries) - var completeness = reconciliation.OutboxCount > 0 - ? (reconciliation.OutboxCount - reconciliation.Mismatches) / (double)reconciliation.OutboxCount * 100 - : 100.0; - - return new ReconciliationMetrics( - reconciliation.Mismatches, - completeness, - reconciliation.OutboxCount, - reconciliation.InboxProcessed, - clock.UtcNow.DateTime); - } - - public async Task GetModelDriftMetricsAsync(CancellationToken ct) - { - await using var connection = await connectionFactory.OpenAsync(ct); - - // Query shadow_run results for OOS performance metrics - var driftMetrics = await connection.QuerySingleAsync<(int Count, double AvgOos, int Degraded)>(""" - SELECT - COUNT(DISTINCT model_id) as Count, - COALESCE(AVG(CAST(validation_gates_json->>'dsr_above_95' AS float)), 0) as AvgOos, - COUNT(CASE - WHEN CAST(validation_gates_json->>'dsr_above_95' AS float) < 0.95 THEN 1 - END) as Degraded - FROM model_operations.shadow_run - WHERE published_at IS NOT NULL - AND published_at >= NOW() - INTERVAL '30 days' - """); - - // Baseline Sharpe (from most recent successful run) - var baseline = await connection.QuerySingleOrDefaultAsync(""" - SELECT CAST(validation_gates_json->>'sharpe' AS float) - FROM model_operations.shadow_run - WHERE status = 'EvaluationComplete' - AND published_at IS NOT NULL - ORDER BY published_at DESC - LIMIT 1 - """); - - return new ModelDriftMetrics( - driftMetrics.Count, - driftMetrics.AvgOos, - driftMetrics.Degraded, - baseline ?? 0.0, - clock.UtcNow.DateTime); - } -} diff --git a/src/KArtSell.Modules.ModelOperations/Observability/StubObservabilityService.cs b/src/KArtSell.Modules.ModelOperations/Observability/StubObservabilityService.cs new file mode 100644 index 00000000..857d8801 --- /dev/null +++ b/src/KArtSell.Modules.ModelOperations/Observability/StubObservabilityService.cs @@ -0,0 +1,41 @@ +namespace KArtSell.Modules.ModelOperations.Observability; + +public sealed class StubObservabilityService : IObservabilityService +{ + public async Task GetMetricsAsync(CancellationToken cancellationToken) + { + await Task.Delay(10, cancellationToken); + + return new ObservabilityMetricsDto + { + BatchSla = new BatchSlaMetrics + { + QueueDepth = 0, + AverageCompletionTimeSeconds = 0, + RetryRate = 0 + }, + DataQuality = new DataQualityMetrics + { + QuarantineCount = 0, + AgeMinutes = 0, + TopFailureReasons = new() + }, + DuplicateDetection = new DuplicateDetectionMetrics + { + ConstraintViolationCount = 0, + LastDetected = DateTime.UtcNow + }, + Reconciliation = new ReconciliationMetrics + { + CompletenessPercentage = 100, + AuditRecordsCount = 0 + }, + ModelDrift = new ModelDriftMetrics + { + OosPerformanceValue = 0, + BaselineComparison = 0, + DegradationFlag = false + } + }; + } +} diff --git a/src/KArtSell.Modules.ModelOperations/ShadowRun/Services/StubKrxDataService.cs b/src/KArtSell.Modules.ModelOperations/ShadowRun/Services/StubKrxDataService.cs new file mode 100644 index 00000000..02c1e69e --- /dev/null +++ b/src/KArtSell.Modules.ModelOperations/ShadowRun/Services/StubKrxDataService.cs @@ -0,0 +1,41 @@ +using KArtSell.Modules.ModelOperations.ShadowRun; +using Microsoft.Extensions.Logging; + +namespace KArtSell.Modules.ModelOperations.ShadowRun.Services; + +/// +/// Stub implementation of KRX data service for local development and testing. +/// Returns empty data sets. Replace with real implementation for production. +/// +public sealed class StubKrxDataService : IKrxDataService +{ + private readonly ILogger _logger; + + public StubKrxDataService(ILogger logger) + { + _logger = logger; + } + + public async Task> GetDailyOhlcvAsync( + string ticker, + DateOnly startDate, + DateOnly endDate, + CancellationToken cancellationToken) + { + _logger.LogWarning("Using stub KRX data service - no real market data available. Ticker: {Ticker}, Period: {Start}..{End}", + ticker, startDate, endDate); + await Task.Delay(10, cancellationToken); + return Array.Empty().AsReadOnly(); + } + + public async Task> GetFeeScheduleAsync( + DateOnly startDate, + DateOnly endDate, + CancellationToken cancellationToken) + { + _logger.LogWarning("Using stub KRX data service - no real fee schedule available. Period: {Start}..{End}", + startDate, endDate); + await Task.Delay(10, cancellationToken); + return Array.Empty().AsReadOnly(); + } +}