From 9a9e4640ead56875115bd8289d7f02910aa89a79 Mon Sep 17 00:00:00 2001 From: Ammar Heidari Date: Sun, 4 Oct 2026 15:42:42 +0330 Subject: [PATCH] fix(analytics): honor configured history limits end to end Signed-off-by: Ammar Heidari --- .../ConfluentKafkaConsumerGroupReadAdapter.cs | 106 +- .../AdoHistoricalMetricMaintenanceStore.cs | 61 +- ...ConsumerLagHistorySamplingHostedService.cs | 357 +++++ .../KafdeckOperationalAnalyticsEndpoints.cs | 217 +++ .../OperationalAnalyticsRuntimeService.cs | 908 ++++++++++++ src/backend/Kafdeck.Api/Program.cs | 60 +- .../Consumers/ConsumerReadContracts.cs | 15 + .../HistoricalMetricMaintenance.cs | 10 + .../OperationalAnalyticsTrends.cs | 258 ++++ .../V08W64OperationalAnalyticsRuntimeTests.cs | 1295 +++++++++++++++++ 10 files changed, 3271 insertions(+), 16 deletions(-) create mode 100644 src/backend/Kafdeck.Api/ConsumerLagHistorySamplingHostedService.cs create mode 100644 src/backend/Kafdeck.Api/KafdeckOperationalAnalyticsEndpoints.cs create mode 100644 src/backend/Kafdeck.Api/OperationalAnalyticsRuntimeService.cs create mode 100644 src/backend/Kafdeck.Core/Observability/OperationalAnalyticsTrends.cs create mode 100644 tests/Kafdeck.Architecture.Tests/V08W64OperationalAnalyticsRuntimeTests.cs diff --git a/src/backend/Infrastructure/Kafdeck.Infrastructure.Kafka/ConfluentKafkaConsumerGroupReadAdapter.cs b/src/backend/Infrastructure/Kafdeck.Infrastructure.Kafka/ConfluentKafkaConsumerGroupReadAdapter.cs index f408bbbe5..22b0d6aea 100644 --- a/src/backend/Infrastructure/Kafdeck.Infrastructure.Kafka/ConfluentKafkaConsumerGroupReadAdapter.cs +++ b/src/backend/Infrastructure/Kafdeck.Infrastructure.Kafka/ConfluentKafkaConsumerGroupReadAdapter.cs @@ -10,7 +10,7 @@ namespace Kafdeck.Infrastructure.Kafka; -public sealed class ConfluentKafkaConsumerGroupReadAdapter : IConsumerGroupReadPort, IDisposable +public sealed class ConfluentKafkaConsumerGroupReadAdapter : IConsumerGroupReadPort, IConsumerGroupSamplingReadPort, IDisposable { private readonly KafkaAdminClientRegistry _clients; private readonly TimeProvider _timeProvider; @@ -68,6 +68,110 @@ public Task>> ListGroupsAsync .ToArray(); }); + public Task> ListGroupPageAsync( + string clusterId, + string? afterGroupId, + int maxItems, + ReadViewOperationContext operation, + CancellationToken cancellationToken) + { + if (maxItems is < 1 or > 2_000) + { + throw new ArgumentOutOfRangeException( + nameof(maxItems)); + } + + if (afterGroupId is not null && + (string.IsNullOrWhiteSpace(afterGroupId) || + !string.Equals( + afterGroupId, + afterGroupId.Trim(), + StringComparison.Ordinal) || + afterGroupId.Any(char.IsControl))) + { + throw new ArgumentException( + "Consumer-group sampling cursor is invalid.", + nameof(afterGroupId)); + } + + return ExecuteAsync( + clusterId, + operation, + cancellationToken, + async (client, timeout, token) => + { + var result = + await client.ListConsumerGroupsAsync( + new ListConsumerGroupsOptions + { + RequestTimeout = timeout, + }) + .WaitAsync(token) + .ConfigureAwait(false); + + var selected = + new SortedSet( + StringComparer.Ordinal); + foreach (var group in result.Valid) + { + var groupId = + group.GroupId; + if (afterGroupId is not null && + string.CompareOrdinal( + groupId, + afterGroupId) <= 0) + { + continue; + } + + selected.Add( + groupId); + if (selected.Count > + maxItems) + { + selected.Remove( + selected.Max!); + } + } + + var items = + selected + .Take(maxItems) + .ToArray(); + + EnsureStringBudget( + items, + operation.MaxResponseBytes); + + string? nextCursor = null; + if (items.Length == + maxItems) + { + var last = + items[^1]; + if (result.Valid.Any(group => + string.CompareOrdinal( + group.GroupId, + last) > 0)) + { + nextCursor = + last; + } + } + + return new ConsumerGroupPage( + items + .Select(groupId => + new ConsumerGroupSummary( + groupId, + CoreConsumerGroupState.Unknown, + null, + false)) + .ToArray(), + nextCursor); + }); + } + public Task> GetGroupAsync( string clusterId, string groupId, diff --git a/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs b/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs index 5c6044694..44128b3b9 100644 --- a/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs +++ b/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs @@ -11,10 +11,12 @@ namespace Kafdeck.Infrastructure.Persistence; public sealed class AdoHistoricalMetricMaintenanceStore : - IHistoricalMetricMaintenanceStore + IHistoricalMetricMaintenanceStore, + IHistoricalMetricSamplingLeaseStore { private const int SchemaVersion = 2; private const int SingletonId = 1; + private const int SamplingSingletonId = 2; private const string Component = "historical-metrics-maintenance"; private const string RollupSource = @@ -460,9 +462,11 @@ EXECUTE FUNCTION kafdeck_enforce_history_exact_writer_fence() .ConfigureAwait(false); } - await using (var seed = - connection.CreateCommand()) + foreach (var singletonId in + new[] { SingletonId, SamplingSingletonId }) { + await using var seed = + connection.CreateCommand(); seed.Transaction = transaction; seed.CommandText = """ @@ -484,7 +488,7 @@ DO NOTHING AddParameter( seed, "@singleton_id", - SingletonId); + singletonId); AddParameter( seed, "@updated_at_utc", @@ -556,12 +560,39 @@ await transaction .ConfigureAwait(false); } - public async Task + public Task TryAcquireLeaseAsync( string ownerId, DateTimeOffset nowUtc, TimeSpan leaseDuration, - CancellationToken cancellationToken = default) + CancellationToken cancellationToken = default) => + TryAcquireLeaseAsync( + SingletonId, + ownerId, + nowUtc, + leaseDuration, + cancellationToken); + + public Task + TryAcquireSamplingLeaseAsync( + string ownerId, + DateTimeOffset nowUtc, + TimeSpan leaseDuration, + CancellationToken cancellationToken = default) => + TryAcquireLeaseAsync( + SamplingSingletonId, + ownerId, + nowUtc, + leaseDuration, + cancellationToken); + + private async Task + TryAcquireLeaseAsync( + int singletonId, + string ownerId, + DateTimeOffset nowUtc, + TimeSpan leaseDuration, + CancellationToken cancellationToken) { ValidateOwner(ownerId); @@ -595,6 +626,7 @@ await connection await LockAndReadLeaseAsync( connection, transaction, + singletonId, cancellationToken) .ConfigureAwait(false); @@ -649,7 +681,7 @@ UPDATE kafdeck_historical_metric_maintenance AddParameter( update, "@singleton_id", - SingletonId); + singletonId); if (await update .ExecuteNonQueryAsync(cancellationToken) @@ -961,9 +993,20 @@ private static CancellationTokenRegistration sqliteConnection); } + private Task LockAndReadLeaseAsync( + DbConnection connection, + DbTransaction transaction, + CancellationToken cancellationToken) => + LockAndReadLeaseAsync( + connection, + transaction, + SingletonId, + cancellationToken); + private async Task LockAndReadLeaseAsync( DbConnection connection, DbTransaction transaction, + int singletonId, CancellationToken cancellationToken) { if (!_connectionFactory.SupportsSelectForUpdate) @@ -980,7 +1023,7 @@ UPDATE kafdeck_historical_metric_maintenance AddParameter( lockCommand, "@singleton_id", - SingletonId); + singletonId); await lockCommand .ExecuteNonQueryAsync(cancellationToken) .ConfigureAwait(false); @@ -1011,7 +1054,7 @@ FROM kafdeck_historical_metric_maintenance AddParameter( command, "@singleton_id", - SingletonId); + singletonId); await using var reader = await command diff --git a/src/backend/Kafdeck.Api/ConsumerLagHistorySamplingHostedService.cs b/src/backend/Kafdeck.Api/ConsumerLagHistorySamplingHostedService.cs new file mode 100644 index 000000000..5e7114840 --- /dev/null +++ b/src/backend/Kafdeck.Api/ConsumerLagHistorySamplingHostedService.cs @@ -0,0 +1,357 @@ +using Kafdeck.Core.Consumers; +using Kafdeck.Core.Observability; +using Kafdeck.Infrastructure.Configuration; +using Kafdeck.Modules.Consumers; + +namespace Kafdeck.Api; + +public sealed record ConsumerLagHistorySamplingPolicy( + TimeSpan Interval, + int MaxGroupsPerCluster, + int MaxConcurrentGroups) +{ + public static ConsumerLagHistorySamplingPolicy Default { get; } = + new( + TimeSpan.FromMinutes(1), + 500, + 4); + + public void Validate() + { + if (Interval < TimeSpan.FromSeconds(10) || + Interval > TimeSpan.FromHours(1) || + MaxGroupsPerCluster is < 1 or > 2_000 || + MaxConcurrentGroups is < 1 or > 16) + { + throw new ArgumentOutOfRangeException( + nameof(Interval)); + } + } +} + +public sealed class ConsumerLagHistorySamplingHostedService : + BackgroundService +{ + private const long LargestExactlyRepresentableInteger = + 9_007_199_254_740_992L; + private static readonly TimeSpan SamplingLeaseDuration = + TimeSpan.FromMinutes(10); + private static readonly TimeSpan LeaseSafetyMargin = + TimeSpan.FromSeconds(5); + + private readonly KafdeckOptions _options; + private readonly ConsumerExplorerService _consumers; + private readonly IConsumerGroupSamplingReadPort _samplingGroups; + private readonly IHistoricalMetricStore _history; + private readonly IHistoricalMetricSamplingLeaseStore + _samplingLeases; + private readonly ConsumerLagHistorySamplingPolicy _policy; + private readonly ILogger< + ConsumerLagHistorySamplingHostedService> _logger; + private readonly TimeProvider _timeProvider; + private readonly string _ownerId; + private readonly Dictionary _clusterCursors = + new(StringComparer.Ordinal); + + public ConsumerLagHistorySamplingHostedService( + KafdeckOptions options, + ConsumerExplorerService consumers, + IConsumerGroupSamplingReadPort samplingGroups, + IHistoricalMetricStore history, + IHistoricalMetricSamplingLeaseStore samplingLeases, + ConsumerLagHistorySamplingPolicy policy, + ILogger logger, + TimeProvider? timeProvider = null) + { + _options = options ?? + throw new ArgumentNullException(nameof(options)); + _consumers = consumers ?? + throw new ArgumentNullException(nameof(consumers)); + _samplingGroups = samplingGroups ?? + throw new ArgumentNullException(nameof(samplingGroups)); + _history = history ?? + throw new ArgumentNullException(nameof(history)); + _samplingLeases = samplingLeases ?? + throw new ArgumentNullException(nameof(samplingLeases)); + _policy = policy ?? + throw new ArgumentNullException(nameof(policy)); + _policy.Validate(); + _logger = logger ?? + throw new ArgumentNullException(nameof(logger)); + _timeProvider = + timeProvider ?? + TimeProvider.System; + _ownerId = + BuildOwnerId(); + } + + protected override async Task ExecuteAsync( + CancellationToken stoppingToken) + { + await SampleOnceAsync( + stoppingToken) + .ConfigureAwait(false); + + using var timer = + new PeriodicTimer( + _policy.Interval, + _timeProvider); + + while (await timer + .WaitForNextTickAsync( + stoppingToken) + .ConfigureAwait(false)) + { + await SampleOnceAsync( + stoppingToken) + .ConfigureAwait(false); + } + } + + internal async Task SampleOnceAsync( + CancellationToken cancellationToken) + { + try + { + var nowUtc = + _timeProvider.GetUtcNow(); + var lease = + await _samplingLeases + .TryAcquireSamplingLeaseAsync( + _ownerId, + nowUtc, + SamplingLeaseDuration, + cancellationToken) + .ConfigureAwait(false); + + if (lease is null) + { + return; + } + + var remaining = + lease.ExpiresAtUtc - + _timeProvider.GetUtcNow() - + LeaseSafetyMargin; + if (remaining <= TimeSpan.Zero) + { + return; + } + + using var cycleDeadline = + CancellationTokenSource.CreateLinkedTokenSource( + cancellationToken); + cycleDeadline.CancelAfter( + remaining); + var cycleToken = + cycleDeadline.Token; + + var observedAtUtc = + AlignToSamplingWindow( + nowUtc, + _policy.Interval); + + foreach (var cluster in + _options.Clusters) + { + cycleToken + .ThrowIfCancellationRequested(); + + _clusterCursors.TryGetValue( + cluster.Id, + out var cursor); + var page = + await _samplingGroups + .ListGroupPageAsync( + cluster.Id, + cursor, + _policy.MaxGroupsPerCluster, + new Kafdeck.Core.ReadViews + .ReadViewOperationContext( + _timeProvider.GetUtcNow() + .AddSeconds(10), + _policy.MaxGroupsPerCluster, + 4 * 1024 * 1024), + cycleToken) + .ConfigureAwait(false); + + if (!page.IsSuccess || + page.Value is null) + { + continue; + } + + var selected = + page.Value.Items + .ToArray(); + + if (selected.Length == 0 && + cursor is not null) + { + page = + await _samplingGroups + .ListGroupPageAsync( + cluster.Id, + afterGroupId: null, + _policy.MaxGroupsPerCluster, + new Kafdeck.Core.ReadViews + .ReadViewOperationContext( + _timeProvider.GetUtcNow() + .AddSeconds(10), + _policy.MaxGroupsPerCluster, + 4 * 1024 * 1024), + cycleToken) + .ConfigureAwait(false); + + if (!page.IsSuccess || + page.Value is null) + { + continue; + } + + selected = + page.Value.Items + .ToArray(); + } + + if (selected.Length == 0) + { + _clusterCursors[cluster.Id] = + null; + continue; + } + + _clusterCursors[cluster.Id] = + page.Value.NextCursor; + + using var slots = + new SemaphoreSlim( + _policy + .MaxConcurrentGroups, + _policy + .MaxConcurrentGroups); + var samples = + new List(); + var gate = + new object(); + + await Task.WhenAll( + selected.Select( + async group => + { + await slots + .WaitAsync( + cycleToken) + .ConfigureAwait(false); + try + { + var lag = + await _consumers + .GetLagAsync( + cluster.Id, + group.GroupId, + cycleToken) + .ConfigureAwait(false); + + if (!lag.IsSuccess || + lag.Value?.TotalLag is + not long total || + total < 0) + { + return; + } + + var exact = + total <= + LargestExactlyRepresentableInteger; + var sample = + HistoricalMetricSample.Gauge( + new HistoricalMetricIdentity( + OperationalMetricHistoryNames + .ConsumerLagTotal, + cluster.Id, + "consumer_group", + group.GroupId), + observedAtUtc, + total, + "consumer-lag-sampler", + lag.Value.IsPartial || + !exact + ? "Partial" + : "Stable"); + + lock (gate) + { + samples.Add( + sample); + } + } + finally + { + slots.Release(); + } + })) + .ConfigureAwait(false); + + if (samples.Count > 0) + { + await _history + .AppendAsync( + samples + .OrderBy(sample => + sample.Identity.ResourceId, + StringComparer.Ordinal) + .ToArray(), + cycleToken) + .ConfigureAwait(false); + } + } + } + catch (OperationCanceledException) + when (cancellationToken + .IsCancellationRequested) + { + throw; + } + catch (OperationCanceledException) + { + _logger.LogWarning( + "Consumer lag history sampling exceeded its fenced lease budget; incomplete evidence remains missing."); + } + catch (Exception exception) + { + _logger.LogWarning( + exception, + "Consumer lag history sampling failed; missing evidence remains missing."); + } + } + + private static DateTimeOffset AlignToSamplingWindow( + DateTimeOffset value, + TimeSpan interval) + { + var utcTicks = + value.UtcDateTime.Ticks; + var alignedTicks = + utcTicks - + utcTicks % + interval.Ticks; + + return new DateTimeOffset( + alignedTicks, + TimeSpan.Zero); + } + + private static string BuildOwnerId() + { + var machine = + Environment.MachineName; + var process = + Environment.ProcessId; + var suffix = + Guid.NewGuid() + .ToString("N"); + + return $"{machine}:{process}:{suffix}"; + } +} diff --git a/src/backend/Kafdeck.Api/KafdeckOperationalAnalyticsEndpoints.cs b/src/backend/Kafdeck.Api/KafdeckOperationalAnalyticsEndpoints.cs new file mode 100644 index 000000000..cb09981f3 --- /dev/null +++ b/src/backend/Kafdeck.Api/KafdeckOperationalAnalyticsEndpoints.cs @@ -0,0 +1,217 @@ +using Kafdeck.Core.Observability; +using Kafdeck.Core.ReadViews; +using Kafdeck.Core.Security; +using Kafdeck.Infrastructure.Configuration; + +namespace Kafdeck.Api; + +public static class KafdeckOperationalAnalyticsEndpoints +{ + public static IEndpointRouteBuilder + MapKafdeckV08OperationalAnalytics( + this IEndpointRouteBuilder endpoints, + KafdeckOptions options) + { + ArgumentNullException.ThrowIfNull(endpoints); + ArgumentNullException.ThrowIfNull(options); + + endpoints.MapGet( + "/api/v1/clusters/{clusterId}/consumer-groups/{groupId}/analytics/live", + async ( + string clusterId, + string groupId, + OperationalAnalyticsRuntimeService analytics, + CancellationToken cancellationToken) => + { + if (!options.Clusters.Any(cluster => + string.Equals( + cluster.Id, + clusterId, + StringComparison.Ordinal))) + { + return ApiResults.Problem( + ApiProblemMapper + .InvalidClusterId( + clusterId)); + } + + var query = + new OperationalAnalyticsQuery( + clusterId, + OperationalResourceKind.ConsumerGroup, + groupId, + [ + OperationalMetricKind.ConsumerLagTotal, + OperationalMetricKind.ConsumerConsumeRecordsPerSecond, + ], + MaxItems: 8); + var result = + await analytics + .QueryAsync( + query, + new ReadViewOperationContext( + DateTimeOffset.UtcNow + .AddSeconds(10), + maxItems: 8, + maxResponseBytes: + 256 * 1024), + cancellationToken) + .ConfigureAwait(false); + + return result.IsSuccess && + result.Value is not null + ? Results.Ok(result.Value) + : ApiResults.Problem( + ApiProblemMapper.FromReadView( + result.Failure!)); + }) + .WithName("v08-consumer-operational-analytics-live") + .RequireKafdeckAuthorization( + AuthorizationAction.ConsumerRead, + "clusterId", + "groupId"); + + endpoints.MapGet( + "/api/v1/clusters/{clusterId}/consumer-groups/{groupId}/analytics/trend", + async ( + string clusterId, + string groupId, + DateTimeOffset? from, + DateTimeOffset? to, + int? maxPoints, + OperationalTrendService trends, + CancellationToken cancellationToken) => + { + if (!options.Clusters.Any(cluster => + string.Equals( + cluster.Id, + clusterId, + StringComparison.Ordinal))) + { + return ApiResults.Problem( + ApiProblemMapper + .InvalidClusterId( + clusterId)); + } + + var toUtc = + to ?? + DateTimeOffset.UtcNow; + var fromUtc = + from ?? + toUtc.AddHours(-1); + + try + { + var result = + await trends.QueryAsync( + new OperationalTrendQuery( + OperationalMetricKind + .ConsumerLagTotal, + new OperationalResourceIdentity( + clusterId, + OperationalResourceKind + .ConsumerGroup, + groupId), + fromUtc, + toUtc, + maxPoints ?? + trends.DefaultMaxPoints), + cancellationToken) + .ConfigureAwait(false); + return Results.Ok(result); + } + catch (ArgumentException exception) + { + return Results.BadRequest( + new + { + code = + "invalid_operational_trend_query", + message = + exception.Message, + }); + } + }) + .WithName("v08-consumer-operational-analytics-trend") + .RequireKafdeckAuthorization( + AuthorizationAction.ConsumerRead, + "clusterId", + "groupId"); + + endpoints.MapGet( + "/api/v1/clusters/{clusterId}/consumer-groups/{groupId}/analytics/slo", + async ( + string clusterId, + string groupId, + double threshold, + double target, + DateTimeOffset? from, + DateTimeOffset? to, + int? maxPoints, + OperationalTrendService trends, + CancellationToken cancellationToken) => + { + if (!options.Clusters.Any(cluster => + string.Equals( + cluster.Id, + clusterId, + StringComparison.Ordinal))) + { + return ApiResults.Problem( + ApiProblemMapper + .InvalidClusterId( + clusterId)); + } + + var toUtc = + to ?? + DateTimeOffset.UtcNow; + var fromUtc = + from ?? + toUtc.AddHours(-1); + + try + { + var result = + await trends.EvaluateSloAsync( + new OperationalSloDefinition( + "consumer-lag", + OperationalMetricKind + .ConsumerLagTotal, + new OperationalResourceIdentity( + clusterId, + OperationalResourceKind + .ConsumerGroup, + groupId), + threshold, + target), + fromUtc, + toUtc, + maxPoints ?? + trends.DefaultMaxPoints, + cancellationToken) + .ConfigureAwait(false); + return Results.Ok(result); + } + catch (ArgumentException exception) + { + return Results.BadRequest( + new + { + code = + "invalid_operational_slo_query", + message = + exception.Message, + }); + } + }) + .WithName("v08-consumer-operational-slo") + .RequireKafdeckAuthorization( + AuthorizationAction.ConsumerRead, + "clusterId", + "groupId"); + + return endpoints; + } +} diff --git a/src/backend/Kafdeck.Api/OperationalAnalyticsRuntimeService.cs b/src/backend/Kafdeck.Api/OperationalAnalyticsRuntimeService.cs new file mode 100644 index 000000000..47c3946ed --- /dev/null +++ b/src/backend/Kafdeck.Api/OperationalAnalyticsRuntimeService.cs @@ -0,0 +1,908 @@ +using Kafdeck.Core.Ecosystem; +using Kafdeck.Core.Observability; +using Kafdeck.Core.ReadViews; +using Kafdeck.Modules.Consumers; + +namespace Kafdeck.Api; + +public sealed class HistoricalConsumerHistoryObservationPort : + IHistoryObservationPort +{ + private readonly IHistoricalMetricStore _store; + private readonly HistoricalMetricStorePolicy _storePolicy; + private readonly TimeProvider _timeProvider; + + public HistoricalConsumerHistoryObservationPort( + IHistoricalMetricStore store, + HistoricalMetricStorePolicy? storePolicy = null, + TimeProvider? timeProvider = null) + { + _store = store ?? + throw new ArgumentNullException(nameof(store)); + _storePolicy = + storePolicy ?? + new HistoricalMetricStorePolicy( + HistoricalMetricQuery.HardMaxRange, + HistoricalMetricQuery.HardMaxSeries, + HistoricalMetricQuery.HardMaxPoints, + TimeSpan.FromSeconds(30), + 1); + _storePolicy.Validate(); + _timeProvider = timeProvider ?? + TimeProvider.System; + } + + public async Task>> + GetConsumerHistoryAsync( + string clusterId, + string groupId, + ReadViewOperationContext operation, + CancellationToken cancellationToken) + { + ArgumentException.ThrowIfNullOrWhiteSpace(clusterId); + ArgumentException.ThrowIfNullOrWhiteSpace(groupId); + ArgumentNullException.ThrowIfNull(operation); + + var now = + _timeProvider.GetUtcNow(); + var historyRange = + _storePolicy.MaxQueryRange < + TimeSpan.FromHours(1) + ? _storePolicy.MaxQueryRange + : TimeSpan.FromHours(1); + var from = + now.Subtract( + historyRange); + var remaining = + operation.DeadlineUtc - + now; + + if (remaining <= TimeSpan.Zero) + { + return ReadViewResult> + .Failed( + new ReadViewFailure( + ReadViewFailureCategory.Timeout, + "consumer_history_deadline_expired", + "Consumer history deadline expired before the provider query.", + true)); + } + + using var deadline = + CancellationTokenSource + .CreateLinkedTokenSource( + cancellationToken); + deadline.CancelAfter( + remaining); + + try + { + var result = + await _store.QueryAsync( + new HistoricalMetricQuery( + OperationalMetricHistoryNames + .ConsumerLagTotal, + clusterId, + "consumer_group", + groupId, + from, + now, + MaxSeries: 1, + MaxPoints: + Math.Min( + Math.Min( + operation.MaxItems, + 512), + _storePolicy + .MaxPointsPerQuery)), + deadline.Token) + .ConfigureAwait(false); + + var history = + result.Series + .SelectMany(series => + series.Points) + .OrderBy(point => + point.ObservedAtUtc) + .Take( + operation.MaxItems) + .Select(point => + { + var average = + point.Average; + long? totalLag = null; + var state = + point.State ?? + "Stable"; + + if (double.IsFinite(average) && + average >= 0 && + average < + 9_223_372_036_854_775_808d) + { + totalLag = + checked( + (long)Math.Round( + average, + MidpointRounding + .AwayFromZero)); + + if (average > + 9_007_199_254_740_992d && + string.Equals( + state, + "Stable", + StringComparison.Ordinal)) + { + state = + "Partial"; + } + } + else + { + state = + "Partial"; + } + + return new ConsumerHistoryObservation( + point.ObservedAtUtc, + state, + totalLag, + point.Source); + }) + .ToArray(); + + var limitations = + result.Truncated + ? new[] + { + new ReadViewLimitation( + "consumer_history_truncated", + result.LimitReason ?? + "Consumer history was truncated by the configured provider limits."), + } + : Array.Empty(); + + return ReadViewResult> + .Success( + Array.AsReadOnly( + history), + limitations); + } + catch (OperationCanceledException) + when (cancellationToken.IsCancellationRequested) + { + throw; + } + catch (OperationCanceledException) + { + return ReadViewResult> + .Failed( + new ReadViewFailure( + ReadViewFailureCategory.Timeout, + "consumer_history_timeout", + "Consumer history query timed out.", + true)); + } + catch (TimeoutException) + { + return ReadViewResult> + .Failed( + new ReadViewFailure( + ReadViewFailureCategory.Timeout, + "consumer_history_timeout", + "Consumer history query timed out.", + true)); + } + catch (Exception) + { + return ReadViewResult> + .Failed( + new ReadViewFailure( + ReadViewFailureCategory.Unavailable, + "consumer_history_unavailable", + "Consumer history provider is unavailable.", + true)); + } + } +} + +public sealed class OperationalAnalyticsRuntimeService : + IOperationalAnalyticsObservationPort +{ + private readonly ConsumerExplorerService _consumers; + private readonly IMetricsObservationPort _metrics; + private readonly TimeProvider _timeProvider; + + public OperationalAnalyticsRuntimeService( + ConsumerExplorerService consumers, + IMetricsObservationPort metrics, + TimeProvider? timeProvider = null) + { + _consumers = consumers ?? + throw new ArgumentNullException(nameof(consumers)); + _metrics = metrics ?? + throw new ArgumentNullException(nameof(metrics)); + _timeProvider = + timeProvider ?? + TimeProvider.System; + } + + public async Task> + QueryAsync( + OperationalAnalyticsQuery query, + ReadViewOperationContext operation, + CancellationToken cancellationToken) + { + ArgumentNullException.ThrowIfNull(query); + ArgumentNullException.ThrowIfNull(operation); + query.Validate(); + + if (query.ResourceKind is null || + string.IsNullOrWhiteSpace( + query.ResourceId)) + { + return ReadViewResult + .Failed( + new ReadViewFailure( + ReadViewFailureCategory.InvalidRequest, + "analytics_target_required", + "Operational analytics runtime queries currently require one explicit resource.", + false)); + } + + var resource = + new OperationalResourceIdentity( + query.ClusterId, + query.ResourceKind.Value, + query.ResourceId); + resource.Validate(); + + var items = + new List( + query.Metrics.Count); + + foreach (var metric in query.Metrics) + { + items.Add( + await ObserveOneAsync( + metric, + resource, + operation, + cancellationToken) + .ConfigureAwait(false)); + } + + var result = + new OperationalAnalyticsResult( + Array.AsReadOnly( + items.ToArray()), + Truncated: false, + LimitReason: null, + query.Metrics.ToArray()); + result.Validate(query); + + return ReadViewResult + .Success(result); + } + + private async Task + ObserveOneAsync( + OperationalMetricKind metric, + OperationalResourceIdentity resource, + ReadViewOperationContext operation, + CancellationToken cancellationToken) + { + if (resource.Kind == + OperationalResourceKind.ConsumerGroup && + metric == + OperationalMetricKind.ConsumerLagTotal) + { + var lag = + await _consumers + .GetLagAsync( + resource.ClusterId, + resource.ResourceId, + cancellationToken) + .ConfigureAwait(false); + + if (!lag.IsSuccess || + lag.Value is null) + { + return Unavailable( + metric, + resource, + "consumer_lag_unavailable"); + } + + if (lag.Value.TotalLag is not long total) + { + return Unknown( + metric, + resource, + "consumer_lag_unknown"); + } + + var precise = + total <= + 9_007_199_254_740_992L; + + return new OperationalMetricEvidence( + metric, + resource, + total, + _timeProvider.GetUtcNow(), + Window: null, + "consumer-lag-read-view", + lag.Value.IsPartial || + !precise + ? OperationalEvidenceState.Partial + : OperationalEvidenceState.Available); + } + + if (resource.Kind == + OperationalResourceKind.ConsumerGroup && + metric == + OperationalMetricKind + .ConsumerConsumeRecordsPerSecond) + { + var rates = + await _metrics + .GetConsumerRateAsync( + resource.ClusterId, + resource.ResourceId, + operation, + cancellationToken) + .ConfigureAwait(false); + + if (!rates.IsSuccess || + rates.Value is null || + rates.Value.ConsumeRecordsPerSecond is + not double value || + !double.IsFinite(value) || + value < 0 || + rates.Value.Window <= + TimeSpan.Zero) + { + return Unavailable( + metric, + resource, + "consumer_rate_provider_unavailable"); + } + + var age = + _timeProvider.GetUtcNow() - + rates.Value.ObservedAt; + if (age < TimeSpan.FromSeconds(-5) || + age > rates.Value.Window) + { + return new OperationalMetricEvidence( + metric, + resource, + value, + rates.Value.ObservedAt, + rates.Value.Window, + rates.Value.Source, + OperationalEvidenceState.Stale); + } + + return new OperationalMetricEvidence( + metric, + resource, + value, + rates.Value.ObservedAt, + rates.Value.Window, + rates.Value.Source, + OperationalEvidenceState.Available); + } + + return Unavailable( + metric, + resource, + "metric_provider_unavailable"); + } + + private static OperationalMetricEvidence Unavailable( + OperationalMetricKind metric, + OperationalResourceIdentity resource, + string source) => + new( + metric, + resource, + Value: null, + ObservedAtUtc: null, + Window: null, + source, + OperationalEvidenceState.Unavailable); + + private static OperationalMetricEvidence Unknown( + OperationalMetricKind metric, + OperationalResourceIdentity resource, + string source) => + new( + metric, + resource, + Value: null, + ObservedAtUtc: null, + Window: null, + source, + OperationalEvidenceState.Unknown); +} + +public sealed class OperationalTrendService +{ + private readonly IHistoricalMetricStore? _history; + private readonly HistoricalMetricStorePolicy _storePolicy; + private readonly TimeSpan _expectedRawCadence; + + public OperationalTrendService( + IHistoricalMetricStore? history = null, + HistoricalMetricStorePolicy? storePolicy = null, + ConsumerLagHistorySamplingPolicy? samplingPolicy = null) + { + _history = history; + _storePolicy = + storePolicy ?? + new HistoricalMetricStorePolicy( + HistoricalMetricQuery.HardMaxRange, + HistoricalMetricQuery.HardMaxSeries, + HistoricalMetricQuery.HardMaxPoints, + TimeSpan.FromSeconds(30), + 1); + _storePolicy.Validate(); + _expectedRawCadence = + (samplingPolicy ?? + ConsumerLagHistorySamplingPolicy.Default) + .Interval; + } + + public int DefaultMaxPoints => + Math.Min( + 1_000, + _storePolicy.MaxPointsPerQuery); + + public TimeSpan MaxQueryRange => + _storePolicy.MaxQueryRange; + + public async Task + QueryAsync( + OperationalTrendQuery query, + CancellationToken cancellationToken) + { + ArgumentNullException.ThrowIfNull(query); + query.Validate(); + + if (query.MaxPoints > + _storePolicy.MaxPointsPerQuery) + { + throw new ArgumentOutOfRangeException( + nameof(query), + "Operational trend query exceeds the configured historical-metrics point limit."); + } + + if (query.ToUtc - query.FromUtc > + _storePolicy.MaxQueryRange) + { + throw new ArgumentOutOfRangeException( + nameof(query), + "Operational trend query exceeds the configured historical-metrics range."); + } + + var metricName = + OperationalMetricHistoryNames.TryGet( + query.Metric); + + if (_history is null || + metricName is null) + { + return Empty( + query, + OperationalTrendState.Unavailable, + "history_provider_unavailable"); + } + + HistoricalMetricQueryResult raw; + try + { + raw = + await _history.QueryAsync( + new HistoricalMetricQuery( + metricName, + query.Resource.ClusterId, + OperationalMetricHistoryNames + .ResourceKind( + query.Resource.Kind), + query.Resource.ResourceId, + query.FromUtc, + query.ToUtc, + MaxSeries: 1, + query.MaxPoints), + cancellationToken) + .ConfigureAwait(false); + } + catch (OperationCanceledException) + when (cancellationToken.IsCancellationRequested) + { + throw; + } + catch (OperationCanceledException) + { + return Empty( + query, + OperationalTrendState.Unavailable, + "history_provider_timeout"); + } + catch (TimeoutException) + { + return Empty( + query, + OperationalTrendState.Unavailable, + "history_provider_timeout"); + } + catch (Exception) + { + return Empty( + query, + OperationalTrendState.Unavailable, + "history_provider_unavailable"); + } + + var points = + raw.Series + .SelectMany(series => + series.Points) + .OrderBy(point => + point.ObservedAtUtc) + .Take( + query.MaxPoints) + .Select(MapPoint) + .ToArray(); + + var state = + points.Length == 0 + ? OperationalTrendState.Unknown + : raw.Truncated || + points.Any(point => + point.State is not + OperationalEvidenceState.Available) + ? OperationalTrendState.Partial + : OperationalTrendState.Available; + + var result = + new OperationalTrendResult( + query.Metric, + query.Resource, + Array.AsReadOnly(points), + raw.Truncated, + raw.LimitReason, + raw.Provider, + state); + result.Validate(query); + return result; + } + + public async Task + EvaluateSloAsync( + OperationalSloDefinition definition, + DateTimeOffset fromUtc, + DateTimeOffset toUtc, + int maxPoints, + CancellationToken cancellationToken) + { + ArgumentNullException.ThrowIfNull(definition); + definition.Validate(); + + var trend = + await QueryAsync( + new OperationalTrendQuery( + definition.Metric, + definition.Resource, + fromUtc, + toUtc, + maxPoints), + cancellationToken) + .ConfigureAwait(false); + + var usable = + trend.Points + .Where(point => + point.State == + OperationalEvidenceState.Available && + point.HasKnownCoverage) + .ToArray(); + + if (usable.Length == 0) + { + var unavailable = + new OperationalSloResult( + definition, + fromUtc, + toUtc, + 0, + 0, + null, + null, + trend.State == + OperationalTrendState.Unavailable + ? OperationalEvidenceState.Unavailable + : OperationalEvidenceState.Unknown, + trend.State == + OperationalTrendState.Unavailable + ? "history_provider_unavailable" + : "no_complete_slo_evidence"); + unavailable.Validate(); + return unavailable; + } + + long weightedGood = 0; + long weightedBad = 0; + var hasAmbiguousAggregate = + false; + + foreach (var point in usable) + { + if (point.Max <= + definition.MaximumGoodValue) + { + weightedGood = + checked( + weightedGood + + point.Count); + continue; + } + + if (point.Min > + definition.MaximumGoodValue) + { + weightedBad = + checked( + weightedBad + + point.Count); + continue; + } + + hasAmbiguousAggregate = + true; + } + + if (hasAmbiguousAggregate) + { + var ambiguous = + new OperationalSloResult( + definition, + fromUtc, + toUtc, + 0, + 0, + null, + null, + OperationalEvidenceState.Partial, + "aggregate_threshold_ambiguous"); + ambiguous.Validate(); + return ambiguous; + } + + var weightedTotal = + checked( + weightedGood + + weightedBad); + if (weightedTotal > + int.MaxValue) + { + var oversized = + new OperationalSloResult( + definition, + fromUtc, + toUtc, + 0, + 0, + null, + null, + OperationalEvidenceState.Partial, + "slo_sample_count_overflow"); + oversized.Validate(); + return oversized; + } + + var evaluated = + (int)weightedTotal; + var good = + (int)weightedGood; + var compliance = + evaluated == 0 + ? 0d + : (double)good / + evaluated; + var badFraction = + 1d - + compliance; + var errorBudget = + 1d - + definition.TargetFraction; + var burn = + badFraction / + errorBudget; + + var completeCoverage = + HasCompleteWindowCoverage( + usable, + fromUtc, + toUtc, + _expectedRawCadence); + + var state = + usable.Length != + trend.Points.Count || + trend.State == + OperationalTrendState.Partial || + trend.Truncated || + !completeCoverage + ? OperationalEvidenceState.Partial + : OperationalEvidenceState.Available; + + var result = + new OperationalSloResult( + definition, + fromUtc, + toUtc, + evaluated, + good, + compliance, + burn, + state, + state == + OperationalEvidenceState.Partial + ? "partial_slo_evidence" + : null); + result.Validate(); + return result; + } + + private static OperationalTrendPoint MapPoint( + HistoricalMetricSample point) + { + var state = + point.State switch + { + "Stable" or null => + OperationalEvidenceState.Available, + "Partial" => + OperationalEvidenceState.Partial, + "Stale" => + OperationalEvidenceState.Stale, + "Unavailable" => + OperationalEvidenceState.Unavailable, + "Unknown" => + OperationalEvidenceState.Unknown, + _ => + OperationalEvidenceState.Unknown, + }; + + return new OperationalTrendPoint( + point.ObservedAtUtc, + point.Min, + point.Max, + point.Average, + point.Count, + point.ResolutionSeconds, + state, + point.HasKnownCoverage, + point.EffectiveFirstObservedAtUtc, + point.EffectiveLastObservedAtUtc); + } + + private static bool HasCompleteWindowCoverage( + IReadOnlyList points, + DateTimeOffset fromUtc, + DateTimeOffset toUtc, + TimeSpan expectedRawCadence) + { + if (points.Count == 0) + { + return false; + } + + var ordered = + points + .OrderBy(point => + point.EffectiveFirstObservedAtUtc) + .ToArray(); + + var first = + ordered[0]; + var firstStart = + first.EffectiveFirstObservedAtUtc; + var lastCovered = + first.EffectiveLastObservedAtUtc; + + if (firstStart is null || + lastCovered is null || + firstStart.Value - fromUtc > + CoverageCadence( + first, + expectedRawCadence)) + { + return false; + } + + for (var index = 1; + index < ordered.Length; + index++) + { + var point = + ordered[index]; + var start = + point.EffectiveFirstObservedAtUtc; + var end = + point.EffectiveLastObservedAtUtc; + + if (start is null || + end is null) + { + return false; + } + + var allowedGap = + Max( + CoverageCadence( + ordered[index - 1], + expectedRawCadence), + CoverageCadence( + point, + expectedRawCadence)); + + if (start.Value - + lastCovered.Value > + allowedGap) + { + return false; + } + + if (end.Value > + lastCovered.Value) + { + lastCovered = + end; + } + } + + return toUtc - + lastCovered.Value <= + CoverageCadence( + ordered[^1], + expectedRawCadence); + } + + private static TimeSpan CoverageCadence( + OperationalTrendPoint point, + TimeSpan expectedRawCadence) => + point.ResolutionSeconds > 0 + ? TimeSpan.FromSeconds( + point.ResolutionSeconds) + : expectedRawCadence; + + private static TimeSpan Max( + TimeSpan left, + TimeSpan right) => + left >= right + ? left + : right; + + private static OperationalTrendResult Empty( + OperationalTrendQuery query, + OperationalTrendState state, + string reason) + { + var result = + new OperationalTrendResult( + query.Metric, + query.Resource, + Array.Empty(), + Truncated: false, + LimitReason: null, + reason, + state); + result.Validate(query); + return result; + } +} diff --git a/src/backend/Kafdeck.Api/Program.cs b/src/backend/Kafdeck.Api/Program.cs index fa677221f..d82afa45c 100644 --- a/src/backend/Kafdeck.Api/Program.cs +++ b/src/backend/Kafdeck.Api/Program.cs @@ -145,11 +145,19 @@ historicalMetricsOptions.ConnectionString is not null kafdeckOptions.Clusters, secretResolver), services.GetRequiredService())); -builder.Services.AddSingleton(services => - new TelemetryConsumerGroupReadPort( +builder.Services.AddSingleton( + _ => new ConfluentKafkaConsumerGroupReadAdapter( kafdeckOptions.Clusters, - secretResolver), + secretResolver)); +builder.Services.AddSingleton( + services => + services.GetRequiredService< + ConfluentKafkaConsumerGroupReadAdapter>()); +builder.Services.AddSingleton(services => + new TelemetryConsumerGroupReadPort( + services.GetRequiredService< + ConfluentKafkaConsumerGroupReadAdapter>(), services.GetRequiredService())); builder.Services.AddSingleton(); @@ -207,19 +215,58 @@ historicalMetricsOptions.ConnectionString is not null .DefaultMaxCycleDuration)); builder.Services.AddSingleton< - IHistoricalMetricMaintenanceStore>( + AdoHistoricalMetricMaintenanceStore>( services => new AdoHistoricalMetricMaintenanceStore( services.GetRequiredService< IHistoricalMetricsDbConnectionFactory>())); + builder.Services.AddSingleton< + IHistoricalMetricMaintenanceStore>( + services => + services.GetRequiredService< + AdoHistoricalMetricMaintenanceStore>()); + builder.Services.AddSingleton< + IHistoricalMetricSamplingLeaseStore>( + services => + services.GetRequiredService< + AdoHistoricalMetricMaintenanceStore>()); builder.Services.AddHostedService< HistoricalMetricMaintenanceHostedService>(); + + builder.Services.AddSingleton( + services => + new HistoricalConsumerHistoryObservationPort( + services.GetRequiredService< + IHistoricalMetricStore>(), + services.GetRequiredService< + HistoricalMetricStorePolicy>())); + builder.Services.AddSingleton( + ConsumerLagHistorySamplingPolicy.Default); + builder.Services.AddHostedService< + ConsumerLagHistorySamplingHostedService>(); +} +else +{ + builder.Services.AddSingleton(); } -builder.Services.AddSingleton(); -builder.Services.AddSingleton(); +builder.Services.AddSingleton(); builder.Services.AddSingleton(); +builder.Services.AddSingleton(); +builder.Services.AddSingleton( + services => + services.GetRequiredService< + OperationalAnalyticsRuntimeService>()); +builder.Services.AddSingleton( + services => + new OperationalTrendService( + services.GetService< + IHistoricalMetricStore>(), + services.GetService< + HistoricalMetricStorePolicy>())); builder.Services.AddSingleton(services => new TelemetryRecordSchemaReadPort( new ConfluentSchemaRegistryReadAdapter( @@ -595,6 +642,7 @@ await historicalMetricMaintenanceStore app.MapKafdeckV07Streaming(kafdeckOptions); app.MapKafdeckV07ControlledSerde(kafdeckOptions); app.MapKafdeckV08Observability(kafdeckOptions); +app.MapKafdeckV08OperationalAnalytics(kafdeckOptions); app.MapKafdeckFleetCapabilities(); app.MapKafdeckV06OpenApi(); app.MapKafdeckV07OpenApi(); diff --git a/src/backend/Kafdeck.Core/Consumers/ConsumerReadContracts.cs b/src/backend/Kafdeck.Core/Consumers/ConsumerReadContracts.cs index 337defd1c..8848da6e9 100644 --- a/src/backend/Kafdeck.Core/Consumers/ConsumerReadContracts.cs +++ b/src/backend/Kafdeck.Core/Consumers/ConsumerReadContracts.cs @@ -79,3 +79,18 @@ Task>> GetOffsetsAsync( ReadViewOperationContext operation, CancellationToken cancellationToken); } + + +public sealed record ConsumerGroupPage( + IReadOnlyList Items, + string? NextCursor); + +public interface IConsumerGroupSamplingReadPort +{ + Task> ListGroupPageAsync( + string clusterId, + string? afterGroupId, + int maxItems, + ReadViewOperationContext operation, + CancellationToken cancellationToken); +} diff --git a/src/backend/Kafdeck.Core/Observability/HistoricalMetricMaintenance.cs b/src/backend/Kafdeck.Core/Observability/HistoricalMetricMaintenance.cs index c831a4791..8eb77d728 100644 --- a/src/backend/Kafdeck.Core/Observability/HistoricalMetricMaintenance.cs +++ b/src/backend/Kafdeck.Core/Observability/HistoricalMetricMaintenance.cs @@ -116,3 +116,13 @@ Task RunCycleAsync( HistoricalMetricMaintenancePolicy policy, CancellationToken cancellationToken = default); } + + +public interface IHistoricalMetricSamplingLeaseStore +{ + Task TryAcquireSamplingLeaseAsync( + string ownerId, + DateTimeOffset nowUtc, + TimeSpan leaseDuration, + CancellationToken cancellationToken = default); +} diff --git a/src/backend/Kafdeck.Core/Observability/OperationalAnalyticsTrends.cs b/src/backend/Kafdeck.Core/Observability/OperationalAnalyticsTrends.cs new file mode 100644 index 000000000..052d01eab --- /dev/null +++ b/src/backend/Kafdeck.Core/Observability/OperationalAnalyticsTrends.cs @@ -0,0 +1,258 @@ +namespace Kafdeck.Core.Observability; + +public enum OperationalTrendState +{ + Available = 1, + Partial = 2, + Unavailable = 3, + Unknown = 4, +} + +public sealed record OperationalTrendQuery( + OperationalMetricKind Metric, + OperationalResourceIdentity Resource, + DateTimeOffset FromUtc, + DateTimeOffset ToUtc, + int MaxPoints) +{ + public const int HardMaxPoints = 10_000; + + public void Validate() + { + if (!Enum.IsDefined(Metric)) + { + throw new ArgumentOutOfRangeException(nameof(Metric)); + } + + ArgumentNullException.ThrowIfNull(Resource); + Resource.Validate(); + OperationalMetricCompatibility.Validate( + Metric, + Resource.Kind); + + if (FromUtc == default || + ToUtc == default || + ToUtc <= FromUtc || + ToUtc - FromUtc > + HistoricalMetricQuery.HardMaxRange) + { + throw new ArgumentOutOfRangeException( + nameof(ToUtc)); + } + + if (MaxPoints is < 1 or > HardMaxPoints) + { + throw new ArgumentOutOfRangeException( + nameof(MaxPoints)); + } + } +} + +public sealed record OperationalTrendPoint( + DateTimeOffset ObservedAtUtc, + double Min, + double Max, + double Average, + long Count, + int ResolutionSeconds, + OperationalEvidenceState State, + bool HasKnownCoverage, + DateTimeOffset? FirstObservedAtUtc = null, + DateTimeOffset? LastObservedAtUtc = null) +{ + public DateTimeOffset? EffectiveFirstObservedAtUtc => + HasKnownCoverage + ? FirstObservedAtUtc ?? ObservedAtUtc + : null; + + public DateTimeOffset? EffectiveLastObservedAtUtc => + HasKnownCoverage + ? LastObservedAtUtc ?? ObservedAtUtc + : null; +} + +public sealed record OperationalTrendResult( + OperationalMetricKind Metric, + OperationalResourceIdentity Resource, + IReadOnlyList Points, + bool Truncated, + string? LimitReason, + string Provider, + OperationalTrendState State) +{ + public void Validate( + OperationalTrendQuery query) + { + ArgumentNullException.ThrowIfNull(query); + query.Validate(); + ArgumentNullException.ThrowIfNull(Points); + ArgumentException.ThrowIfNullOrWhiteSpace(Provider); + + if (Metric != query.Metric || + Resource != query.Resource || + Points.Count > query.MaxPoints || + !Enum.IsDefined(State)) + { + throw new ArgumentException( + "Operational trend result is inconsistent with its query."); + } + + if (Truncated != + !string.IsNullOrWhiteSpace(LimitReason)) + { + throw new ArgumentException( + "Truncated operational trends require an explicit limit reason."); + } + + if (State is + OperationalTrendState.Unavailable or + OperationalTrendState.Unknown && + Points.Count != 0) + { + throw new ArgumentException( + "Unavailable/unknown operational trends cannot fabricate points."); + } + + foreach (var point in Points) + { + if (point.ObservedAtUtc == default || + !double.IsFinite(point.Min) || + !double.IsFinite(point.Max) || + !double.IsFinite(point.Average) || + point.Min > point.Max || + point.Count < 1 || + point.ResolutionSeconds < 0 || + !Enum.IsDefined(point.State) || + ((point.FirstObservedAtUtc is null) != + (point.LastObservedAtUtc is null)) || + (point.FirstObservedAtUtc is not null && + (point.FirstObservedAtUtc.Value == default || + point.LastObservedAtUtc!.Value < + point.FirstObservedAtUtc.Value))) + { + throw new ArgumentException( + "Operational trend contains invalid evidence."); + } + } + } +} + +public sealed record OperationalSloDefinition( + string Id, + OperationalMetricKind Metric, + OperationalResourceIdentity Resource, + double MaximumGoodValue, + double TargetFraction) +{ + public const int MaxIdLength = 128; + + public void Validate() + { + ArgumentException.ThrowIfNullOrWhiteSpace(Id); + + if (!string.Equals( + Id, + Id.Trim(), + StringComparison.Ordinal) || + Id.Length > MaxIdLength || + Id.Any(char.IsControl) || + !double.IsFinite(MaximumGoodValue) || + MaximumGoodValue < 0 || + !double.IsFinite(TargetFraction) || + TargetFraction <= 0 || + TargetFraction >= 1) + { + throw new ArgumentException( + "Operational SLO definition is invalid."); + } + + ArgumentNullException.ThrowIfNull(Resource); + Resource.Validate(); + OperationalMetricCompatibility.Validate( + Metric, + Resource.Kind); + } +} + +public sealed record OperationalSloResult( + OperationalSloDefinition Definition, + DateTimeOffset FromUtc, + DateTimeOffset ToUtc, + int EvaluatedPoints, + int GoodPoints, + double? ComplianceFraction, + double? BurnRate, + OperationalEvidenceState State, + string? ReasonCode) +{ + public void Validate() + { + ArgumentNullException.ThrowIfNull(Definition); + Definition.Validate(); + + if (FromUtc == default || + ToUtc <= FromUtc || + EvaluatedPoints < 0 || + GoodPoints < 0 || + GoodPoints > EvaluatedPoints || + !Enum.IsDefined(State)) + { + throw new ArgumentException( + "Operational SLO result is invalid."); + } + + if (EvaluatedPoints == 0) + { + if (ComplianceFraction is not null || + BurnRate is not null) + { + throw new ArgumentException( + "SLO without evaluated evidence cannot carry numeric compliance."); + } + + return; + } + + if (ComplianceFraction is null || + BurnRate is null || + !double.IsFinite(ComplianceFraction.Value) || + !double.IsFinite(BurnRate.Value) || + ComplianceFraction.Value is < 0 or > 1 || + BurnRate.Value < 0) + { + throw new ArgumentException( + "Operational SLO numeric evidence is invalid."); + } + } +} + +public static class OperationalMetricHistoryNames +{ + public const string ConsumerLagTotal = + "consumer.lag.total"; + + public static string? TryGet( + OperationalMetricKind metric) => + metric switch + { + OperationalMetricKind.ConsumerLagTotal => + ConsumerLagTotal, + _ => null, + }; + + public static string ResourceKind( + OperationalResourceKind kind) => + kind switch + { + OperationalResourceKind.ConsumerGroup => + "consumer_group", + OperationalResourceKind.Topic => + "topic", + OperationalResourceKind.Broker => + "broker", + OperationalResourceKind.Operation => + "operation", + _ => throw new ArgumentOutOfRangeException( + nameof(kind)), + }; +} diff --git a/tests/Kafdeck.Architecture.Tests/V08W64OperationalAnalyticsRuntimeTests.cs b/tests/Kafdeck.Architecture.Tests/V08W64OperationalAnalyticsRuntimeTests.cs new file mode 100644 index 000000000..5c38a96fd --- /dev/null +++ b/tests/Kafdeck.Architecture.Tests/V08W64OperationalAnalyticsRuntimeTests.cs @@ -0,0 +1,1295 @@ +using Kafdeck.Api; +using Kafdeck.Core.Consumers; +using Kafdeck.Core.Ecosystem; +using Kafdeck.Core.Observability; +using Kafdeck.Core.ReadViews; +using Kafdeck.Infrastructure.Configuration; +using Kafdeck.Modules.Consumers; +using Microsoft.Extensions.Logging.Abstractions; +using Xunit; + +namespace Kafdeck.Architecture.Tests; + +public sealed class V08W64OperationalAnalyticsRuntimeTests +{ + [Fact] + public async Task Historical_trend_reads_W63_consumer_lag_series() + { + var now = + DateTimeOffset.UtcNow; + var store = + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + [ + new HistoricalMetricSeries( + new HistoricalMetricIdentity( + OperationalMetricHistoryNames + .ConsumerLagTotal, + "prod", + "consumer_group", + "group-a"), + [ + HistoricalMetricSample.Gauge( + new HistoricalMetricIdentity( + OperationalMetricHistoryNames + .ConsumerLagTotal, + "prod", + "consumer_group", + "group-a"), + now.AddMinutes(-2), + 12, + "test", + "Stable"), + HistoricalMetricSample.Gauge( + new HistoricalMetricIdentity( + OperationalMetricHistoryNames + .ConsumerLagTotal, + "prod", + "consumer_group", + "group-a"), + now.AddMinutes(-1), + 4, + "test", + "Stable"), + ]), + ], + Truncated: false, + LimitReason: null, + now.AddHours(-1), + now, + "test")); + + var service = + new OperationalTrendService( + store); + var query = + new OperationalTrendQuery( + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + now.AddHours(-1), + now, + 100); + + var result = + await service.QueryAsync( + query, + CancellationToken.None); + + Assert.Equal( + OperationalTrendState.Available, + result.State); + Assert.Equal( + 2, + result.Points.Count); + Assert.Equal( + 12, + result.Points[0].Average); + Assert.Equal( + 4, + result.Points[1].Average); + Assert.NotNull( + store.LastQuery); + Assert.Equal( + "consumer_group", + store.LastQuery!.ResourceKind); + } + + [Fact] + public async Task Slo_uses_only_complete_available_history() + { + var now = + DateTimeOffset.UtcNow; + var identity = + new HistoricalMetricIdentity( + OperationalMetricHistoryNames + .ConsumerLagTotal, + "prod", + "consumer_group", + "group-a"); + var store = + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + [ + new HistoricalMetricSeries( + identity, + [ + HistoricalMetricSample.Gauge( + identity, + now.AddMinutes(-2), + 5, + "test", + "Stable"), + HistoricalMetricSample.Gauge( + identity, + now.AddMinutes(-1), + 15, + "test", + "Stable"), + ]), + ], + false, + null, + now.AddHours(-1), + now, + "test")); + var service = + new OperationalTrendService( + store); + + var result = + await service.EvaluateSloAsync( + new OperationalSloDefinition( + "consumer-lag", + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + MaximumGoodValue: 10, + TargetFraction: 0.9), + now.AddMinutes(-3), + now, + 100, + CancellationToken.None); + + Assert.Equal( + 2, + result.EvaluatedPoints); + Assert.Equal( + 1, + result.GoodPoints); + Assert.Equal( + 0.5, + result.ComplianceFraction); + Assert.NotNull( + result.BurnRate); + Assert.Equal( + 5d, + result.BurnRate!.Value, + precision: 10); + Assert.Equal( + OperationalEvidenceState.Available, + result.State); + } + + [Fact] + public async Task Trend_is_explicitly_unavailable_without_history_provider() + { + var now = + DateTimeOffset.UtcNow; + var service = + new OperationalTrendService(); + + var result = + await service.QueryAsync( + new OperationalTrendQuery( + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + now.AddMinutes(-30), + now, + 100), + CancellationToken.None); + + Assert.Equal( + OperationalTrendState.Unavailable, + result.State); + Assert.Empty( + result.Points); + } + + [Fact] + public async Task Live_consumer_lag_uses_Kafka_evidence_and_keeps_rate_unavailable() + { + var consumerPort = + new FakeConsumerGroupReadPort(); + var consumers = + new ConsumerExplorerService( + consumerPort); + var service = + new OperationalAnalyticsRuntimeService( + consumers, + new UnavailableMetricsObservationPort()); + + var query = + new OperationalAnalyticsQuery( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a", + [ + OperationalMetricKind.ConsumerLagTotal, + OperationalMetricKind.ConsumerConsumeRecordsPerSecond, + ], + 8); + + var result = + await service.QueryAsync( + query, + new ReadViewOperationContext( + DateTimeOffset.UtcNow + .AddSeconds(10), + 8, + 256 * 1024), + CancellationToken.None); + + Assert.True( + result.IsSuccess); + var value = + Assert.IsType( + result.Value); + Assert.Collection( + value.Items, + lag => + { + Assert.Equal( + OperationalMetricKind.ConsumerLagTotal, + lag.Metric); + Assert.Equal( + 10d, + lag.Value); + Assert.Equal( + OperationalEvidenceState.Available, + lag.State); + }, + rate => + { + Assert.Equal( + OperationalMetricKind.ConsumerConsumeRecordsPerSecond, + rate.Metric); + Assert.Null( + rate.Value); + Assert.Equal( + OperationalEvidenceState.Unavailable, + rate.State); + }); + } + + [Fact] + public async Task Live_consumer_lag_marks_imprecise_large_integer_as_partial() + { + var consumers = + new ConsumerExplorerService( + new FakeConsumerGroupReadPort( + 9_007_199_254_740_993L)); + var service = + new OperationalAnalyticsRuntimeService( + consumers, + new UnavailableMetricsObservationPort()); + + var result = + await service.QueryAsync( + new OperationalAnalyticsQuery( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a", + [OperationalMetricKind.ConsumerLagTotal], + 4), + new ReadViewOperationContext( + DateTimeOffset.UtcNow.AddSeconds(10), + 4, + 128 * 1024), + CancellationToken.None); + + var evidence = + Assert.Single( + Assert.IsType( + result.Value).Items); + + Assert.Equal( + OperationalEvidenceState.Partial, + evidence.State); + } + + [Fact] + public async Task Live_consumer_rate_marks_stale_provider_evidence_as_stale() + { + var now = + DateTimeOffset.UtcNow; + var consumers = + new ConsumerExplorerService( + new FakeConsumerGroupReadPort()); + var service = + new OperationalAnalyticsRuntimeService( + consumers, + new StaticMetricsObservationPort( + new ConsumerRateObservation( + ProduceRecordsPerSecond: null, + ConsumeRecordsPerSecond: 5, + Window: TimeSpan.FromMinutes(1), + ObservedAt: + now.AddMinutes(-5), + Source: "test")), + new FixedTimeProvider(now)); + + var result = + await service.QueryAsync( + new OperationalAnalyticsQuery( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a", + [OperationalMetricKind.ConsumerConsumeRecordsPerSecond], + 4), + new ReadViewOperationContext( + now.AddSeconds(10), + 4, + 128 * 1024), + CancellationToken.None); + + var evidence = + Assert.Single( + Assert.IsType( + result.Value).Items); + + Assert.Equal( + OperationalEvidenceState.Stale, + evidence.State); + } + + [Fact] + public async Task Trend_with_unknown_points_is_reported_partial() + { + var now = + DateTimeOffset.UtcNow; + var identity = + new HistoricalMetricIdentity( + OperationalMetricHistoryNames.ConsumerLagTotal, + "prod", + "consumer_group", + "group-a"); + var store = + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + [ + new HistoricalMetricSeries( + identity, + [ + new HistoricalMetricSample( + identity, + now.AddMinutes(-1), + 1, + 1, + 1, + 1, + 300, + "legacy", + "Unknown"), + ]), + ], + false, + null, + now.AddHours(-1), + now, + "test")); + var service = + new OperationalTrendService( + store); + + var result = + await service.QueryAsync( + new OperationalTrendQuery( + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + now.AddHours(-1), + now, + 10), + CancellationToken.None); + + Assert.Equal( + OperationalTrendState.Partial, + result.State); + } + + [Fact] + public async Task Consumer_history_rejects_rounded_two_to_the_63_before_long_cast() + { + var now = + DateTimeOffset.UtcNow; + var identity = + new HistoricalMetricIdentity( + OperationalMetricHistoryNames.ConsumerLagTotal, + "prod", + "consumer_group", + "group-a"); + var store = + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + [ + new HistoricalMetricSeries( + identity, + [ + HistoricalMetricSample.Gauge( + identity, + now.AddMinutes(-1), + (double)long.MaxValue, + "test", + "Partial"), + ]), + ], + false, + null, + now.AddHours(-1), + now, + "test")); + var port = + new HistoricalConsumerHistoryObservationPort( + store, + timeProvider: + new FixedTimeProvider(now)); + + var result = + await port.GetConsumerHistoryAsync( + "prod", + "group-a", + new ReadViewOperationContext( + now.AddSeconds(10), + 100, + 128 * 1024), + CancellationToken.None); + + Assert.True( + result.IsSuccess); + var point = + Assert.Single( + result.Value!); + Assert.Null( + point.TotalLag); + Assert.Equal( + "Partial", + point.State); + } + + [Fact] + public async Task Slo_weights_rollups_by_sample_count_and_rejects_mixed_threshold_rollup() + { + var now = + DateTimeOffset.UtcNow; + var identity = + new HistoricalMetricIdentity( + OperationalMetricHistoryNames.ConsumerLagTotal, + "prod", + "consumer_group", + "group-a"); + var goodRollup = + new HistoricalMetricSample( + identity, + now.AddMinutes(-5), + 1, + 5, + 15, + 5, + 300, + "rollup", + "Stable", + now.AddMinutes(-10), + now.AddMinutes(-5)); + var badRaw = + HistoricalMetricSample.Gauge( + identity, + now.AddMinutes(-1), + 20, + "raw", + "Stable"); + var store = + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + [ + new HistoricalMetricSeries( + identity, + [goodRollup, badRaw]), + ], + false, + null, + now.AddMinutes(-10), + now, + "test")); + var service = + new OperationalTrendService( + store); + + var result = + await service.EvaluateSloAsync( + new OperationalSloDefinition( + "weighted", + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + 10, + 0.9), + now.AddMinutes(-10), + now, + 100, + CancellationToken.None); + + Assert.Equal( + 6, + result.EvaluatedPoints); + Assert.Equal( + 5, + result.GoodPoints); + + var mixed = + new HistoricalMetricSample( + identity, + now.AddMinutes(-5), + 1, + 20, + 50, + 5, + 300, + "rollup", + "Stable", + now.AddMinutes(-10), + now.AddMinutes(-5)); + var mixedService = + new OperationalTrendService( + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + [ + new HistoricalMetricSeries( + identity, + [mixed]), + ], + false, + null, + now.AddMinutes(-10), + now, + "test"))); + + var ambiguous = + await mixedService.EvaluateSloAsync( + new OperationalSloDefinition( + "mixed", + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + 10, + 0.9), + now.AddMinutes(-10), + now.AddMinutes(-5), + 100, + CancellationToken.None); + + Assert.Equal( + OperationalEvidenceState.Partial, + ambiguous.State); + Assert.Equal( + "aggregate_threshold_ambiguous", + ambiguous.ReasonCode); + Assert.Null( + ambiguous.ComplianceFraction); + } + + [Fact] + public async Task Trend_rejects_explicit_query_beyond_configured_history_limits() + { + var now = + DateTimeOffset.UtcNow; + var policy = + new HistoricalMetricStorePolicy( + TimeSpan.FromMinutes(30), + MaxSeriesPerQuery: 1, + MaxPointsPerQuery: 50, + MaxQueryDuration: + TimeSpan.FromSeconds(5), + MaxConcurrentQueries: 1); + var service = + new OperationalTrendService( + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + Array.Empty(), + false, + null, + now.AddMinutes(-30), + now, + "test")), + policy); + + Assert.Equal( + 50, + service.DefaultMaxPoints); + + await Assert.ThrowsAsync( + () => service.QueryAsync( + new OperationalTrendQuery( + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + now.AddHours(-1), + now, + 50), + CancellationToken.None)); + + await Assert.ThrowsAsync( + () => service.QueryAsync( + new OperationalTrendQuery( + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + now.AddMinutes(-20), + now, + 51), + CancellationToken.None)); + } + + [Fact] + public async Task Consumer_history_clamps_fixed_query_to_configured_point_limit() + { + var now = + DateTimeOffset.UtcNow; + var store = + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + Array.Empty(), + false, + null, + now.AddMinutes(-30), + now, + "test")); + var policy = + new HistoricalMetricStorePolicy( + TimeSpan.FromMinutes(30), + MaxSeriesPerQuery: 1, + MaxPointsPerQuery: 25, + MaxQueryDuration: + TimeSpan.FromSeconds(5), + MaxConcurrentQueries: 1); + var port = + new HistoricalConsumerHistoryObservationPort( + store, + policy, + new FixedTimeProvider(now)); + + var result = + await port.GetConsumerHistoryAsync( + "prod", + "group-a", + new ReadViewOperationContext( + now.AddSeconds(10), + 500, + 128 * 1024), + CancellationToken.None); + + Assert.True( + result.IsSuccess); + Assert.NotNull( + store.LastQuery); + Assert.Equal( + 25, + store.LastQuery!.MaxPoints); + Assert.Equal( + TimeSpan.FromMinutes(30), + store.LastQuery.ToUtc - + store.LastQuery.FromUtc); + } + + [Fact] + public async Task Truncated_history_is_partial_and_cannot_produce_available_slo() + { + var now = + DateTimeOffset.UtcNow; + var identity = + new HistoricalMetricIdentity( + OperationalMetricHistoryNames.ConsumerLagTotal, + "prod", + "consumer_group", + "group-a"); + var store = + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + [ + new HistoricalMetricSeries( + identity, + [ + HistoricalMetricSample.Gauge( + identity, + now.AddMinutes(-2), + 5, + "test", + "Stable"), + HistoricalMetricSample.Gauge( + identity, + now.AddMinutes(-1), + 6, + "test", + "Stable"), + ]), + ], + Truncated: true, + LimitReason: "max_points", + now.AddMinutes(-3), + now, + "test")); + var service = + new OperationalTrendService(store); + + var trend = + await service.QueryAsync( + new OperationalTrendQuery( + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + now.AddMinutes(-3), + now, + 2), + CancellationToken.None); + + Assert.Equal( + OperationalTrendState.Partial, + trend.State); + + var slo = + await service.EvaluateSloAsync( + new OperationalSloDefinition( + "lag-slo", + OperationalMetricKind.ConsumerLagTotal, + trend.Resource, + 10, + 0.99), + now.AddMinutes(-3), + now, + 2, + CancellationToken.None); + + Assert.Equal( + OperationalEvidenceState.Partial, + slo.State); + } + + [Fact] + public async Task Sparse_history_window_is_partial_slo_evidence() + { + var now = + DateTimeOffset.UtcNow; + var identity = + new HistoricalMetricIdentity( + OperationalMetricHistoryNames.ConsumerLagTotal, + "prod", + "consumer_group", + "group-a"); + var store = + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + [ + new HistoricalMetricSeries( + identity, + [ + HistoricalMetricSample.Gauge( + identity, + now.AddMinutes(-30), + 5, + "test", + "Stable"), + ]), + ], + false, + null, + now.AddHours(-1), + now, + "test")); + var service = + new OperationalTrendService(store); + + var slo = + await service.EvaluateSloAsync( + new OperationalSloDefinition( + "lag-slo", + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + 10, + 0.99), + now.AddHours(-1), + now, + 100, + CancellationToken.None); + + Assert.Equal( + OperationalEvidenceState.Partial, + slo.State); + Assert.Equal( + "partial_slo_evidence", + slo.ReasonCode); + } + + [Fact] + public async Task Trend_preserves_stale_provider_state() + { + var now = + DateTimeOffset.UtcNow; + var identity = + new HistoricalMetricIdentity( + OperationalMetricHistoryNames.ConsumerLagTotal, + "prod", + "consumer_group", + "group-a"); + var store = + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + [ + new HistoricalMetricSeries( + identity, + [ + HistoricalMetricSample.Gauge( + identity, + now.AddMinutes(-1), + 5, + "test", + "Stale"), + ]), + ], + false, + null, + now.AddMinutes(-2), + now, + "test")); + var service = + new OperationalTrendService(store); + + var trend = + await service.QueryAsync( + new OperationalTrendQuery( + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + now.AddMinutes(-2), + now, + 10), + CancellationToken.None); + + Assert.Equal( + OperationalTrendState.Partial, + trend.State); + Assert.Equal( + OperationalEvidenceState.Stale, + Assert.Single(trend.Points).State); + } + + [Fact] + public async Task History_provider_failure_is_explicitly_unavailable() + { + var now = + DateTimeOffset.UtcNow; + var service = + new OperationalTrendService( + new ThrowingHistoricalMetricStore()); + + var trend = + await service.QueryAsync( + new OperationalTrendQuery( + OperationalMetricKind.ConsumerLagTotal, + new OperationalResourceIdentity( + "prod", + OperationalResourceKind.ConsumerGroup, + "group-a"), + now.AddMinutes(-5), + now, + 10), + CancellationToken.None); + + Assert.Equal( + OperationalTrendState.Unavailable, + trend.State); + Assert.Empty( + trend.Points); + } + + [Fact] + public async Task Sampler_uses_fenced_lease_rotates_groups_and_marks_imprecise_lag_partial() + { + var now = + new DateTimeOffset( + 2026, + 9, + 28, + 20, + 0, + 0, + TimeSpan.Zero); + var time = + new MutableTimeProvider(now); + var consumerPort = + new FakeConsumerGroupReadPort( + 9_007_199_254_740_993L, + ["group-a", "group-b", "group-c"]); + var consumers = + new ConsumerExplorerService( + consumerPort); + var store = + new FakeHistoricalMetricStore( + new HistoricalMetricQueryResult( + Array.Empty(), + false, + null, + now.AddMinutes(-5), + now, + "test")); + var leases = + new FakeSamplingLeaseStore(); + var options = + new KafdeckOptions( + new DeploymentOptions( + "http://127.0.0.1:8080", + null), + [ + new ClusterProfile( + "prod", + ["localhost:9092"], + KafkaSecurityProtocol.Plaintext, + null, + null), + ]); + var service = + new ConsumerLagHistorySamplingHostedService( + options, + consumers, + consumerPort, + store, + leases, + new ConsumerLagHistorySamplingPolicy( + TimeSpan.FromMinutes(1), + MaxGroupsPerCluster: 1, + MaxConcurrentGroups: 1), + NullLogger< + ConsumerLagHistorySamplingHostedService>.Instance, + time); + + await service.SampleOnceAsync( + CancellationToken.None); + + time.Advance( + TimeSpan.FromMinutes(1)); + await service.SampleOnceAsync( + CancellationToken.None); + + Assert.Equal( + 2, + consumerPort.RequestedGroups + .Distinct(StringComparer.Ordinal) + .Count()); + Assert.All( + store.AppendedSamples, + sample => + Assert.Equal( + "Partial", + sample.State)); + + leases.DenyNext = true; + var readsBefore = + consumerPort.RequestedGroups.Count; + time.Advance( + TimeSpan.FromMinutes(1)); + await service.SampleOnceAsync( + CancellationToken.None); + + Assert.Equal( + readsBefore, + consumerPort.RequestedGroups.Count); + } + + private sealed class FakeHistoricalMetricStore : + IHistoricalMetricStore + { + private readonly HistoricalMetricQueryResult + _result; + + public FakeHistoricalMetricStore( + HistoricalMetricQueryResult result) + { + _result = result; + } + + public HistoricalMetricQuery? LastQuery { get; private set; } + + public List AppendedSamples { get; } = + new(); + + public Task InitializeAsync( + CancellationToken cancellationToken = default) => + Task.CompletedTask; + + public Task AppendAsync( + IReadOnlyList samples, + CancellationToken cancellationToken = default) + { + AppendedSamples.AddRange( + samples); + return Task.CompletedTask; + } + + public Task QueryAsync( + HistoricalMetricQuery query, + CancellationToken cancellationToken = default) + { + LastQuery = query; + return Task.FromResult( + _result); + } + + public Task DeleteExpiredAsync( + DateTimeOffset rawBeforeUtc, + DateTimeOffset rollupBeforeUtc, + CancellationToken cancellationToken = default) => + Task.FromResult( + new HistoricalMetricRetentionResult( + 0, + 0)); + } + + private sealed class ThrowingHistoricalMetricStore : + IHistoricalMetricStore + { + public Task InitializeAsync( + CancellationToken cancellationToken = default) => + Task.CompletedTask; + + public Task AppendAsync( + IReadOnlyList samples, + CancellationToken cancellationToken = default) => + Task.CompletedTask; + + public Task QueryAsync( + HistoricalMetricQuery query, + CancellationToken cancellationToken = default) => + throw new TimeoutException( + "test timeout"); + + public Task DeleteExpiredAsync( + DateTimeOffset rawBeforeUtc, + DateTimeOffset rollupBeforeUtc, + CancellationToken cancellationToken = default) => + Task.FromResult( + new HistoricalMetricRetentionResult( + 0, + 0)); + } + + private sealed class FakeSamplingLeaseStore : + IHistoricalMetricSamplingLeaseStore + { + public bool DenyNext { get; set; } + + public Task + TryAcquireSamplingLeaseAsync( + string ownerId, + DateTimeOffset nowUtc, + TimeSpan leaseDuration, + CancellationToken cancellationToken = default) + { + if (DenyNext) + { + DenyNext = false; + return Task.FromResult< + HistoricalMetricMaintenanceLease?>( + null); + } + + return Task.FromResult< + HistoricalMetricMaintenanceLease?>( + new HistoricalMetricMaintenanceLease( + ownerId, + 1, + nowUtc.Add(leaseDuration))); + } + } + + private sealed class StaticMetricsObservationPort : + IMetricsObservationPort + { + private readonly ConsumerRateObservation _value; + + public StaticMetricsObservationPort( + ConsumerRateObservation value) + { + _value = value; + } + + public Task> + GetConsumerRateAsync( + string clusterId, + string groupId, + ReadViewOperationContext operation, + CancellationToken cancellationToken) => + Task.FromResult( + ReadViewResult + .Success(_value)); + } + + private sealed class MutableTimeProvider : + TimeProvider + { + private DateTimeOffset _utcNow; + + public MutableTimeProvider( + DateTimeOffset utcNow) + { + _utcNow = utcNow; + } + + public override DateTimeOffset GetUtcNow() => + _utcNow; + + public void Advance( + TimeSpan value) + { + _utcNow = + _utcNow.Add(value); + } + } + + private sealed class FixedTimeProvider : + TimeProvider + { + private readonly DateTimeOffset _utcNow; + + public FixedTimeProvider( + DateTimeOffset utcNow) + { + _utcNow = utcNow; + } + + public override DateTimeOffset GetUtcNow() => + _utcNow; + } + + private sealed class FakeConsumerGroupReadPort : + IConsumerGroupReadPort, + IConsumerGroupSamplingReadPort + { + private readonly long _lag; + private readonly IReadOnlyList _groups; + + public FakeConsumerGroupReadPort( + long lag = 10, + IReadOnlyList? groups = null) + { + _lag = lag; + _groups = + groups ?? + ["group-a"]; + } + + public List RequestedGroups { get; } = + new(); + public Task> + ListGroupPageAsync( + string clusterId, + string? afterGroupId, + int maxItems, + ReadViewOperationContext operation, + CancellationToken cancellationToken) + { + var page = + _groups + .Where(group => + afterGroupId is null || + string.CompareOrdinal( + group, + afterGroupId) > 0) + .OrderBy(group => + group, + StringComparer.Ordinal) + .Take(maxItems + 1) + .ToArray(); + var items = + page + .Take(maxItems) + .Select(group => + new ConsumerGroupSummary( + group, + ConsumerGroupState.Stable, + 1, + false)) + .ToArray(); + var nextCursor = + page.Length > maxItems + ? items[^1].GroupId + : null; + + return Task.FromResult( + ReadViewResult + .Success( + new ConsumerGroupPage( + items, + nextCursor))); + } + + public Task>> + ListGroupsAsync( + string clusterId, + ReadViewOperationContext operation, + CancellationToken cancellationToken) => + Task.FromResult( + ReadViewResult> + .Success( + _groups + .Select(group => + new ConsumerGroupSummary( + group, + ConsumerGroupState.Stable, + 1, + false)) + .ToArray())); + + public Task> + GetGroupAsync( + string clusterId, + string groupId, + ReadViewOperationContext operation, + CancellationToken cancellationToken) => + Task.FromResult( + ReadViewResult.Success( + new ConsumerGroupDetail( + groupId, + ConsumerGroupState.Stable, + null, + null, + null, + Array.Empty()))); + + public Task>> + GetOffsetsAsync( + string clusterId, + string groupId, + ReadViewOperationContext operation, + CancellationToken cancellationToken) + { + RequestedGroups.Add( + groupId); + + return Task.FromResult( + ReadViewResult> + .Success( + [ + new ConsumerOffsetProjection( + "orders", + 0, + 90, + 100, + _lag, + ConsumerOffsetState.Observed), + ])); + } + } +}