diff --git a/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs b/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs new file mode 100644 index 00000000..0c62a9d5 --- /dev/null +++ b/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs @@ -0,0 +1,1381 @@ +using System.Diagnostics; +using System.Text; +using System.Data; +using System.Data.Common; +using System.Globalization; +using Kafdeck.Core.Observability; +using Microsoft.Data.Sqlite; +using SQLitePCL; + +namespace Kafdeck.Infrastructure.Persistence; + +public sealed class AdoHistoricalMetricMaintenanceStore : + IHistoricalMetricMaintenanceStore +{ + private const int SchemaVersion = 1; + private const int SingletonId = 1; + private const string Component = + "historical-metrics-maintenance"; + private const string RollupSource = + "kafdeck-rollup"; + internal const int MaxRawSamplesPerRollupBatch = 32; + internal const int MaxRollupBatchesPerCycle = 24; + + private sealed record RawSampleRow( + HistoricalMetricIdentity Identity, + DateTimeOffset ObservedAtUtc, + double Min, + double Max, + double Sum, + long Count, + string? State, + DateTimeOffset FirstObservedAtUtc, + DateTimeOffset LastObservedAtUtc); + private const long MaintenanceMigrationLockKey = + 4_839_176_502_110_873_342L; + + private readonly IHistoricalMetricsDbConnectionFactory + _connectionFactory; + private readonly TimeProvider _timeProvider; + + public AdoHistoricalMetricMaintenanceStore( + IHistoricalMetricsDbConnectionFactory connectionFactory, + TimeProvider? timeProvider = null) + { + _connectionFactory = + connectionFactory ?? + throw new ArgumentNullException( + nameof(connectionFactory)); + _timeProvider = + timeProvider ?? + TimeProvider.System; + } + + public async Task InitializeAsync( + CancellationToken cancellationToken = default) + { + await using var connection = + await _connectionFactory + .OpenAsync(cancellationToken) + .ConfigureAwait(false); + await using var transaction = + await connection + .BeginTransactionAsync(cancellationToken) + .ConfigureAwait(false); + + if (_connectionFactory.SupportsSelectForUpdate) + { + await using var lockCommand = + connection.CreateCommand(); + lockCommand.Transaction = transaction; + lockCommand.CommandText = + "SELECT pg_advisory_xact_lock(@lock_key)"; + AddParameter( + lockCommand, + "@lock_key", + MaintenanceMigrationLockKey); + await lockCommand + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + + await ExecuteAsync( + connection, + transaction, + """ + CREATE TABLE IF NOT EXISTS kafdeck_schema_info ( + component TEXT PRIMARY KEY, + schema_version INTEGER NOT NULL + ) + """, + cancellationToken) + .ConfigureAwait(false); + + var existingVersion = + await ReadSchemaVersionAsync( + connection, + transaction, + cancellationToken) + .ConfigureAwait(false); + + if (existingVersion is not null && + existingVersion.Value != SchemaVersion) + { + throw new InvalidOperationException( + $"Historical metric maintenance schema version {existingVersion.Value} is unsupported by this binary (expected {SchemaVersion})."); + } + + await ExecuteAsync( + connection, + transaction, + """ + CREATE TABLE IF NOT EXISTS kafdeck_historical_metric_maintenance ( + singleton_id INTEGER PRIMARY KEY, + lease_owner TEXT NULL, + lease_expires_at_utc TEXT NULL, + fencing_token BIGINT NOT NULL, + updated_at_utc TEXT NOT NULL + ) + """, + cancellationToken) + .ConfigureAwait(false); + + await ExecuteAsync( + connection, + transaction, + """ + CREATE TABLE IF NOT EXISTS kafdeck_historical_metric_raw_identities ( + metric_name TEXT NOT NULL, + cluster_id TEXT NOT NULL, + resource_kind TEXT NOT NULL, + resource_id TEXT NOT NULL, + observed_at_utc TEXT NOT NULL, + expires_at_utc TEXT NOT NULL, + PRIMARY KEY ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc) + ) + """, + cancellationToken) + .ConfigureAwait(false); + + await ExecuteAsync( + connection, + transaction, + """ + CREATE INDEX IF NOT EXISTS ix_kafdeck_history_raw_identity_expiry + ON kafdeck_historical_metric_raw_identities ( + expires_at_utc) + """, + cancellationToken) + .ConfigureAwait(false); + + await using (var seed = + connection.CreateCommand()) + { + seed.Transaction = transaction; + seed.CommandText = + """ + INSERT INTO kafdeck_historical_metric_maintenance ( + singleton_id, + lease_owner, + lease_expires_at_utc, + fencing_token, + updated_at_utc) + VALUES ( + @singleton_id, + NULL, + NULL, + 0, + @updated_at_utc) + ON CONFLICT (singleton_id) + DO NOTHING + """; + AddParameter( + seed, + "@singleton_id", + SingletonId); + AddParameter( + seed, + "@updated_at_utc", + DateTimeOffset.UnixEpoch + .ToString("O")); + await seed + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + + if (existingVersion is null) + { + await using var versionInsert = + connection.CreateCommand(); + versionInsert.Transaction = transaction; + versionInsert.CommandText = + """ + INSERT INTO kafdeck_schema_info ( + component, + schema_version) + VALUES ( + @component, + @schema_version) + """; + AddParameter( + versionInsert, + "@component", + Component); + AddParameter( + versionInsert, + "@schema_version", + SchemaVersion); + await versionInsert + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + + await transaction + .CommitAsync(cancellationToken) + .ConfigureAwait(false); + } + + public async Task + TryAcquireLeaseAsync( + string ownerId, + DateTimeOffset nowUtc, + TimeSpan leaseDuration, + CancellationToken cancellationToken = default) + { + ValidateOwner(ownerId); + + if (nowUtc == default) + { + throw new ArgumentOutOfRangeException( + nameof(nowUtc)); + } + + if (leaseDuration <= TimeSpan.Zero || + leaseDuration > + TimeSpan.FromMinutes(10)) + { + throw new ArgumentOutOfRangeException( + nameof(leaseDuration)); + } + + ownerId = ownerId.Trim(); + nowUtc = nowUtc.ToUniversalTime(); + + await using var connection = + await _connectionFactory + .OpenAsync(cancellationToken) + .ConfigureAwait(false); + await using var transaction = + await connection + .BeginTransactionAsync(cancellationToken) + .ConfigureAwait(false); + + var current = + await LockAndReadLeaseAsync( + connection, + transaction, + cancellationToken) + .ConfigureAwait(false); + + if (current.OwnerId is not null && + current.ExpiresAtUtc is not null && + current.ExpiresAtUtc.Value > nowUtc && + !string.Equals( + current.OwnerId, + ownerId, + StringComparison.Ordinal)) + { + await transaction + .CommitAsync(cancellationToken) + .ConfigureAwait(false); + return null; + } + + var nextToken = + checked(current.FencingToken + 1); + var expiresAtUtc = + nowUtc.Add(leaseDuration); + + await using var update = + connection.CreateCommand(); + update.Transaction = transaction; + update.CommandText = + """ + UPDATE kafdeck_historical_metric_maintenance + SET + lease_owner = @lease_owner, + lease_expires_at_utc = @lease_expires_at_utc, + fencing_token = @fencing_token, + updated_at_utc = @updated_at_utc + WHERE singleton_id = @singleton_id + """; + AddParameter( + update, + "@lease_owner", + ownerId); + AddParameter( + update, + "@lease_expires_at_utc", + expiresAtUtc.ToString("O")); + AddParameter( + update, + "@fencing_token", + nextToken); + AddParameter( + update, + "@updated_at_utc", + nowUtc.ToString("O")); + AddParameter( + update, + "@singleton_id", + SingletonId); + + if (await update + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false) != 1) + { + throw new InvalidOperationException( + "Historical metric maintenance lease row is missing."); + } + + await transaction + .CommitAsync(cancellationToken) + .ConfigureAwait(false); + + return new HistoricalMetricMaintenanceLease( + ownerId, + nextToken, + expiresAtUtc); + } + + public async Task + RunCycleAsync( + HistoricalMetricMaintenanceLease lease, + DateTimeOffset nowUtc, + HistoricalMetricMaintenancePolicy policy, + CancellationToken cancellationToken = default) + { + ArgumentNullException.ThrowIfNull(lease); + ArgumentNullException.ThrowIfNull(policy); + policy.Validate(); + ValidateOwner(lease.OwnerId); + + if (lease.FencingToken <= 0) + { + throw new ArgumentOutOfRangeException( + nameof(lease)); + } + + if (nowUtc == default) + { + throw new ArgumentOutOfRangeException( + nameof(nowUtc)); + } + + nowUtc = nowUtc.ToUniversalTime(); + var rawCutoffUtc = + nowUtc.Subtract(policy.RawRetention); + var rollupCutoffUtc = + nowUtc.Subtract(policy.RollupRetention); + + using var deadline = + CancellationTokenSource.CreateLinkedTokenSource( + cancellationToken); + deadline.CancelAfter( + policy.MaxCycleDuration); + var cycleToken = + deadline.Token; + + try + { + await using var connection = + await _connectionFactory + .OpenAsync(cycleToken) + .ConfigureAwait(false); + using var providerCancellation = + RegisterProviderCancellation( + connection, + cycleToken); + + var rolledWindows = 0; + var rollupBatches = 0; + var completedWindows = + new HashSet(); + long rawDeleted = 0; + + while (rolledWindows < + policy.MaxRollupWindowsPerCycle && + rollupBatches < + MaxRollupBatchesPerCycle) + { + await using var rollupTransaction = + await BeginMaintenanceTransactionAsync( + connection, + cycleToken) + .ConfigureAwait(false); + + var current = + await LockAndReadLeaseAsync( + connection, + rollupTransaction, + cycleToken) + .ConfigureAwait(false); + + var leaseCheckUtc = + _timeProvider.GetUtcNow(); + + if (!LeaseIsValid( + current, + lease, + leaseCheckUtc)) + { + await rollupTransaction + .CommitAsync(cycleToken) + .ConfigureAwait(false); + + return new HistoricalMetricMaintenanceResult( + LeaseValid: false, + RolledWindows: rolledWindows, + RawDeleted: rawDeleted, + RollupDeleted: 0); + } + + var oldestRaw = + await ReadOldestEligibleRawTimestampAsync( + connection, + rollupTransaction, + rawCutoffUtc, + cycleToken) + .ConfigureAwait(false); + + if (oldestRaw is null) + { + await rollupTransaction + .CommitAsync(cycleToken) + .ConfigureAwait(false); + break; + } + + var windowStartUtc = + FloorToWindow( + oldestRaw.Value, + policy.RollupResolutionSeconds); + var windowEndUtc = + windowStartUtc.AddSeconds( + policy.RollupResolutionSeconds); + + if (windowEndUtc > rawCutoffUtc) + { + await rollupTransaction + .CommitAsync(cycleToken) + .ConfigureAwait(false); + break; + } + + var processed = + await ProcessRawRollupBatchAsync( + connection, + rollupTransaction, + windowStartUtc, + windowEndUtc, + policy.RollupResolutionSeconds, + MaxRawSamplesPerRollupBatch, + cycleToken) + .ConfigureAwait(false); + + var windowHasRemainingRows = + await HasRawRowsInWindowAsync( + connection, + rollupTransaction, + windowStartUtc, + windowEndUtc, + cycleToken) + .ConfigureAwait(false); + + await rollupTransaction + .CommitAsync(cycleToken) + .ConfigureAwait(false); + + rawDeleted += processed; + rollupBatches++; + + if (!windowHasRemainingRows && + completedWindows.Add( + windowStartUtc)) + { + rolledWindows = + completedWindows.Count; + } + } + + long rollupDeleted = 0; + await using (var retentionTransaction = + await BeginMaintenanceTransactionAsync( + connection, + cycleToken) + .ConfigureAwait(false)) + { + var current = + await LockAndReadLeaseAsync( + connection, + retentionTransaction, + cycleToken) + .ConfigureAwait(false); + + var leaseCheckUtc = + _timeProvider.GetUtcNow(); + + if (!LeaseIsValid( + current, + lease, + leaseCheckUtc)) + { + await retentionTransaction + .CommitAsync(cycleToken) + .ConfigureAwait(false); + + return new HistoricalMetricMaintenanceResult( + LeaseValid: false, + RolledWindows: rolledWindows, + RawDeleted: rawDeleted, + RollupDeleted: 0); + } + + rollupDeleted = + await DeleteExpiredRollupsAsync( + connection, + retentionTransaction, + rollupCutoffUtc, + policy.MaxRollupDeletesPerCycle, + cycleToken) + .ConfigureAwait(false); + + _ = await DeleteExpiredRawIdentitiesAsync( + connection, + retentionTransaction, + _timeProvider.GetUtcNow(), + policy.MaxRollupDeletesPerCycle, + cycleToken) + .ConfigureAwait(false); + + await retentionTransaction + .CommitAsync(cycleToken) + .ConfigureAwait(false); + } + + return new HistoricalMetricMaintenanceResult( + LeaseValid: true, + RolledWindows: rolledWindows, + RawDeleted: rawDeleted, + RollupDeleted: rollupDeleted); + } + catch (SqliteException exception) + when (exception.SqliteErrorCode == + raw.SQLITE_INTERRUPT && + deadline.IsCancellationRequested) + { + if (cancellationToken.IsCancellationRequested) + { + throw new OperationCanceledException( + "Historical metric maintenance was cancelled.", + exception, + cancellationToken); + } + + throw new TimeoutException( + "Historical metric maintenance exceeded the configured cycle duration.", + exception); + } + catch (OperationCanceledException) + when (!cancellationToken.IsCancellationRequested && + deadline.IsCancellationRequested) + { + throw new TimeoutException( + "Historical metric maintenance exceeded the configured cycle duration."); + } + } + + internal ValueTask + BeginMaintenanceTransactionAsync( + DbConnection connection, + CancellationToken cancellationToken) => + connection.BeginTransactionAsync( + _connectionFactory.SupportsSelectForUpdate + ? IsolationLevel.RepeatableRead + : IsolationLevel.Serializable, + cancellationToken); + + private static bool LeaseIsValid( + LeaseRow current, + HistoricalMetricMaintenanceLease lease, + DateTimeOffset nowUtc) => + string.Equals( + current.OwnerId, + lease.OwnerId, + StringComparison.Ordinal) && + current.FencingToken == + lease.FencingToken && + current.ExpiresAtUtc is not null && + current.ExpiresAtUtc.Value > nowUtc; + + private static CancellationTokenRegistration + RegisterProviderCancellation( + DbConnection connection, + CancellationToken cancellationToken) + { + if (connection is not SqliteConnection sqliteConnection || + !cancellationToken.CanBeCanceled) + { + return default; + } + + return cancellationToken.UnsafeRegister( + static state => + { + var sqlite = + (SqliteConnection)state!; + raw.sqlite3_interrupt( + sqlite.Handle); + }, + sqliteConnection); + } + + private async Task LockAndReadLeaseAsync( + DbConnection connection, + DbTransaction transaction, + CancellationToken cancellationToken) + { + if (!_connectionFactory.SupportsSelectForUpdate) + { + await using var lockCommand = + connection.CreateCommand(); + lockCommand.Transaction = transaction; + lockCommand.CommandText = + """ + UPDATE kafdeck_historical_metric_maintenance + SET updated_at_utc = updated_at_utc + WHERE singleton_id = @singleton_id + """; + AddParameter( + lockCommand, + "@singleton_id", + SingletonId); + await lockCommand + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = + _connectionFactory.SupportsSelectForUpdate + ? """ + SELECT + lease_owner, + lease_expires_at_utc, + fencing_token + FROM kafdeck_historical_metric_maintenance + WHERE singleton_id = @singleton_id + FOR UPDATE + """ + : """ + SELECT + lease_owner, + lease_expires_at_utc, + fencing_token + FROM kafdeck_historical_metric_maintenance + WHERE singleton_id = @singleton_id + """; + AddParameter( + command, + "@singleton_id", + SingletonId); + + await using var reader = + await command + .ExecuteReaderAsync(cancellationToken) + .ConfigureAwait(false); + + if (!await reader + .ReadAsync(cancellationToken) + .ConfigureAwait(false)) + { + throw new InvalidOperationException( + "Historical metric maintenance lease row is missing."); + } + + var owner = + reader.IsDBNull(0) + ? null + : reader.GetString(0); + DateTimeOffset? expires = null; + if (!reader.IsDBNull(1)) + { + expires = + DateTimeOffset.Parse( + reader.GetString(1), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind); + } + + return new LeaseRow( + owner, + expires, + Convert.ToInt64( + reader.GetValue(2), + CultureInfo.InvariantCulture)); + } + + private static async Task + ReadOldestEligibleRawTimestampAsync( + DbConnection connection, + DbTransaction transaction, + DateTimeOffset rawCutoffUtc, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = + """ + SELECT observed_at_utc + FROM kafdeck_historical_metric_samples + WHERE resolution_seconds = 0 + AND observed_at_utc < @raw_cutoff_utc + ORDER BY observed_at_utc + LIMIT 1 + """; + AddParameter( + command, + "@raw_cutoff_utc", + rawCutoffUtc.ToString("O")); + + var value = + await command + .ExecuteScalarAsync(cancellationToken) + .ConfigureAwait(false); + + return value is null || + value is DBNull + ? null + : DateTimeOffset.Parse( + Convert.ToString( + value, + CultureInfo.InvariantCulture)!, + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind); + } + + private static async Task HasRawRowsInWindowAsync( + DbConnection connection, + DbTransaction transaction, + DateTimeOffset windowStartUtc, + DateTimeOffset windowEndUtc, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = + """ + SELECT 1 + FROM kafdeck_historical_metric_samples + WHERE resolution_seconds = 0 + AND observed_at_utc >= @window_start_utc + AND observed_at_utc < @window_end_utc + LIMIT 1 + """; + AddParameter( + command, + "@window_start_utc", + windowStartUtc.ToString("O")); + AddParameter( + command, + "@window_end_utc", + windowEndUtc.ToString("O")); + + return await command + .ExecuteScalarAsync(cancellationToken) + .ConfigureAwait(false) is not null; + } + + private static async Task ProcessRawRollupBatchAsync( + DbConnection connection, + DbTransaction transaction, + DateTimeOffset windowStartUtc, + DateTimeOffset windowEndUtc, + int resolutionSeconds, + int maxRawSamples, + CancellationToken cancellationToken) + { + if (maxRawSamples is < 1 or > 10_000) + { + throw new ArgumentOutOfRangeException( + nameof(maxRawSamples)); + } + + var rows = + await ReadRawBatchAsync( + connection, + transaction, + windowStartUtc, + windowEndUtc, + maxRawSamples, + cancellationToken) + .ConfigureAwait(false); + + if (rows.Count == 0) + { + return 0; + } + + foreach (var group in rows.GroupBy( + row => row.Identity)) + { + var min = + group.Min(row => row.Min); + var max = + group.Max(row => row.Max); + var sum = + group.Sum(row => row.Sum); + var count = + group.Sum(row => row.Count); + var states = + group + .Select(row => row.State) + .Distinct(StringComparer.Ordinal) + .ToArray(); + var state = + states.Length == 1 + ? states[0] + : "Partial"; + var firstObservedAtUtc = + group.Min( + row => + row.FirstObservedAtUtc); + var lastObservedAtUtc = + group.Max( + row => + row.LastObservedAtUtc); + + await UpsertRollupAsync( + connection, + transaction, + new HistoricalMetricSample( + group.Key, + windowStartUtc, + min, + max, + sum, + count, + resolutionSeconds, + RollupSource, + state, + firstObservedAtUtc, + lastObservedAtUtc), + cancellationToken) + .ConfigureAwait(false); + } + + var deleted = + await DeleteRawBatchAsync( + connection, + transaction, + rows, + cancellationToken) + .ConfigureAwait(false); + + if (deleted != rows.Count) + { + throw new InvalidOperationException( + "Historical metric raw rollup rows changed during a fenced maintenance batch."); + } + + return rows.Count; + } + + private static async Task> + ReadRawBatchAsync( + DbConnection connection, + DbTransaction transaction, + DateTimeOffset windowStartUtc, + DateTimeOffset windowEndUtc, + int maxRawSamples, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = + """ + SELECT + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + min_value, + max_value, + sum_value, + sample_count, + state, + first_observed_at_utc, + last_observed_at_utc + FROM kafdeck_historical_metric_samples + WHERE resolution_seconds = 0 + AND observed_at_utc >= @window_start_utc + AND observed_at_utc < @window_end_utc + ORDER BY + observed_at_utc, + metric_name, + cluster_id, + resource_kind, + resource_id + LIMIT @raw_limit + """; + AddParameter( + command, + "@window_start_utc", + windowStartUtc.ToString("O")); + AddParameter( + command, + "@window_end_utc", + windowEndUtc.ToString("O")); + AddParameter( + command, + "@raw_limit", + maxRawSamples); + + var rows = + new List( + maxRawSamples); + await using var reader = + await command + .ExecuteReaderAsync(cancellationToken) + .ConfigureAwait(false); + + while (await reader + .ReadAsync(cancellationToken) + .ConfigureAwait(false)) + { + var observedAtUtc = + DateTimeOffset.Parse( + reader.GetString(4), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind); + rows.Add( + new RawSampleRow( + new HistoricalMetricIdentity( + reader.GetString(0), + reader.GetString(1), + reader.GetString(2), + reader.GetString(3)), + observedAtUtc, + Convert.ToDouble( + reader.GetValue(5), + CultureInfo.InvariantCulture), + Convert.ToDouble( + reader.GetValue(6), + CultureInfo.InvariantCulture), + Convert.ToDouble( + reader.GetValue(7), + CultureInfo.InvariantCulture), + Convert.ToInt64( + reader.GetValue(8), + CultureInfo.InvariantCulture), + reader.IsDBNull(9) + ? null + : reader.GetString(9), + reader.IsDBNull(10) + ? observedAtUtc + : DateTimeOffset.Parse( + reader.GetString(10), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind), + reader.IsDBNull(11) + ? observedAtUtc + : DateTimeOffset.Parse( + reader.GetString(11), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind))); + } + + return rows; + } + + private static async Task UpsertRollupAsync( + DbConnection connection, + DbTransaction transaction, + HistoricalMetricSample sample, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = + """ + INSERT INTO kafdeck_historical_metric_samples ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + resolution_seconds, + min_value, + max_value, + sum_value, + sample_count, + source, + state, + first_observed_at_utc, + last_observed_at_utc) + VALUES ( + @metric_name, + @cluster_id, + @resource_kind, + @resource_id, + @observed_at_utc, + @resolution_seconds, + @min_value, + @max_value, + @sum_value, + @sample_count, + @source, + @state, + @first_observed_at_utc, + @last_observed_at_utc) + ON CONFLICT ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + resolution_seconds) + DO UPDATE SET + min_value = CASE + WHEN excluded.min_value < + min_value + THEN excluded.min_value + ELSE min_value + END, + max_value = CASE + WHEN excluded.max_value > + max_value + THEN excluded.max_value + ELSE max_value + END, + sum_value = + sum_value + + excluded.sum_value, + sample_count = + sample_count + + excluded.sample_count, + source = excluded.source, + state = CASE + WHEN state = excluded.state + OR (state IS NULL AND + excluded.state IS NULL) + THEN state + ELSE 'Partial' + END, + first_observed_at_utc = CASE + WHEN first_observed_at_utc IS NULL + THEN NULL + WHEN excluded.first_observed_at_utc < + first_observed_at_utc + THEN excluded.first_observed_at_utc + ELSE first_observed_at_utc + END, + last_observed_at_utc = CASE + WHEN last_observed_at_utc IS NULL + THEN NULL + WHEN excluded.last_observed_at_utc > + last_observed_at_utc + THEN excluded.last_observed_at_utc + ELSE last_observed_at_utc + END + """; + AddParameter(command, "@metric_name", sample.Identity.MetricName); + AddParameter(command, "@cluster_id", sample.Identity.ClusterId); + AddParameter(command, "@resource_kind", sample.Identity.ResourceKind); + AddParameter(command, "@resource_id", sample.Identity.ResourceId); + AddParameter( + command, + "@observed_at_utc", + sample.ObservedAtUtc.ToUniversalTime().ToString("O")); + AddParameter(command, "@resolution_seconds", sample.ResolutionSeconds); + AddParameter(command, "@min_value", sample.Min); + AddParameter(command, "@max_value", sample.Max); + AddParameter(command, "@sum_value", sample.Sum); + AddParameter(command, "@sample_count", sample.Count); + AddParameter(command, "@source", sample.Source); + AddParameter( + command, + "@state", + sample.State is null + ? DBNull.Value + : sample.State); + AddParameter( + command, + "@first_observed_at_utc", + sample.FirstObservedAtUtc!.Value + .ToUniversalTime() + .ToString("O")); + AddParameter( + command, + "@last_observed_at_utc", + sample.LastObservedAtUtc!.Value + .ToUniversalTime() + .ToString("O")); + + await command + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + + private static async Task DeleteRawBatchAsync( + DbConnection connection, + DbTransaction transaction, + IReadOnlyList rows, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + + var predicates = + new List( + rows.Count); + for (var index = 0; + index < rows.Count; + index++) + { + var row = + rows[index]; + var suffix = + index.ToString( + CultureInfo.InvariantCulture); + predicates.Add( + $"(metric_name = @metric_name_{suffix} AND cluster_id = @cluster_id_{suffix} AND resource_kind = @resource_kind_{suffix} AND resource_id = @resource_id_{suffix} AND observed_at_utc = @observed_at_utc_{suffix})"); + AddParameter(command, $"@metric_name_{suffix}", row.Identity.MetricName); + AddParameter(command, $"@cluster_id_{suffix}", row.Identity.ClusterId); + AddParameter(command, $"@resource_kind_{suffix}", row.Identity.ResourceKind); + AddParameter(command, $"@resource_id_{suffix}", row.Identity.ResourceId); + AddParameter( + command, + $"@observed_at_utc_{suffix}", + row.ObservedAtUtc.ToUniversalTime().ToString("O")); + } + + command.CommandText = + $""" + DELETE FROM kafdeck_historical_metric_samples + WHERE resolution_seconds = 0 + AND ({string.Join(" OR ", predicates)}) + """; + + return await command + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + + private static async Task DeleteExpiredRollupsAsync( + DbConnection connection, + DbTransaction transaction, + DateTimeOffset rollupCutoffUtc, + int maxDeletes, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = + """ + DELETE FROM kafdeck_historical_metric_samples + WHERE ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + resolution_seconds) + IN ( + SELECT + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + resolution_seconds + FROM kafdeck_historical_metric_samples + WHERE resolution_seconds > 0 + AND observed_at_utc < @rollup_cutoff_utc + ORDER BY + observed_at_utc, + metric_name, + cluster_id, + resource_kind, + resource_id, + resolution_seconds + LIMIT @delete_limit + ) + """; + AddParameter( + command, + "@rollup_cutoff_utc", + rollupCutoffUtc.ToString("O")); + AddParameter( + command, + "@delete_limit", + maxDeletes); + + return await command + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + + private static async Task DeleteExpiredRawIdentitiesAsync( + DbConnection connection, + DbTransaction transaction, + DateTimeOffset nowUtc, + int maxDeletes, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = + """ + DELETE FROM kafdeck_historical_metric_raw_identities + WHERE ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc) + IN ( + SELECT + identity.metric_name, + identity.cluster_id, + identity.resource_kind, + identity.resource_id, + identity.observed_at_utc + FROM kafdeck_historical_metric_raw_identities identity + WHERE identity.expires_at_utc < @now_utc + AND NOT EXISTS ( + SELECT 1 + FROM kafdeck_historical_metric_samples sample + WHERE sample.metric_name = identity.metric_name + AND sample.cluster_id = identity.cluster_id + AND sample.resource_kind = identity.resource_kind + AND sample.resource_id = identity.resource_id + AND sample.observed_at_utc = identity.observed_at_utc + AND sample.resolution_seconds = 0) + ORDER BY identity.expires_at_utc + LIMIT @delete_limit + ) + """; + AddParameter( + command, + "@now_utc", + nowUtc.ToUniversalTime().ToString("O")); + AddParameter( + command, + "@delete_limit", + maxDeletes); + + return await command + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + + private static DateTimeOffset FloorToWindow( + DateTimeOffset value, + int resolutionSeconds) + { + var unixSeconds = + value.ToUniversalTime() + .ToUnixTimeSeconds(); + var remainder = + unixSeconds % + resolutionSeconds; + if (remainder < 0) + { + remainder += + resolutionSeconds; + } + + return DateTimeOffset + .FromUnixTimeSeconds( + unixSeconds - remainder); + } + + private static async Task ExecuteAsync( + DbConnection connection, + DbTransaction transaction, + string sql, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = sql; + await command + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + + private static async Task ReadSchemaVersionAsync( + DbConnection connection, + DbTransaction transaction, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = + """ + SELECT schema_version + FROM kafdeck_schema_info + WHERE component = @component + """; + AddParameter( + command, + "@component", + Component); + + var value = + await command + .ExecuteScalarAsync(cancellationToken) + .ConfigureAwait(false); + + return value is null || + value is DBNull + ? null + : Convert.ToInt32( + value, + CultureInfo.InvariantCulture); + } + + private static void ValidateOwner( + string ownerId) + { + ArgumentException.ThrowIfNullOrWhiteSpace( + ownerId); + + if (!string.Equals( + ownerId, + ownerId.Trim(), + StringComparison.Ordinal) || + ownerId.Length > 256 || + ownerId.Any(char.IsControl)) + { + throw new ArgumentException( + "Historical metric maintenance owner must be trimmed, contain no control characters, and not exceed 256 characters.", + nameof(ownerId)); + } + } + + private static void AddParameter( + DbCommand command, + string name, + object value) + { + var parameter = + command.CreateParameter(); + parameter.ParameterName = name; + parameter.Value = value; + command.Parameters.Add(parameter); + } + + private sealed record LeaseRow( + string? OwnerId, + DateTimeOffset? ExpiresAtUtc, + long FencingToken); +} diff --git a/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricStore.cs b/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricStore.cs index 40462bcf..2453088b 100644 --- a/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricStore.cs +++ b/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricStore.cs @@ -9,9 +9,18 @@ namespace Kafdeck.Infrastructure.Persistence; public sealed class AdoHistoricalMetricStore : IHistoricalMetricStore { - private const int SchemaVersion = 1; + private const int SchemaVersion = 2; private const string Component = "historical-metrics"; + internal static readonly TimeSpan RawIdentityRetention = + TimeSpan.FromDays(90); + + private readonly record struct RawSampleIdentity( + string MetricName, + string ClusterId, + string ResourceKind, + string ResourceId, + string ObservedAtUtc); private readonly IHistoricalMetricsDbConnectionFactory _connectionFactory; @@ -85,13 +94,14 @@ await ReadSchemaVersionAsync( .ConfigureAwait(false); if (existingVersion is not null && - existingVersion.Value != SchemaVersion) + (existingVersion.Value < 1 || + existingVersion.Value > SchemaVersion)) { throw new InvalidOperationException( - $"Historical metrics schema version {existingVersion.Value} is unsupported by this binary (expected {SchemaVersion})."); + $"Historical metrics schema version {existingVersion.Value} is unsupported by this binary (expected 1..{SchemaVersion})."); } - string[] versionOneStatements = + string[] currentSchemaStatements = [ """ CREATE TABLE IF NOT EXISTS kafdeck_historical_metric_samples ( @@ -107,6 +117,8 @@ CREATE TABLE IF NOT EXISTS kafdeck_historical_metric_samples ( sample_count BIGINT NOT NULL, source TEXT NOT NULL, state TEXT NULL, + first_observed_at_utc TEXT NULL, + last_observed_at_utc TEXT NULL, PRIMARY KEY ( metric_name, cluster_id, @@ -131,9 +143,30 @@ ON kafdeck_historical_metric_samples ( resolution_seconds, observed_at_utc) """, + """ + CREATE TABLE IF NOT EXISTS kafdeck_historical_metric_raw_identities ( + metric_name TEXT NOT NULL, + cluster_id TEXT NOT NULL, + resource_kind TEXT NOT NULL, + resource_id TEXT NOT NULL, + observed_at_utc TEXT NOT NULL, + expires_at_utc TEXT NOT NULL, + PRIMARY KEY ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc) + ) + """, + """ + CREATE INDEX IF NOT EXISTS ix_kafdeck_history_raw_identity_expiry + ON kafdeck_historical_metric_raw_identities ( + expires_at_utc) + """, ]; - foreach (var statement in versionOneStatements) + foreach (var statement in currentSchemaStatements) { await ExecuteInitializationStatementAsync( connection, @@ -143,7 +176,74 @@ await ExecuteInitializationStatementAsync( .ConfigureAwait(false); } - if (existingVersion is null) + if (existingVersion == 1) + { + string[] versionTwoMigrationStatements = + [ + """ + ALTER TABLE kafdeck_historical_metric_samples + ADD COLUMN first_observed_at_utc TEXT NULL + """, + """ + ALTER TABLE kafdeck_historical_metric_samples + ADD COLUMN last_observed_at_utc TEXT NULL + """, + """ + UPDATE kafdeck_historical_metric_samples + SET + first_observed_at_utc = observed_at_utc, + last_observed_at_utc = observed_at_utc + WHERE resolution_seconds = 0 + AND ( + first_observed_at_utc IS NULL + OR last_observed_at_utc IS NULL) + """, + """ + UPDATE kafdeck_historical_metric_samples + SET + first_observed_at_utc = NULL, + last_observed_at_utc = NULL, + state = CASE + WHEN state = 'Partial' + THEN 'Partial' + ELSE 'Unknown' + END + WHERE resolution_seconds > 0 + """, + ]; + + foreach (var statement in versionTwoMigrationStatements) + { + await ExecuteInitializationStatementAsync( + connection, + transaction, + statement, + cancellationToken) + .ConfigureAwait(false); + } + + await using var versionUpdate = + connection.CreateCommand(); + versionUpdate.Transaction = transaction; + versionUpdate.CommandText = + """ + UPDATE kafdeck_schema_info + SET schema_version = @schema_version + WHERE component = @component + """; + AddParameter( + versionUpdate, + "@schema_version", + SchemaVersion); + AddParameter( + versionUpdate, + "@component", + Component); + await versionUpdate + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + else if (existingVersion is null) { await using var versionInsert = connection.CreateCommand(); @@ -170,6 +270,14 @@ await versionInsert .ConfigureAwait(false); } + await BackfillRawIdentitiesAsync( + connection, + transaction, + DateTimeOffset.UtcNow.Add( + RawIdentityRetention), + cancellationToken) + .ConfigureAwait(false); + await transaction .CommitAsync(cancellationToken) .ConfigureAwait(false); @@ -193,6 +301,50 @@ await command .ConfigureAwait(false); } + private static async Task BackfillRawIdentitiesAsync( + DbConnection connection, + DbTransaction transaction, + DateTimeOffset expiresAtUtc, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = + """ + INSERT INTO kafdeck_historical_metric_raw_identities ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + expires_at_utc) + SELECT + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + @expires_at_utc + FROM kafdeck_historical_metric_samples + WHERE resolution_seconds = 0 + ON CONFLICT ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc) + DO NOTHING + """; + AddParameter( + command, + "@expires_at_utc", + expiresAtUtc.ToUniversalTime().ToString("O")); + await command + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false); + } + private static async Task ReadSchemaVersionAsync( DbConnection connection, DbTransaction transaction, @@ -254,8 +406,25 @@ await connection cancellationToken) .ConfigureAwait(false); + var reservedRawIdentities = + await ReserveRawIdentitiesAsync( + connection, + transaction, + samples, + DateTimeOffset.UtcNow.Add( + RawIdentityRetention), + cancellationToken) + .ConfigureAwait(false); + foreach (var sample in samples) { + if (sample.ResolutionSeconds == 0 && + !reservedRawIdentities.Contains( + ToRawIdentity(sample))) + { + continue; + } + await using var command = connection.CreateCommand(); command.Transaction = @@ -274,7 +443,9 @@ INSERT INTO kafdeck_historical_metric_samples ( sum_value, sample_count, source, - state) + state, + first_observed_at_utc, + last_observed_at_utc) VALUES ( @metric_name, @cluster_id, @@ -287,7 +458,9 @@ INSERT INTO kafdeck_historical_metric_samples ( @sum_value, @sample_count, @source, - @state) + @state, + @first_observed_at_utc, + @last_observed_at_utc) ON CONFLICT ( metric_name, cluster_id, @@ -313,6 +486,132 @@ await transaction .ConfigureAwait(false); } + private static async Task> + ReserveRawIdentitiesAsync( + DbConnection connection, + DbTransaction transaction, + IReadOnlyList samples, + DateTimeOffset expiresAtUtc, + CancellationToken cancellationToken) + { + var rawSamples = + samples + .Where(sample => + sample.ResolutionSeconds == 0) + .ToArray(); + if (rawSamples.Length == 0) + { + return []; + } + + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + + var values = + new List( + rawSamples.Length); + for (var index = 0; + index < rawSamples.Length; + index++) + { + var sample = + rawSamples[index]; + var suffix = + index.ToString( + CultureInfo.InvariantCulture); + values.Add( + $"(@metric_name_{suffix}, @cluster_id_{suffix}, @resource_kind_{suffix}, @resource_id_{suffix}, @observed_at_utc_{suffix}, @expires_at_utc)"); + + AddParameter( + command, + $"@metric_name_{suffix}", + sample.Identity.MetricName); + AddParameter( + command, + $"@cluster_id_{suffix}", + sample.Identity.ClusterId); + AddParameter( + command, + $"@resource_kind_{suffix}", + sample.Identity.ResourceKind); + AddParameter( + command, + $"@resource_id_{suffix}", + sample.Identity.ResourceId); + AddParameter( + command, + $"@observed_at_utc_{suffix}", + sample.ObservedAtUtc + .ToUniversalTime() + .ToString("O")); + } + + AddParameter( + command, + "@expires_at_utc", + expiresAtUtc + .ToUniversalTime() + .ToString("O")); + + command.CommandText = + $""" + INSERT INTO kafdeck_historical_metric_raw_identities ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + expires_at_utc) + VALUES {string.Join(", ", values)} + ON CONFLICT ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc) + DO NOTHING + RETURNING + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc + """; + + var reserved = + new HashSet(); + await using var reader = + await command + .ExecuteReaderAsync(cancellationToken) + .ConfigureAwait(false); + while (await reader + .ReadAsync(cancellationToken) + .ConfigureAwait(false)) + { + reserved.Add( + new RawSampleIdentity( + reader.GetString(0), + reader.GetString(1), + reader.GetString(2), + reader.GetString(3), + reader.GetString(4))); + } + + return reserved; + } + + private static RawSampleIdentity ToRawIdentity( + HistoricalMetricSample sample) => + new( + sample.Identity.MetricName, + sample.Identity.ClusterId, + sample.Identity.ResourceKind, + sample.Identity.ResourceId, + sample.ObservedAtUtc + .ToUniversalTime() + .ToString("O")); + public async Task QueryAsync( HistoricalMetricQuery query, @@ -363,7 +662,9 @@ query.ResourceId is null sum_value, sample_count, source, - state + state, + first_observed_at_utc, + last_observed_at_utc FROM kafdeck_historical_metric_samples WHERE metric_name = @metric_name AND cluster_id = @cluster_id @@ -386,7 +687,9 @@ LIMIT @row_limit sum_value, sample_count, source, - state + state, + first_observed_at_utc, + last_observed_at_utc FROM kafdeck_historical_metric_samples WHERE metric_name = @metric_name AND cluster_id = @cluster_id @@ -519,7 +822,33 @@ await command reader.GetString(7), reader.IsDBNull(8) ? null - : reader.GetString(8))); + : reader.GetString(8), + reader.IsDBNull(9) + ? (Convert.ToInt32( + reader.GetValue(2), + CultureInfo.InvariantCulture) == 0 + ? DateTimeOffset.Parse( + reader.GetString(1), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind) + : null) + : DateTimeOffset.Parse( + reader.GetString(9), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind), + reader.IsDBNull(10) + ? (Convert.ToInt32( + reader.GetValue(2), + CultureInfo.InvariantCulture) == 0 + ? DateTimeOffset.Parse( + reader.GetString(1), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind) + : null) + : DateTimeOffset.Parse( + reader.GetString(10), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind))); totalPoints++; } @@ -740,6 +1069,22 @@ private static void BindSample( sample.State is null ? DBNull.Value : sample.State); + AddParameter( + command, + "@first_observed_at_utc", + sample.EffectiveFirstObservedAtUtc is { } firstObservedAtUtc + ? firstObservedAtUtc + .ToUniversalTime() + .ToString("O") + : DBNull.Value); + AddParameter( + command, + "@last_observed_at_utc", + sample.EffectiveLastObservedAtUtc is { } lastObservedAtUtc + ? lastObservedAtUtc + .ToUniversalTime() + .ToString("O") + : DBNull.Value); } private static string ProviderName( diff --git a/src/backend/Kafdeck.Api/HistoricalMetricMaintenanceHostedService.cs b/src/backend/Kafdeck.Api/HistoricalMetricMaintenanceHostedService.cs new file mode 100644 index 00000000..4c4b1c04 --- /dev/null +++ b/src/backend/Kafdeck.Api/HistoricalMetricMaintenanceHostedService.cs @@ -0,0 +1,138 @@ +using Kafdeck.Core.Observability; + +namespace Kafdeck.Api; + +public sealed class HistoricalMetricMaintenanceHostedService : + BackgroundService +{ + private readonly IHistoricalMetricMaintenanceStore + _maintenanceStore; + private readonly HistoricalMetricMaintenancePolicy + _policy; + private readonly ILogger< + HistoricalMetricMaintenanceHostedService> _logger; + private readonly TimeProvider _timeProvider; + private readonly string _ownerId; + + public HistoricalMetricMaintenanceHostedService( + IHistoricalMetricMaintenanceStore maintenanceStore, + HistoricalMetricMaintenancePolicy policy, + ILogger logger, + TimeProvider? timeProvider = null) + { + _maintenanceStore = + maintenanceStore ?? + throw new ArgumentNullException( + nameof(maintenanceStore)); + _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 RunOneCycleAsync( + stoppingToken) + .ConfigureAwait(false); + + using var timer = + new PeriodicTimer( + _policy.CycleInterval, + _timeProvider); + + while (await timer + .WaitForNextTickAsync(stoppingToken) + .ConfigureAwait(false)) + { + await RunOneCycleAsync( + stoppingToken) + .ConfigureAwait(false); + } + } + + private async Task RunOneCycleAsync( + CancellationToken cancellationToken) + { + try + { + var nowUtc = + _timeProvider.GetUtcNow(); + var lease = + await _maintenanceStore + .TryAcquireLeaseAsync( + _ownerId, + nowUtc, + _policy.LeaseDuration, + cancellationToken) + .ConfigureAwait(false); + + if (lease is null) + { + return; + } + + var result = + await _maintenanceStore + .RunCycleAsync( + lease, + nowUtc, + _policy, + cancellationToken) + .ConfigureAwait(false); + + if (!result.LeaseValid) + { + _logger.LogWarning( + "Historical metrics maintenance lease was fenced before execution."); + return; + } + + if (result.RolledWindows > 0 || + result.RawDeleted > 0 || + result.RollupDeleted > 0) + { + _logger.LogInformation( + "Historical metrics maintenance completed: {RolledWindows} rollup windows, {RawDeleted} raw samples deleted, {RollupDeleted} expired rollups deleted.", + result.RolledWindows, + result.RawDeleted, + result.RollupDeleted); + } + } + catch (OperationCanceledException) + when (cancellationToken.IsCancellationRequested) + { + throw; + } + catch (Exception exception) + { + _logger.LogWarning( + exception, + "Historical metrics maintenance cycle failed; retained data remains authoritative and the next bounded cycle will retry."); + } + } + + 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/Program.cs b/src/backend/Kafdeck.Api/Program.cs index 4dc51a6d..fa677221 100644 --- a/src/backend/Kafdeck.Api/Program.cs +++ b/src/backend/Kafdeck.Api/Program.cs @@ -186,6 +186,35 @@ historicalMetricsOptions.ConnectionString is not null IHistoricalMetricsDbConnectionFactory>(), services.GetRequiredService< HistoricalMetricStorePolicy>())); + + builder.Services.AddSingleton( + new HistoricalMetricMaintenancePolicy( + TimeSpan.FromHours( + historicalMetricsOptions.RawRetentionHours), + TimeSpan.FromDays( + historicalMetricsOptions.RollupRetentionDays), + HistoricalMetricMaintenancePolicy + .DefaultRollupResolutionSeconds, + HistoricalMetricMaintenancePolicy + .DefaultMaxRollupWindowsPerCycle, + HistoricalMetricMaintenancePolicy + .DefaultMaxRollupDeletesPerCycle, + HistoricalMetricMaintenancePolicy + .DefaultLeaseDuration, + HistoricalMetricMaintenancePolicy + .DefaultCycleInterval, + HistoricalMetricMaintenancePolicy + .DefaultMaxCycleDuration)); + + builder.Services.AddSingleton< + IHistoricalMetricMaintenanceStore>( + services => + new AdoHistoricalMetricMaintenanceStore( + services.GetRequiredService< + IHistoricalMetricsDbConnectionFactory>())); + + builder.Services.AddHostedService< + HistoricalMetricMaintenanceHostedService>(); } builder.Services.AddSingleton(); @@ -447,6 +476,13 @@ await historicalMetricStore .InitializeAsync() .ConfigureAwait(false); + var historicalMetricMaintenanceStore = + app.Services.GetRequiredService< + IHistoricalMetricMaintenanceStore>(); + await historicalMetricMaintenanceStore + .InitializeAsync() + .ConfigureAwait(false); + app.Logger.LogInformation( "Kafdeck historical metrics provider initialized with provider {Provider}, execution mode {ExecutionMode}, raw retention {RawRetentionHours}h and rollup retention {RollupRetentionDays}d.", historicalMetricsOptions.Provider, diff --git a/src/backend/Kafdeck.Core/Observability/HistoricalMetricMaintenance.cs b/src/backend/Kafdeck.Core/Observability/HistoricalMetricMaintenance.cs new file mode 100644 index 00000000..c831a479 --- /dev/null +++ b/src/backend/Kafdeck.Core/Observability/HistoricalMetricMaintenance.cs @@ -0,0 +1,118 @@ +namespace Kafdeck.Core.Observability; + +public sealed record HistoricalMetricMaintenancePolicy( + TimeSpan RawRetention, + TimeSpan RollupRetention, + int RollupResolutionSeconds, + int MaxRollupWindowsPerCycle, + int MaxRollupDeletesPerCycle, + TimeSpan LeaseDuration, + TimeSpan CycleInterval, + TimeSpan MaxCycleDuration) +{ + public const int DefaultRollupResolutionSeconds = 300; + public const int DefaultMaxRollupWindowsPerCycle = 24; + public const int HardMaxRollupWindowsPerCycle = 288; + public const int DefaultMaxRollupDeletesPerCycle = 1_000; + public const int HardMaxRollupDeletesPerCycle = 10_000; + + public static readonly TimeSpan DefaultLeaseDuration = + TimeSpan.FromMinutes(2); + public static readonly TimeSpan DefaultCycleInterval = + TimeSpan.FromMinutes(1); + public static readonly TimeSpan DefaultMaxCycleDuration = + TimeSpan.FromSeconds(30); + + public void Validate() + { + if (RawRetention <= TimeSpan.Zero || + RawRetention > + TimeSpan.FromDays(7)) + { + throw new ArgumentOutOfRangeException( + nameof(RawRetention)); + } + + if (RollupRetention <= TimeSpan.Zero || + RollupRetention > + TimeSpan.FromDays(90)) + { + throw new ArgumentOutOfRangeException( + nameof(RollupRetention)); + } + + if (RollupResolutionSeconds is < 60 or > 86_400 || + 86_400 % RollupResolutionSeconds != 0) + { + throw new ArgumentOutOfRangeException( + nameof(RollupResolutionSeconds)); + } + + if (MaxRollupWindowsPerCycle is < 1 or > + HardMaxRollupWindowsPerCycle) + { + throw new ArgumentOutOfRangeException( + nameof(MaxRollupWindowsPerCycle)); + } + + if (MaxRollupDeletesPerCycle is < 1 or > + HardMaxRollupDeletesPerCycle) + { + throw new ArgumentOutOfRangeException( + nameof(MaxRollupDeletesPerCycle)); + } + + if (LeaseDuration <= TimeSpan.Zero || + LeaseDuration > + TimeSpan.FromMinutes(10)) + { + throw new ArgumentOutOfRangeException( + nameof(LeaseDuration)); + } + + if (CycleInterval <= TimeSpan.Zero || + CycleInterval > + TimeSpan.FromMinutes(10)) + { + throw new ArgumentOutOfRangeException( + nameof(CycleInterval)); + } + + if (MaxCycleDuration <= TimeSpan.Zero || + MaxCycleDuration > + TimeSpan.FromMinutes(1)) + { + throw new ArgumentOutOfRangeException( + nameof(MaxCycleDuration)); + } + } +} + +public sealed record HistoricalMetricMaintenanceLease( + string OwnerId, + long FencingToken, + DateTimeOffset ExpiresAtUtc); + +public sealed record HistoricalMetricMaintenanceResult( + bool LeaseValid, + int RolledWindows, + long RawDeleted, + long RollupDeleted); + +public interface IHistoricalMetricMaintenanceStore +{ + Task InitializeAsync( + CancellationToken cancellationToken = default); + + Task TryAcquireLeaseAsync( + string ownerId, + DateTimeOffset nowUtc, + TimeSpan leaseDuration, + CancellationToken cancellationToken = default); + + Task RunCycleAsync( + HistoricalMetricMaintenanceLease lease, + DateTimeOffset nowUtc, + HistoricalMetricMaintenancePolicy policy, + CancellationToken cancellationToken = default); +} diff --git a/src/backend/Kafdeck.Core/Observability/HistoricalMetrics.cs b/src/backend/Kafdeck.Core/Observability/HistoricalMetrics.cs index 85dfb269..ec461cb1 100644 --- a/src/backend/Kafdeck.Core/Observability/HistoricalMetrics.cs +++ b/src/backend/Kafdeck.Core/Observability/HistoricalMetrics.cs @@ -75,7 +75,9 @@ public sealed record HistoricalMetricSample( long Count, int ResolutionSeconds, string Source, - string? State = null) + string? State = null, + DateTimeOffset? FirstObservedAtUtc = null, + DateTimeOffset? LastObservedAtUtc = null) { public const int MaxSourceLength = 128; public const int MaxStateLength = 64; @@ -83,6 +85,22 @@ public sealed record HistoricalMetricSample( public double Average => Count == 0 ? 0 : Sum / Count; + public DateTimeOffset? EffectiveFirstObservedAtUtc => + FirstObservedAtUtc ?? + (ResolutionSeconds == 0 + ? ObservedAtUtc + : null); + + public DateTimeOffset? EffectiveLastObservedAtUtc => + LastObservedAtUtc ?? + (ResolutionSeconds == 0 + ? ObservedAtUtc + : null); + + public bool HasKnownCoverage => + EffectiveFirstObservedAtUtc is not null && + EffectiveLastObservedAtUtc is not null; + public static HistoricalMetricSample Gauge( HistoricalMetricIdentity identity, DateTimeOffset observedAtUtc, @@ -139,6 +157,57 @@ public void Validate() MaxStateLength, nameof(State)); } + + if ((FirstObservedAtUtc is null) != + (LastObservedAtUtc is null)) + { + throw new ArgumentException( + "Explicit historical metric observation coverage must provide both first/last timestamps or neither."); + } + + var firstObservedAtUtc = + EffectiveFirstObservedAtUtc; + var lastObservedAtUtc = + EffectiveLastObservedAtUtc; + + if ((firstObservedAtUtc is null) != + (lastObservedAtUtc is null)) + { + throw new ArgumentException( + "Historical metric observation coverage must provide both first/last timestamps or neither."); + } + + if (firstObservedAtUtc is not null && + (firstObservedAtUtc.Value == default || + lastObservedAtUtc!.Value == default || + lastObservedAtUtc.Value < + firstObservedAtUtc.Value)) + { + throw new ArgumentException( + "Historical metric observation coverage must contain a valid ordered first/last timestamp."); + } + + if (ResolutionSeconds == 0 && + firstObservedAtUtc is null) + { + throw new ArgumentException( + "Raw historical metric samples require known observation coverage."); + } + + if (ResolutionSeconds > 0 && + firstObservedAtUtc is null && + !string.Equals( + State, + "Unknown", + StringComparison.Ordinal) && + !string.Equals( + State, + "Partial", + StringComparison.Ordinal)) + { + throw new ArgumentException( + "Historical rollups with unknown coverage must carry an explicit Unknown or Partial state."); + } } private static void ValidateSafeText( diff --git a/tests/Kafdeck.Architecture.Tests/V08W63HistoricalMetricsPersistenceTests.cs b/tests/Kafdeck.Architecture.Tests/V08W63HistoricalMetricsPersistenceTests.cs index 21b4fc96..c12ab210 100644 --- a/tests/Kafdeck.Architecture.Tests/V08W63HistoricalMetricsPersistenceTests.cs +++ b/tests/Kafdeck.Architecture.Tests/V08W63HistoricalMetricsPersistenceTests.cs @@ -1,3 +1,4 @@ +using System.Data; using System.Data.Common; using Kafdeck.Core.Observability; using Kafdeck.Infrastructure.Configuration; @@ -129,7 +130,7 @@ INSERT INTO kafdeck_schema_info ( schema_version) VALUES ( 'historical-metrics', - 2); + 3); """; await command.ExecuteNonQueryAsync(); } @@ -167,6 +168,172 @@ FROM sqlite_master } } + [Fact] + public async Task Version_one_rollup_migration_preserves_unknown_coverage() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-v1-rollup-{Guid.NewGuid():N}.db"); + + try + { + var rawObserved = + new DateTimeOffset( + 2026, + 9, + 28, + 9, + 1, + 0, + TimeSpan.Zero); + var rollupBucket = + new DateTimeOffset( + 2026, + 9, + 28, + 8, + 0, + 0, + TimeSpan.Zero); + + await using (var connection = + new SqliteConnection( + $"Data Source={path}")) + { + await connection.OpenAsync(); + await using var command = + connection.CreateCommand(); + command.CommandText = + """ + CREATE TABLE kafdeck_schema_info ( + component TEXT PRIMARY KEY, + schema_version INTEGER NOT NULL + ); + INSERT INTO kafdeck_schema_info ( + component, + schema_version) + VALUES ( + 'historical-metrics', + 1); + + CREATE TABLE kafdeck_historical_metric_samples ( + metric_name TEXT NOT NULL, + cluster_id TEXT NOT NULL, + resource_kind TEXT NOT NULL, + resource_id TEXT NOT NULL, + observed_at_utc TEXT NOT NULL, + resolution_seconds INTEGER NOT NULL, + min_value REAL NOT NULL, + max_value REAL NOT NULL, + sum_value REAL NOT NULL, + sample_count BIGINT NOT NULL, + source TEXT NOT NULL, + state TEXT NULL, + PRIMARY KEY ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + resolution_seconds) + ); + + INSERT INTO kafdeck_historical_metric_samples + VALUES ( + 'consumer.lag.total', + 'prod', + 'consumer_group', + 'group-a', + @raw_observed, + 0, + 10, + 10, + 10, + 1, + 'legacy', + 'Stable'); + + INSERT INTO kafdeck_historical_metric_samples + VALUES ( + 'consumer.lag.total', + 'prod', + 'consumer_group', + 'group-a', + @rollup_bucket, + 300, + 5, + 20, + 45, + 3, + 'legacy-rollup', + 'Stable'); + """; + command.Parameters.AddWithValue( + "@raw_observed", + rawObserved.ToString("O")); + command.Parameters.AddWithValue( + "@rollup_bucket", + rollupBucket.ToString("O")); + await command.ExecuteNonQueryAsync(); + } + + var store = + new AdoHistoricalMetricStore( + new SqliteHistoricalMetricsDbConnectionFactory( + path), + TestPolicy()); + + await store.InitializeAsync(); + + var query = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + rollupBucket.AddHours(-1), + rawObserved.AddHours(1), + MaxSeries: 1, + MaxPoints: 10)); + + var points = + Assert.Single(query.Series).Points; + + var raw = + Assert.Single( + points, + point => + point.ResolutionSeconds == 0); + Assert.True(raw.HasKnownCoverage); + Assert.Equal( + rawObserved, + raw.FirstObservedAtUtc); + Assert.Equal( + rawObserved, + raw.LastObservedAtUtc); + + var rollup = + Assert.Single( + points, + point => + point.ResolutionSeconds == 300); + Assert.False( + rollup.HasKnownCoverage); + Assert.Null( + rollup.FirstObservedAtUtc); + Assert.Null( + rollup.LastObservedAtUtc); + Assert.Equal( + "Unknown", + rollup.State); + } + finally + { + DeleteSqliteFiles(path); + } + } + [Fact] public async Task PostgreSql_initialization_is_serialized_across_replicas() { @@ -260,6 +427,115 @@ FROM information_schema.tables } } + [Fact] + public async Task Version_one_partial_rollup_migration_preserves_partial_state() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-v1-partial-{Guid.NewGuid():N}.db"); + + try + { + var bucket = + DateTimeOffset.UtcNow.AddHours(-2); + + await using (var connection = + new SqliteConnection( + $"Data Source={path}")) + { + await connection.OpenAsync(); + await using var command = + connection.CreateCommand(); + command.CommandText = + """ + CREATE TABLE kafdeck_schema_info ( + component TEXT PRIMARY KEY, + schema_version INTEGER NOT NULL + ); + INSERT INTO kafdeck_schema_info ( + component, + schema_version) + VALUES ( + 'historical-metrics', + 1); + + CREATE TABLE kafdeck_historical_metric_samples ( + metric_name TEXT NOT NULL, + cluster_id TEXT NOT NULL, + resource_kind TEXT NOT NULL, + resource_id TEXT NOT NULL, + observed_at_utc TEXT NOT NULL, + resolution_seconds INTEGER NOT NULL, + min_value REAL NOT NULL, + max_value REAL NOT NULL, + sum_value REAL NOT NULL, + sample_count BIGINT NOT NULL, + source TEXT NOT NULL, + state TEXT NULL, + PRIMARY KEY ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + resolution_seconds) + ); + + INSERT INTO kafdeck_historical_metric_samples + VALUES ( + 'consumer.lag.total', + 'prod', + 'consumer_group', + 'group-a', + @bucket, + 300, + 5, + 20, + 45, + 3, + 'legacy-rollup', + 'Partial'); + """; + command.Parameters.AddWithValue( + "@bucket", + bucket.ToString("O")); + await command.ExecuteNonQueryAsync(); + } + + var store = + new AdoHistoricalMetricStore( + new SqliteHistoricalMetricsDbConnectionFactory( + path), + TestPolicy()); + await store.InitializeAsync(); + + var query = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + bucket.AddHours(-1), + bucket.AddHours(1), + MaxSeries: 1, + MaxPoints: 10)); + + var rollup = + Assert.Single( + Assert.Single(query.Series).Points); + Assert.Equal( + "Partial", + rollup.State); + Assert.False( + rollup.HasKnownCoverage); + } + finally + { + DeleteSqliteFiles(path); + } + } + [Fact] public async Task Sqlite_store_append_is_idempotent_and_query_is_bounded() { @@ -344,6 +620,33 @@ public async Task Sqlite_store_append_is_idempotent_and_query_is_bounded() } } + [Fact] + public void Raw_sample_rejects_one_sided_explicit_coverage() + { + var observedAt = + DateTimeOffset.UtcNow; + var sample = + new HistoricalMetricSample( + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"), + observedAt, + 10, + 10, + 10, + 1, + ResolutionSeconds: 0, + Source: "consumer_observer", + State: "Stable", + FirstObservedAtUtc: observedAt, + LastObservedAtUtc: null); + + Assert.Throws( + sample.Validate); + } + [Fact] public async Task Sqlite_store_applies_raw_and_rollup_retention_separately() { @@ -386,7 +689,7 @@ public async Task Sqlite_store_applies_raw_and_rollup_retention_separately() 4, ResolutionSeconds: 300, Source: "rollup", - State: "Stable"); + State: "Unknown"); var newRollup = new HistoricalMetricSample( identity, now.AddDays(-1), @@ -396,7 +699,7 @@ public async Task Sqlite_store_applies_raw_and_rollup_retention_separately() 4, ResolutionSeconds: 300, Source: "rollup", - State: "Stable"); + State: "Unknown"); await store.AppendAsync( [oldRaw, newRaw, oldRollup, newRollup]); @@ -511,7 +814,9 @@ WHERE value < 1000000000 CAST(SUM(value) AS REAL) AS sum_value, COUNT(*) AS sample_count, 'test' AS source, - 'Stable' AS state + 'Stable' AS state, + '2026-01-01T00:00:00.0000000+00:00' AS first_observed_at_utc, + '2026-01-01T00:00:00.0000000+00:00' AS last_observed_at_utc FROM counter """; await schema.ExecuteNonQueryAsync(); @@ -560,6 +865,1218 @@ await Assert.ThrowsAsync( CancellationToken.None)); } + [Fact] + public async Task Maintenance_rolls_expired_raw_before_deletion() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-maint-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await store.InitializeAsync(); + await maintenance.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var rollupWindowStart = + DateTimeOffset.FromUnixTimeSeconds( + now.AddHours(-3) + .ToUnixTimeSeconds() / + 300 * + 300); + var identity = + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"); + + await store.AppendAsync( + [ + HistoricalMetricSample.Gauge( + identity, + rollupWindowStart + .AddMinutes(1), + 10, + "consumer_observer", + "Stable"), + HistoricalMetricSample.Gauge( + identity, + rollupWindowStart + .AddMinutes(2), + 30, + "consumer_observer", + "Stable"), + HistoricalMetricSample.Gauge( + identity, + now.AddMinutes(-30), + 40, + "consumer_observer", + "Stable"), + new HistoricalMetricSample( + identity, + now.AddDays(-8), + 1, + 1, + 1, + 1, + ResolutionSeconds: 300, + Source: "rollup", + State: "Unknown"), + ]); + + var policy = + TestMaintenancePolicy(); + var lease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + + var result = + await maintenance.RunCycleAsync( + lease, + now, + policy); + + Assert.True(result.LeaseValid); + Assert.Equal(1, result.RolledWindows); + Assert.Equal(2, result.RawDeleted); + Assert.Equal(1, result.RollupDeleted); + + var query = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + now.AddHours(-4), + now.AddMinutes(1), + MaxSeries: 1, + MaxPoints: 10)); + + var points = + Assert.Single(query.Series) + .Points; + Assert.Equal(2, points.Count); + + var rollup = + Assert.Single( + points, + point => + point.ResolutionSeconds == + 300); + Assert.Equal(10, rollup.Min); + Assert.Equal(30, rollup.Max); + Assert.Equal(40, rollup.Sum); + Assert.Equal(2, rollup.Count); + Assert.Equal(20, rollup.Average); + + var recent = + Assert.Single( + points, + point => + point.ResolutionSeconds == + 0); + Assert.Equal(40, recent.Sum); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_merges_late_raw_into_existing_rollup() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-late-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await store.InitializeAsync(); + await maintenance.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var rollupWindowStart = + DateTimeOffset.FromUnixTimeSeconds( + now.AddHours(-3) + .ToUnixTimeSeconds() / + 300 * + 300); + var identity = + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"); + var policy = + TestMaintenancePolicy(); + + await store.AppendAsync( + [ + HistoricalMetricSample.Gauge( + identity, + rollupWindowStart + .AddMinutes(1), + 10, + "consumer_observer"), + HistoricalMetricSample.Gauge( + identity, + rollupWindowStart + .AddMinutes(2), + 30, + "consumer_observer"), + ]); + + var firstLease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + _ = await maintenance.RunCycleAsync( + firstLease, + now, + policy); + + await store.AppendAsync( + [ + HistoricalMetricSample.Gauge( + identity, + rollupWindowStart + .AddMinutes(3), + 50, + "consumer_observer"), + ]); + + var secondNow = + now.AddSeconds(30); + var secondLease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + secondNow, + policy.LeaseDuration)); + _ = await maintenance.RunCycleAsync( + secondLease, + secondNow, + policy); + + var query = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + now.AddHours(-4), + now, + MaxSeries: 1, + MaxPoints: 10)); + + var rollup = + Assert.Single( + Assert.Single(query.Series) + .Points, + point => + point.ResolutionSeconds == + 300); + + Assert.Equal(10, rollup.Min); + Assert.Equal(50, rollup.Max); + Assert.Equal(90, rollup.Sum); + Assert.Equal(3, rollup.Count); + Assert.Equal(30, rollup.Average); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_preserves_partial_and_unknown_truth_in_rollup() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-truth-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await store.InitializeAsync(); + await maintenance.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var rollupWindowStart = + DateTimeOffset.FromUnixTimeSeconds( + now.AddHours(-3) + .ToUnixTimeSeconds() / + 300 * + 300); + var identity = + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"); + var policy = + TestMaintenancePolicy(); + + await store.AppendAsync( + [ + HistoricalMetricSample.Gauge( + identity, + rollupWindowStart + .AddMinutes(1), + 10, + "consumer_observer", + "Unknown"), + HistoricalMetricSample.Gauge( + identity, + rollupWindowStart + .AddMinutes(2), + 20, + "consumer_observer", + "Partial"), + ]); + + var lease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + + var result = + await maintenance.RunCycleAsync( + lease, + now, + policy); + + Assert.True(result.LeaseValid); + + var query = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + now.AddHours(-4), + now, + MaxSeries: 1, + MaxPoints: 10)); + + var rollup = + Assert.Single( + Assert.Single(query.Series) + .Points, + point => + point.ResolutionSeconds == + policy.RollupResolutionSeconds); + + Assert.Equal( + "Partial", + rollup.State); + Assert.NotNull( + rollup.FirstObservedAtUtc); + Assert.NotNull( + rollup.LastObservedAtUtc); + Assert.True( + rollup.FirstObservedAtUtc <= + rollup.LastObservedAtUtc); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_rejects_retry_of_already_rolled_raw_identity() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-retry-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await store.InitializeAsync(); + await maintenance.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var observedAt = + now.AddHours(-3) + .AddMinutes(1); + var identity = + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"); + var sample = + HistoricalMetricSample.Gauge( + identity, + observedAt, + 10, + "consumer_observer", + "Stable"); + var policy = + TestMaintenancePolicy(); + + await store.AppendAsync([sample]); + + var firstLease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + _ = await maintenance.RunCycleAsync( + firstLease, + now, + policy); + + await store.AppendAsync([sample]); + + var secondNow = + DateTimeOffset.UtcNow; + var secondLease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + secondNow, + policy.LeaseDuration)); + _ = await maintenance.RunCycleAsync( + secondLease, + secondNow, + policy); + + var result = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + observedAt.AddHours(-1), + secondNow.AddMinutes(1), + MaxSeries: 1, + MaxPoints: 10)); + + var points = + Assert.Single(result.Series).Points; + var rollup = + Assert.Single( + points, + point => + point.ResolutionSeconds == 300); + + Assert.Equal(1, rollup.Count); + Assert.Equal(10, rollup.Sum); + Assert.DoesNotContain( + points, + point => + point.ResolutionSeconds == 0 && + point.ObservedAtUtc == observedAt); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_keeps_legacy_rollup_coverage_unknown_after_late_merge() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-legacy-late-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await store.InitializeAsync(); + await maintenance.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var bucket = + DateTimeOffset.FromUnixTimeSeconds( + now.AddHours(-3) + .ToUnixTimeSeconds() / + 300 * + 300); + var identity = + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"); + + await store.AppendAsync( + [ + new HistoricalMetricSample( + identity, + bucket, + 5, + 20, + 45, + 3, + 300, + "legacy-rollup", + "Unknown"), + HistoricalMetricSample.Gauge( + identity, + bucket.AddMinutes(1), + 30, + "consumer_observer", + "Stable"), + ]); + + // Simulate migrated legacy coverage: aggregate exists but original bounds are unknowable. + await using (var connection = + new SqliteConnection( + $"Data Source={path}")) + { + await connection.OpenAsync(); + await using var command = + connection.CreateCommand(); + command.CommandText = + """ + UPDATE kafdeck_historical_metric_samples + SET + first_observed_at_utc = NULL, + last_observed_at_utc = NULL + WHERE resolution_seconds = 300 + """; + await command.ExecuteNonQueryAsync(); + } + + var policy = + TestMaintenancePolicy(); + var lease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + _ = await maintenance.RunCycleAsync( + lease, + now, + policy); + + var result = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + bucket.AddHours(-1), + now.AddMinutes(1), + MaxSeries: 1, + MaxPoints: 10)); + + var rollup = + Assert.Single( + Assert.Single(result.Series).Points, + point => + point.ResolutionSeconds == 300); + Assert.False( + rollup.HasKnownCoverage); + Assert.Null( + rollup.FirstObservedAtUtc); + Assert.Null( + rollup.LastObservedAtUtc); + Assert.Equal( + "Partial", + rollup.State); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_rejects_lease_expired_before_cycle_execution() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-expired-lease-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await store.InitializeAsync(); + await maintenance.InitializeAsync(); + + var acquiredAt = + DateTimeOffset.UtcNow.AddMinutes(-3); + var lease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + acquiredAt, + TimeSpan.FromMinutes(2))); + + var result = + await maintenance.RunCycleAsync( + lease, + acquiredAt, + TestMaintenancePolicy()); + + Assert.False( + result.LeaseValid); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_counts_exact_full_batch_as_completed_window() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-full-batch-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await store.InitializeAsync(); + await maintenance.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var windowStart = + DateTimeOffset.FromUnixTimeSeconds( + now.AddHours(-3) + .ToUnixTimeSeconds() / + 300 * + 300); + var identity = + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"); + + var samples = + Enumerable.Range( + 0, + AdoHistoricalMetricMaintenanceStore + .MaxRawSamplesPerRollupBatch) + .Select( + index => + HistoricalMetricSample.Gauge( + identity, + windowStart.AddSeconds(index), + index + 1, + "consumer_observer", + "Stable")) + .Append( + HistoricalMetricSample.Gauge( + identity, + windowStart + .AddMinutes(5) + .AddSeconds(1), + 999, + "consumer_observer", + "Stable")) + .ToArray(); + + await store.AppendAsync(samples); + + var policy = + TestMaintenancePolicy() with + { + MaxRollupWindowsPerCycle = 1, + }; + var lease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + + var result = + await maintenance.RunCycleAsync( + lease, + now, + policy); + + Assert.True(result.LeaseValid); + Assert.Equal(1, result.RolledWindows); + Assert.Equal( + AdoHistoricalMetricMaintenanceStore + .MaxRawSamplesPerRollupBatch, + result.RawDeleted); + + var query = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + windowStart.AddMinutes(-1), + now.AddMinutes(1), + MaxSeries: 1, + MaxPoints: 100)); + + var points = + Assert.Single(query.Series).Points; + var rollup = + Assert.Single( + points, + point => + point.ResolutionSeconds == 300); + Assert.Equal( + AdoHistoricalMetricMaintenanceStore + .MaxRawSamplesPerRollupBatch, + rollup.Count); + Assert.Contains( + points, + point => + point.ResolutionSeconds == 0 && + point.ObservedAtUtc == + windowStart + .AddMinutes(5) + .AddSeconds(1)); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_lease_fences_stale_worker() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-fence-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await store.InitializeAsync(); + await maintenance.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var policy = + TestMaintenancePolicy(); + + var first = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + + Assert.Null( + await maintenance.TryAcquireLeaseAsync( + "node-b", + now.AddSeconds(30), + policy.LeaseDuration)); + + var takeoverAt = + now.AddMinutes(3); + var second = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-b", + takeoverAt, + policy.LeaseDuration)); + + Assert.True( + second.FencingToken > + first.FencingToken); + + var staleResult = + await maintenance.RunCycleAsync( + first, + takeoverAt, + policy); + Assert.False( + staleResult.LeaseValid); + + var activeResult = + await maintenance.RunCycleAsync( + second, + takeoverAt, + policy); + Assert.True( + activeResult.LeaseValid); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task PostgreSql_maintenance_lease_has_single_active_owner() + { + var baseConnectionString = + Environment.GetEnvironmentVariable( + "KAFDECK_TEST_POSTGRES"); + if (string.IsNullOrWhiteSpace( + baseConnectionString)) + { + return; + } + + var schema = + $"w63_maint_{Guid.NewGuid():N}"; + var adminBuilder = + new NpgsqlConnectionStringBuilder( + baseConnectionString); + await using var admin = + new NpgsqlConnection( + adminBuilder.ConnectionString); + await admin.OpenAsync(); + + try + { + await using (var create = + admin.CreateCommand()) + { + create.CommandText = + $"CREATE SCHEMA \"{schema}\""; + await create.ExecuteNonQueryAsync(); + } + + var scopedBuilder = + new NpgsqlConnectionStringBuilder( + baseConnectionString) + { + SearchPath = schema, + Pooling = false, + }; + var factory = + new PostgreSqlHistoricalMetricsDbConnectionFactory( + scopedBuilder.ConnectionString); + var history = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + var firstStore = + new AdoHistoricalMetricMaintenanceStore( + factory); + var secondStore = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await history.InitializeAsync(); + await firstStore.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var leaseDuration = + TimeSpan.FromMinutes(2); + + var leases = + await Task.WhenAll( + firstStore.TryAcquireLeaseAsync( + "node-a", + now, + leaseDuration), + secondStore.TryAcquireLeaseAsync( + "node-b", + now, + leaseDuration)); + + Assert.Equal( + 1, + leases.Count( + lease => + lease is not null)); + } + finally + { + await using var drop = + admin.CreateCommand(); + drop.CommandText = + $"DROP SCHEMA IF EXISTS \"{schema}\" CASCADE"; + await drop.ExecuteNonQueryAsync(); + } + } + + [Fact] + public async Task Maintenance_batches_expired_rollup_deletion() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-retention-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await store.InitializeAsync(); + await maintenance.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var identity = + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"); + + await store.AppendAsync( + [ + new HistoricalMetricSample( + identity, + now.AddDays(-10), + 1, + 1, + 1, + 1, + 300, + "rollup", + "Unknown"), + new HistoricalMetricSample( + identity, + now.AddDays(-9), + 2, + 2, + 2, + 1, + 300, + "rollup", + "Unknown"), + new HistoricalMetricSample( + identity, + now.AddDays(-8), + 3, + 3, + 3, + 1, + 300, + "rollup", + "Unknown"), + ]); + + var policy = + TestMaintenancePolicy() with + { + MaxRollupDeletesPerCycle = 2, + }; + + var firstLease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + var first = + await maintenance.RunCycleAsync( + firstLease, + now, + policy); + + Assert.True(first.LeaseValid); + Assert.Equal(2, first.RollupDeleted); + + var secondNow = + now.AddSeconds(30); + var secondLease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + secondNow, + policy.LeaseDuration)); + var second = + await maintenance.RunCycleAsync( + secondLease, + secondNow, + policy); + + Assert.True(second.LeaseValid); + Assert.Equal(1, second.RollupDeleted); + + var remaining = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + now.AddDays(-11), + now, + MaxSeries: 1, + MaxPoints: 10)); + + Assert.Empty(remaining.Series); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task PostgreSql_maintenance_uses_repeatable_read_snapshot() + { + var baseConnectionString = + Environment.GetEnvironmentVariable( + "KAFDECK_TEST_POSTGRES"); + if (string.IsNullOrWhiteSpace( + baseConnectionString)) + { + return; + } + + var schema = + $"w63_snapshot_{Guid.NewGuid():N}"; + var adminBuilder = + new NpgsqlConnectionStringBuilder( + baseConnectionString); + await using var admin = + new NpgsqlConnection( + adminBuilder.ConnectionString); + await admin.OpenAsync(); + + try + { + await using (var create = + admin.CreateCommand()) + { + create.CommandText = + $"CREATE SCHEMA \"{schema}\""; + await create.ExecuteNonQueryAsync(); + } + + var scopedBuilder = + new NpgsqlConnectionStringBuilder( + baseConnectionString) + { + SearchPath = schema, + Pooling = false, + }; + var factory = + new PostgreSqlHistoricalMetricsDbConnectionFactory( + scopedBuilder.ConnectionString); + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + + await using var connection = + await factory.OpenAsync(); + await using var transaction = + await maintenance.BeginMaintenanceTransactionAsync( + connection, + CancellationToken.None); + + Assert.Equal( + IsolationLevel.RepeatableRead, + transaction.IsolationLevel); + + await transaction.RollbackAsync(); + } + finally + { + await using var drop = + admin.CreateCommand(); + drop.CommandText = + $"DROP SCHEMA IF EXISTS \"{schema}\" CASCADE"; + await drop.ExecuteNonQueryAsync(); + } + } + + [Fact] + public void Maintenance_has_hard_raw_batch_and_cycle_batch_caps() + { + Assert.InRange( + AdoHistoricalMetricMaintenanceStore + .MaxRawSamplesPerRollupBatch, + 1, + 10_000); + Assert.InRange( + AdoHistoricalMetricMaintenanceStore + .MaxRollupBatchesPerCycle, + 1, + HistoricalMetricMaintenancePolicy + .HardMaxRollupWindowsPerCycle); + + var root = + FindRepositoryRoot(); + var source = + File.ReadAllText( + Path.Combine( + root, + "src", + "backend", + "Infrastructure", + "Kafdeck.Infrastructure.Persistence", + "AdoHistoricalMetricMaintenanceStore.cs")); + + Assert.Contains( + "LIMIT @raw_limit", + source, + StringComparison.Ordinal); + Assert.Contains( + "rollupBatches <", + source, + StringComparison.Ordinal); + Assert.Contains( + "rollupBatches++;", + source, + StringComparison.Ordinal); + } + + [Fact] + public void Maintenance_commits_each_rollup_window_in_its_own_transaction() + { + var root = + FindRepositoryRoot(); + var source = + File.ReadAllText( + Path.Combine( + root, + "src", + "backend", + "Infrastructure", + "Kafdeck.Infrastructure.Persistence", + "AdoHistoricalMetricMaintenanceStore.cs")); + + var loopIndex = + source.IndexOf( + "while (rolledWindows <", + StringComparison.Ordinal); + var transactionIndex = + source.IndexOf( + "await using var rollupTransaction", + loopIndex, + StringComparison.Ordinal); + var commitIndex = + source.IndexOf( + "await rollupTransaction", + transactionIndex, + StringComparison.Ordinal); + var retentionIndex = + source.IndexOf( + "long rollupDeleted = 0;", + commitIndex, + StringComparison.Ordinal); + + Assert.True(loopIndex >= 0); + Assert.True( + transactionIndex > loopIndex); + Assert.True( + commitIndex > transactionIndex); + Assert.True( + retentionIndex > commitIndex); + } + [Fact] public void Query_rejects_range_above_hard_cap() { @@ -579,6 +2096,18 @@ public void Query_rejects_range_above_hard_cap() query.Validate); } + private static HistoricalMetricMaintenancePolicy + TestMaintenancePolicy() => + new( + RawRetention: TimeSpan.FromHours(1), + RollupRetention: TimeSpan.FromDays(7), + RollupResolutionSeconds: 300, + MaxRollupWindowsPerCycle: 24, + MaxRollupDeletesPerCycle: 1_000, + LeaseDuration: TimeSpan.FromMinutes(2), + CycleInterval: TimeSpan.FromMinutes(1), + MaxCycleDuration: TimeSpan.FromSeconds(30)); + private static HistoricalMetricStorePolicy TestPolicy() => new( HistoricalMetricQuery.HardMaxRange, @@ -649,6 +2178,27 @@ private static KafdeckOptions Load( .AddInMemoryCollection(values) .Build()); + private static string FindRepositoryRoot() + { + DirectoryInfo? current = + new(AppContext.BaseDirectory); + while (current is not null) + { + if (File.Exists( + Path.Combine( + current.FullName, + "Kafdeck.slnx"))) + { + return current.FullName; + } + + current = current.Parent; + } + + throw new DirectoryNotFoundException( + "Unable to locate Kafdeck repository root."); + } + private static void DeleteSqliteFiles(string path) { foreach (var candidate in new[]