diff --git a/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs b/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs index 0c62a9d5..5c604469 100644 --- a/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs +++ b/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricMaintenanceStore.cs @@ -3,6 +3,7 @@ using System.Data; using System.Data.Common; using System.Globalization; +using System.Numerics; using Kafdeck.Core.Observability; using Microsoft.Data.Sqlite; using SQLitePCL; @@ -12,7 +13,7 @@ namespace Kafdeck.Infrastructure.Persistence; public sealed class AdoHistoricalMetricMaintenanceStore : IHistoricalMetricMaintenanceStore { - private const int SchemaVersion = 1; + private const int SchemaVersion = 2; private const int SingletonId = 1; private const string Component = "historical-metrics-maintenance"; @@ -31,6 +32,244 @@ private sealed record RawSampleRow( string? State, DateTimeOffset FirstObservedAtUtc, DateTimeOffset LastObservedAtUtc); + private readonly record struct ExactBinarySum( + BigInteger Significand, + int Exponent) + { + public static ExactBinarySum Zero => + new( + BigInteger.Zero, + 0); + + public static ExactBinarySum FromDouble( + double value) + { + if (!double.IsFinite(value)) + { + throw new ArgumentOutOfRangeException( + nameof(value)); + } + + if (value == 0) + { + return Zero; + } + + var bits = + unchecked( + (ulong)BitConverter + .DoubleToInt64Bits(value)); + var negative = + (bits >> 63) != 0; + var exponentBits = + (int)((bits >> 52) & 0x7ff); + var fraction = + bits & + 0x000f_ffff_ffff_ffffUL; + + BigInteger significand; + int exponent; + + if (exponentBits == 0) + { + significand = + new BigInteger( + fraction); + exponent = -1074; + } + else + { + significand = + new BigInteger( + (1UL << 52) | + fraction); + exponent = + exponentBits - + 1023 - + 52; + } + + if (negative) + { + significand = + -significand; + } + + return Normalize( + new ExactBinarySum( + significand, + exponent)); + } + + public ExactBinarySum Add( + ExactBinarySum other) + { + if (Significand.IsZero) + { + return other; + } + + if (other.Significand.IsZero) + { + return this; + } + + var commonExponent = + Math.Min( + Exponent, + other.Exponent); + var left = + Significand << + (Exponent - commonExponent); + var right = + other.Significand << + (other.Exponent - commonExponent); + + return Normalize( + new ExactBinarySum( + left + right, + commonExponent)); + } + + public double ToFiniteDouble( + out bool bounded) + { + bounded = false; + + if (Significand.IsZero) + { + return 0; + } + + var sign = + Significand.Sign; + var magnitude = + BigInteger.Abs( + Significand); + var max = + FromDouble( + double.MaxValue); + + if (CompareMagnitude( + this, + max) > 0) + { + bounded = true; + return sign < 0 + ? -double.MaxValue + : double.MaxValue; + } + + var bitLength = + magnitude.GetBitLength(); + var shift = + Math.Max( + 0, + checked( + (int)bitLength - 53)); + var top = + magnitude >> shift; + + if (shift > 0) + { + var remainder = + magnitude - + (top << shift); + var halfway = + BigInteger.One << + (shift - 1); + + if (remainder > halfway || + (remainder == halfway && + !top.IsEven)) + { + top += + BigInteger.One; + + if (top.GetBitLength() > 53) + { + top >>= 1; + shift++; + } + } + } + + var value = + Math.ScaleB( + (double)top, + checked( + Exponent + shift)); + + if (!double.IsFinite(value)) + { + bounded = true; + value = + double.MaxValue; + } + + return sign < 0 + ? -value + : value; + } + + private static int CompareMagnitude( + ExactBinarySum left, + ExactBinarySum right) + { + var commonExponent = + Math.Min( + left.Exponent, + right.Exponent); + var leftMagnitude = + BigInteger.Abs( + left.Significand) << + (left.Exponent - commonExponent); + var rightMagnitude = + BigInteger.Abs( + right.Significand) << + (right.Exponent - commonExponent); + + return leftMagnitude.CompareTo( + rightMagnitude); + } + + private static ExactBinarySum Normalize( + ExactBinarySum value) + { + var significand = + value.Significand; + var exponent = + value.Exponent; + + if (significand.IsZero) + { + return Zero; + } + + while (significand.IsEven) + { + significand >>= 1; + exponent++; + } + + return new ExactBinarySum( + significand, + exponent); + } + } + + private sealed record ExistingRollup( + double Min, + double Max, + string? State, + DateTimeOffset? FirstObservedAtUtc, + DateTimeOffset? LastObservedAtUtc, + ExactBinarySum ExactSum, + BigInteger ExactCount, + string? ExactState, + bool ExactSumIsLossy, + bool ExactCountIsLossy); + private const long MaintenanceMigrationLockKey = 4_839_176_502_110_873_342L; @@ -99,10 +338,24 @@ await ReadSchemaVersionAsync( .ConfigureAwait(false); if (existingVersion is not null && - existingVersion.Value != SchemaVersion) + (existingVersion.Value < 1 || + existingVersion.Value > SchemaVersion)) + { + throw new InvalidOperationException( + $"Historical metric maintenance schema version {existingVersion.Value} is unsupported by this binary (expected 1..{SchemaVersion})."); + } + + var historyVersion = + await ReadComponentSchemaVersionAsync( + connection, + transaction, + "historical-metrics", + cancellationToken) + .ConfigureAwait(false); + if (historyVersion != 3) { throw new InvalidOperationException( - $"Historical metric maintenance schema version {existingVersion.Value} is unsupported by this binary (expected {SchemaVersion})."); + $"Historical metric maintenance requires historical-metrics schema version 3; found {historyVersion?.ToString(CultureInfo.InvariantCulture) ?? "missing"}."); } await ExecuteAsync( @@ -153,6 +406,60 @@ ON kafdeck_historical_metric_raw_identities ( cancellationToken) .ConfigureAwait(false); + if (_connectionFactory.SupportsSelectForUpdate) + { + await ExecuteAsync( + connection, + transaction, + """ + CREATE OR REPLACE FUNCTION kafdeck_enforce_history_exact_writer_fence() + RETURNS trigger + LANGUAGE plpgsql + AS $function$ + BEGIN + IF OLD.resolution_seconds > 0 + AND OLD.exact_sum_significand IS NOT NULL + AND ( + NEW.sum_value IS DISTINCT FROM OLD.sum_value + OR NEW.sample_count IS DISTINCT FROM OLD.sample_count) + AND NEW.exact_sum_significand IS NOT DISTINCT FROM OLD.exact_sum_significand + AND NEW.exact_sum_exponent IS NOT DISTINCT FROM OLD.exact_sum_exponent + AND NEW.exact_sample_count IS NOT DISTINCT FROM OLD.exact_sample_count + AND NEW.exact_state IS NOT DISTINCT FROM OLD.exact_state + THEN + RAISE EXCEPTION 'KAFDECK_HISTORY_EXACT_WRITER_FENCE'; + END IF; + + RETURN NEW; + END; + $function$ + """, + cancellationToken) + .ConfigureAwait(false); + + await ExecuteAsync( + connection, + transaction, + """ + DROP TRIGGER IF EXISTS trg_kafdeck_history_exact_writer_fence + ON kafdeck_historical_metric_samples + """, + cancellationToken) + .ConfigureAwait(false); + + await ExecuteAsync( + connection, + transaction, + """ + CREATE TRIGGER trg_kafdeck_history_exact_writer_fence + BEFORE UPDATE ON kafdeck_historical_metric_samples + FOR EACH ROW + EXECUTE FUNCTION kafdeck_enforce_history_exact_writer_fence() + """, + cancellationToken) + .ConfigureAwait(false); + } + await using (var seed = connection.CreateCommand()) { @@ -215,6 +522,35 @@ await versionInsert .ConfigureAwait(false); } + if (existingVersion == 1) + { + 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); + + if (await versionUpdate + .ExecuteNonQueryAsync(cancellationToken) + .ConfigureAwait(false) != 1) + { + throw new InvalidOperationException( + "Historical metric maintenance schema version was upgraded without a durable component row."); + } + } + await transaction .CommitAsync(cancellationToken) .ConfigureAwait(false); @@ -822,19 +1158,45 @@ await ReadRawBatchAsync( group.Min(row => row.Min); var max = group.Max(row => row.Max); + var exactSum = + ExactBinarySum.Zero; + var exactCount = + BigInteger.Zero; + + foreach (var row in group) + { + exactSum = + exactSum.Add( + ExactBinarySum.FromDouble( + row.Sum)); + exactCount += + row.Count; + } + var sum = - group.Sum(row => row.Sum); + exactSum.ToFiniteDouble( + out var sumWasBounded); + var countWasBounded = + exactCount > + long.MaxValue; var count = - group.Sum(row => row.Count); + countWasBounded + ? long.MaxValue + : (long)exactCount; var states = group .Select(row => row.State) .Distinct(StringComparer.Ordinal) .ToArray(); - var state = + var baseState = states.Length == 1 ? states[0] : "Partial"; + var state = + sumWasBounded || + countWasBounded + ? "Partial" + : baseState; var firstObservedAtUtc = group.Min( row => @@ -859,6 +1221,9 @@ await UpsertRollupAsync( state, firstObservedAtUtc, lastObservedAtUtc), + exactSum, + exactCount, + baseState, cancellationToken) .ConfigureAwait(false); } @@ -989,12 +1354,348 @@ await command return rows; } + private static BigInteger ReadLegacyExactCount( + object value, + out bool wasLossy) + { + wasLossy = false; + + switch (value) + { + case byte byteValue: + return new BigInteger( + byteValue); + case short shortValue + when shortValue >= 0: + return new BigInteger( + shortValue); + case int intValue + when intValue >= 0: + return new BigInteger( + intValue); + case long longValue + when longValue >= 0: + return new BigInteger( + longValue); + case decimal decimalValue + when decimalValue >= 0: + { + var truncated = + decimal.Truncate( + decimalValue); + wasLossy = + truncated != + decimalValue; + return new BigInteger( + truncated); + } + case double doubleValue: + { + if (!double.IsFinite( + doubleValue) || + doubleValue < 0) + { + wasLossy = true; + return new BigInteger( + long.MaxValue) + + BigInteger.One; + } + + var truncated = + Math.Truncate( + doubleValue); + wasLossy = + truncated != + doubleValue; + return new BigInteger( + truncated); + } + case float floatValue: + { + if (!float.IsFinite( + floatValue) || + floatValue < 0) + { + wasLossy = true; + return new BigInteger( + long.MaxValue) + + BigInteger.One; + } + + var truncated = + MathF.Truncate( + floatValue); + wasLossy = + truncated != + floatValue; + return new BigInteger( + truncated); + } + case string text + when BigInteger.TryParse( + text, + NumberStyles.Integer, + CultureInfo.InvariantCulture, + out var parsed) && + parsed >= + BigInteger.Zero: + return parsed; + default: + wasLossy = true; + return new BigInteger( + long.MaxValue) + + BigInteger.One; + } + } + + private static async Task ReadExistingRollupAsync( + DbConnection connection, + DbTransaction transaction, + HistoricalMetricSample sample, + CancellationToken cancellationToken) + { + await using var command = + connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = + """ + SELECT + min_value, + max_value, + sum_value, + sample_count, + state, + first_observed_at_utc, + last_observed_at_utc, + exact_sum_significand, + exact_sum_exponent, + exact_sample_count, + exact_state + FROM kafdeck_historical_metric_samples + WHERE metric_name = @metric_name + AND cluster_id = @cluster_id + AND resource_kind = @resource_kind + AND resource_id = @resource_id + AND observed_at_utc = @observed_at_utc + AND resolution_seconds = @resolution_seconds + """; + 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); + + await using var reader = + await command + .ExecuteReaderAsync(cancellationToken) + .ConfigureAwait(false); + if (!await reader + .ReadAsync(cancellationToken) + .ConfigureAwait(false)) + { + return null; + } + + var storedSum = + Convert.ToDouble( + reader.GetValue(2), + CultureInfo.InvariantCulture); + var storedCount = + ReadLegacyExactCount( + reader.GetValue(3), + out var exactCountIsLossy); + + var hasExactSum = + !reader.IsDBNull(7) && + !reader.IsDBNull(8); + var exactSumIsLossy = + !hasExactSum && + !double.IsFinite( + storedSum); + var safeStoredSum = + double.IsFinite(storedSum) + ? storedSum + : storedSum > 0 + ? double.MaxValue + : storedSum < 0 + ? -double.MaxValue + : 0d; + var exactSum = + hasExactSum + ? new ExactBinarySum( + BigInteger.Parse( + reader.GetString(7), + CultureInfo.InvariantCulture), + Convert.ToInt32( + reader.GetValue(8), + CultureInfo.InvariantCulture)) + : ExactBinarySum.FromDouble( + safeStoredSum); + var exactCount = + !reader.IsDBNull(9) + ? BigInteger.Parse( + reader.GetString(9), + CultureInfo.InvariantCulture) + : storedCount; + + return new ExistingRollup( + Convert.ToDouble( + reader.GetValue(0), + CultureInfo.InvariantCulture), + Convert.ToDouble( + reader.GetValue(1), + CultureInfo.InvariantCulture), + reader.IsDBNull(4) + ? null + : reader.GetString(4), + reader.IsDBNull(5) + ? null + : DateTimeOffset.Parse( + reader.GetString(5), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind), + reader.IsDBNull(6) + ? null + : DateTimeOffset.Parse( + reader.GetString(6), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind), + exactSum, + exactCount, + reader.IsDBNull(10) + ? exactSumIsLossy + ? "Partial:LegacyNonFinite" + : null + : reader.GetString(10), + exactSumIsLossy, + exactCountIsLossy); + } + private static async Task UpsertRollupAsync( DbConnection connection, DbTransaction transaction, HistoricalMetricSample sample, + ExactBinarySum incomingExactSum, + BigInteger incomingExactCount, + string? incomingExactState, CancellationToken cancellationToken) { + var existing = + await ReadExistingRollupAsync( + connection, + transaction, + sample, + cancellationToken) + .ConfigureAwait(false); + + var exactSum = + existing is null + ? incomingExactSum + : existing.ExactSumIsLossy + ? existing.ExactSum + : existing.ExactSum.Add( + incomingExactSum); + var exactCount = + existing is null + ? incomingExactCount + : existing.ExactCountIsLossy + ? existing.ExactCount + : existing.ExactCount + + incomingExactCount; + var exactState = + existing is null + ? incomingExactState + : existing.ExactSumIsLossy + ? "Partial:LegacyNonFinite" + : string.Equals( + existing.ExactState ?? + existing.State, + incomingExactState, + StringComparison.Ordinal) + ? existing.ExactState ?? + existing.State + : "Partial"; + + var sum = + exactSum.ToFiniteDouble( + out var sumWasBounded); + var countWasBounded = + exactCount > + long.MaxValue; + var count = + countWasBounded + ? long.MaxValue + : (long)exactCount; + var state = + sumWasBounded || + countWasBounded || + existing?.ExactSumIsLossy == true || + existing?.ExactCountIsLossy == true + ? "Partial" + : exactState; + + var min = + existing is null + ? sample.Min + : Math.Min( + existing.Min, + sample.Min); + var max = + existing is null + ? sample.Max + : Math.Max( + existing.Max, + sample.Max); + + DateTimeOffset? firstObservedAtUtc = + sample.FirstObservedAtUtc; + DateTimeOffset? lastObservedAtUtc = + sample.LastObservedAtUtc; + + if (existing is not null) + { + if (existing.FirstObservedAtUtc is null || + existing.LastObservedAtUtc is null || + firstObservedAtUtc is null || + lastObservedAtUtc is null) + { + firstObservedAtUtc = null; + lastObservedAtUtc = null; + } + else + { + firstObservedAtUtc = + existing.FirstObservedAtUtc.Value < + firstObservedAtUtc.Value + ? existing.FirstObservedAtUtc + : firstObservedAtUtc; + lastObservedAtUtc = + existing.LastObservedAtUtc.Value > + lastObservedAtUtc.Value + ? existing.LastObservedAtUtc + : lastObservedAtUtc; + } + } + + var finalSample = + new HistoricalMetricSample( + sample.Identity, + sample.ObservedAtUtc, + min, + max, + sum, + count, + sample.ResolutionSeconds, + sample.Source, + state, + firstObservedAtUtc, + lastObservedAtUtc); + finalSample.Validate(); + await using var command = connection.CreateCommand(); command.Transaction = transaction; @@ -1014,7 +1715,11 @@ INSERT INTO kafdeck_historical_metric_samples ( source, state, first_observed_at_utc, - last_observed_at_utc) + last_observed_at_utc, + exact_sum_significand, + exact_sum_exponent, + exact_sample_count, + exact_state) VALUES ( @metric_name, @cluster_id, @@ -1029,7 +1734,11 @@ INSERT INTO kafdeck_historical_metric_samples ( @source, @state, @first_observed_at_utc, - @last_observed_at_utc) + @last_observed_at_utc, + @exact_sum_significand, + @exact_sum_exponent, + @exact_sample_count, + @exact_state) ON CONFLICT ( metric_name, cluster_id, @@ -1038,81 +1747,75 @@ ON CONFLICT ( 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, + min_value = excluded.min_value, + max_value = excluded.max_value, + sum_value = excluded.sum_value, + 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 + state = excluded.state, + first_observed_at_utc = excluded.first_observed_at_utc, + last_observed_at_utc = excluded.last_observed_at_utc, + exact_sum_significand = excluded.exact_sum_significand, + exact_sum_exponent = excluded.exact_sum_exponent, + exact_sample_count = excluded.exact_sample_count, + exact_state = excluded.exact_state """; - 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, "@metric_name", finalSample.Identity.MetricName); + AddParameter(command, "@cluster_id", finalSample.Identity.ClusterId); + AddParameter(command, "@resource_kind", finalSample.Identity.ResourceKind); + AddParameter(command, "@resource_id", finalSample.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); + finalSample.ObservedAtUtc.ToUniversalTime().ToString("O")); + AddParameter(command, "@resolution_seconds", finalSample.ResolutionSeconds); + AddParameter(command, "@min_value", finalSample.Min); + AddParameter(command, "@max_value", finalSample.Max); + AddParameter(command, "@sum_value", finalSample.Sum); + AddParameter(command, "@sample_count", finalSample.Count); + AddParameter(command, "@source", finalSample.Source); AddParameter( command, "@state", - sample.State is null + finalSample.State is null ? DBNull.Value - : sample.State); + : finalSample.State); AddParameter( command, "@first_observed_at_utc", - sample.FirstObservedAtUtc!.Value - .ToUniversalTime() - .ToString("O")); + finalSample.FirstObservedAtUtc is null + ? DBNull.Value + : finalSample.FirstObservedAtUtc.Value + .ToUniversalTime() + .ToString("O")); AddParameter( command, "@last_observed_at_utc", - sample.LastObservedAtUtc!.Value - .ToUniversalTime() - .ToString("O")); + finalSample.LastObservedAtUtc is null + ? DBNull.Value + : finalSample.LastObservedAtUtc.Value + .ToUniversalTime() + .ToString("O")); + AddParameter( + command, + "@exact_sum_significand", + exactSum.Significand.ToString( + CultureInfo.InvariantCulture)); + AddParameter( + command, + "@exact_sum_exponent", + exactSum.Exponent); + AddParameter( + command, + "@exact_sample_count", + exactCount.ToString( + CultureInfo.InvariantCulture)); + AddParameter( + command, + "@exact_state", + exactState is null + ? DBNull.Value + : exactState); await command .ExecuteNonQueryAsync(cancellationToken) @@ -1311,6 +2014,39 @@ await command .ConfigureAwait(false); } + private static async Task ReadComponentSchemaVersionAsync( + DbConnection connection, + DbTransaction transaction, + string component, + 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 async Task ReadSchemaVersionAsync( DbConnection connection, DbTransaction transaction, diff --git a/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricStore.cs b/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricStore.cs index 2453088b..559e07fb 100644 --- a/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricStore.cs +++ b/src/backend/Infrastructure/Kafdeck.Infrastructure.Persistence/AdoHistoricalMetricStore.cs @@ -9,7 +9,7 @@ namespace Kafdeck.Infrastructure.Persistence; public sealed class AdoHistoricalMetricStore : IHistoricalMetricStore { - private const int SchemaVersion = 2; + private const int SchemaVersion = 3; private const string Component = "historical-metrics"; internal static readonly TimeSpan RawIdentityRetention = @@ -119,6 +119,10 @@ CREATE TABLE IF NOT EXISTS kafdeck_historical_metric_samples ( state TEXT NULL, first_observed_at_utc TEXT NULL, last_observed_at_utc TEXT NULL, + exact_sum_significand TEXT NULL, + exact_sum_exponent INTEGER NULL, + exact_sample_count TEXT NULL, + exact_state TEXT NULL, PRIMARY KEY ( metric_name, cluster_id, @@ -221,6 +225,39 @@ await ExecuteInitializationStatementAsync( cancellationToken) .ConfigureAwait(false); } + } + + if (existingVersion is 1 or 2) + { + string[] versionThreeMigrationStatements = + [ + """ + ALTER TABLE kafdeck_historical_metric_samples + ADD COLUMN exact_sum_significand TEXT NULL + """, + """ + ALTER TABLE kafdeck_historical_metric_samples + ADD COLUMN exact_sum_exponent INTEGER NULL + """, + """ + ALTER TABLE kafdeck_historical_metric_samples + ADD COLUMN exact_sample_count TEXT NULL + """, + """ + ALTER TABLE kafdeck_historical_metric_samples + ADD COLUMN exact_state TEXT NULL + """, + ]; + + foreach (var statement in versionThreeMigrationStatements) + { + await ExecuteInitializationStatementAsync( + connection, + transaction, + statement, + cancellationToken) + .ConfigureAwait(false); + } await using var versionUpdate = connection.CreateCommand(); diff --git a/tests/Kafdeck.Architecture.Tests/V08W63HistoricalMetricsPersistenceTests.cs b/tests/Kafdeck.Architecture.Tests/V08W63HistoricalMetricsPersistenceTests.cs index c12ab210..1fb271eb 100644 --- a/tests/Kafdeck.Architecture.Tests/V08W63HistoricalMetricsPersistenceTests.cs +++ b/tests/Kafdeck.Architecture.Tests/V08W63HistoricalMetricsPersistenceTests.cs @@ -130,7 +130,7 @@ INSERT INTO kafdeck_schema_info ( schema_version) VALUES ( 'historical-metrics', - 3); + 4); """; await command.ExecuteNonQueryAsync(); } @@ -168,6 +168,111 @@ FROM sqlite_master } } + [Fact] + public async Task Version_two_schema_migrates_exact_accumulator_columns() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-v2-exact-{Guid.NewGuid():N}.db"); + + try + { + 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', + 2); + + 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, + first_observed_at_utc TEXT NULL, + last_observed_at_utc TEXT NULL, + PRIMARY KEY ( + metric_name, + cluster_id, + resource_kind, + resource_id, + observed_at_utc, + resolution_seconds) + ); + """; + await command.ExecuteNonQueryAsync(); + } + + var store = + new AdoHistoricalMetricStore( + new SqliteHistoricalMetricsDbConnectionFactory( + path), + TestPolicy()); + await store.InitializeAsync(); + + await using var verify = + new SqliteConnection( + $"Data Source={path}"); + await verify.OpenAsync(); + + await using var version = + verify.CreateCommand(); + version.CommandText = + """ + SELECT schema_version + FROM kafdeck_schema_info + WHERE component = 'historical-metrics' + """; + Assert.Equal( + 3L, + Convert.ToInt64( + await version.ExecuteScalarAsync())); + + await using var columns = + verify.CreateCommand(); + columns.CommandText = + """ + SELECT COUNT(*) + FROM pragma_table_info( + 'kafdeck_historical_metric_samples') + WHERE name IN ( + 'exact_sum_significand', + 'exact_sum_exponent', + 'exact_sample_count', + 'exact_state') + """; + Assert.Equal( + 4L, + Convert.ToInt64( + await columns.ExecuteScalarAsync())); + } + finally + { + DeleteSqliteFiles(path); + } + } + [Fact] public async Task Version_one_rollup_migration_preserves_unknown_coverage() { @@ -1636,11 +1741,11 @@ await store.QueryAsync( } [Fact] - public async Task Maintenance_lease_fences_stale_worker() + public async Task Maintenance_bounds_unrepresentable_batch_aggregate_as_partial() { var path = Path.Combine( Path.GetTempPath(), - $"kafdeck-history-fence-{Guid.NewGuid():N}.db"); + $"kafdeck-history-overflow-batch-{Guid.NewGuid():N}.db"); try { @@ -1660,151 +1765,94 @@ public async Task Maintenance_lease_fences_stale_worker() 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"); + + await store.AppendAsync( + [ + new HistoricalMetricSample( + identity, + windowStart.AddSeconds(1), + double.MaxValue, + double.MaxValue, + double.MaxValue, + long.MaxValue, + 0, + "consumer_observer", + "Stable"), + new HistoricalMetricSample( + identity, + windowStart.AddSeconds(2), + double.MaxValue, + double.MaxValue, + double.MaxValue, + 1, + 0, + "consumer_observer", + "Stable"), + ]); + var policy = TestMaintenancePolicy(); - - var first = + var lease = 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 = + var cycle = await maintenance.RunCycleAsync( - second, - takeoverAt, + lease, + now, 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(); + Assert.True(cycle.LeaseValid); + Assert.Equal(2, cycle.RawDeleted); - var now = - DateTimeOffset.UtcNow; - var leaseDuration = - TimeSpan.FromMinutes(2); + var query = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + windowStart.AddMinutes(-1), + now.AddMinutes(1), + MaxSeries: 1, + MaxPoints: 10)); - var leases = - await Task.WhenAll( - firstStore.TryAcquireLeaseAsync( - "node-a", - now, - leaseDuration), - secondStore.TryAcquireLeaseAsync( - "node-b", - now, - leaseDuration)); + var rollup = + Assert.Single( + Assert.Single(query.Series).Points); - Assert.Equal( - 1, - leases.Count( - lease => - lease is not null)); + Assert.Equal(long.MaxValue, rollup.Count); + Assert.Equal(double.MaxValue, rollup.Sum); + Assert.True(double.IsFinite(rollup.Sum)); + Assert.Equal("Partial", rollup.State); } finally { - await using var drop = - admin.CreateCommand(); - drop.CommandText = - $"DROP SCHEMA IF EXISTS \"{schema}\" CASCADE"; - await drop.ExecuteNonQueryAsync(); + DeleteSqliteFiles(path); } } [Fact] - public async Task Maintenance_batches_expired_rollup_deletion() + public async Task Maintenance_exact_sum_rounds_to_nearest_even() { var path = Path.Combine( Path.GetTempPath(), - $"kafdeck-history-retention-{Guid.NewGuid():N}.db"); + $"kafdeck-history-round-even-{Guid.NewGuid():N}.db"); try { @@ -1824,11 +1872,946 @@ public async Task Maintenance_batches_expired_rollup_deletion() var now = DateTimeOffset.UtcNow; - var identity = - new HistoricalMetricIdentity( - "consumer.lag.total", - "prod", - "consumer_group", + var windowStart = + 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, + windowStart.AddSeconds(1), + 1d, + "consumer_observer", + "Stable"), + HistoricalMetricSample.Gauge( + identity, + windowStart.AddSeconds(2), + Math.ScaleB(3d, -54), + "consumer_observer", + "Stable"), + ]); + + var policy = + TestMaintenancePolicy(); + var lease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + _ = await maintenance.RunCycleAsync( + lease, + now, + policy); + + var query = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + windowStart.AddMinutes(-1), + now.AddMinutes(1), + MaxSeries: 1, + MaxPoints: 10)); + + var rollup = + Assert.Single( + Assert.Single(query.Series).Points); + + Assert.Equal( + Math.BitIncrement(1d), + rollup.Sum); + Assert.Equal( + "Stable", + rollup.State); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_exact_batch_accumulator_recovers_representable_sum() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-exact-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"); + + await store.AppendAsync( + [ + new HistoricalMetricSample( + identity, + windowStart.AddSeconds(1), + double.MaxValue, + double.MaxValue, + double.MaxValue, + 1, + 0, + "consumer_observer", + "Stable"), + new HistoricalMetricSample( + identity, + windowStart.AddSeconds(2), + double.MaxValue, + double.MaxValue, + double.MaxValue, + 1, + 0, + "consumer_observer", + "Stable"), + new HistoricalMetricSample( + identity, + windowStart.AddSeconds(3), + -double.MaxValue, + -double.MaxValue, + -double.MaxValue, + 1, + 0, + "consumer_observer", + "Stable"), + ]); + + var policy = + TestMaintenancePolicy(); + var lease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + _ = await maintenance.RunCycleAsync( + lease, + now, + policy); + + var query = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + windowStart.AddMinutes(-1), + now.AddMinutes(1), + MaxSeries: 1, + MaxPoints: 10)); + + var rollup = + Assert.Single( + Assert.Single(query.Series).Points); + + Assert.Equal( + double.MaxValue, + rollup.Sum); + Assert.Equal( + 3, + rollup.Count); + Assert.Equal( + "Stable", + rollup.State); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_bounds_late_merge_overflow_as_partial() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-overflow-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 windowStart = + 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( + [ + new HistoricalMetricSample( + identity, + windowStart.AddSeconds(1), + double.MaxValue, + double.MaxValue, + double.MaxValue, + long.MaxValue, + 0, + "consumer_observer", + "Stable"), + ]); + + var firstLease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + _ = await maintenance.RunCycleAsync( + firstLease, + now, + policy); + + await store.AppendAsync( + [ + new HistoricalMetricSample( + identity, + windowStart.AddSeconds(2), + double.MaxValue, + double.MaxValue, + double.MaxValue, + 1, + 0, + "consumer_observer", + "Stable"), + ]); + + var secondNow = + DateTimeOffset.UtcNow; + 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", + windowStart.AddMinutes(-1), + secondNow.AddMinutes(1), + MaxSeries: 1, + MaxPoints: 10)); + + var rollup = + Assert.Single( + Assert.Single(query.Series).Points); + + Assert.Equal(long.MaxValue, rollup.Count); + Assert.Equal(double.MaxValue, rollup.Sum); + Assert.True(double.IsFinite(rollup.Sum)); + Assert.Equal("Partial", rollup.State); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_recovers_forward_progress_from_legacy_non_finite_rollup() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-legacy-infinity-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + + await store.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var windowStart = + DateTimeOffset.FromUnixTimeSeconds( + now.AddHours(-3) + .ToUnixTimeSeconds() / + 300 * + 300); + + await using (var connection = + new SqliteConnection( + $"Data Source={path}")) + { + await connection.OpenAsync(); + await using var command = + connection.CreateCommand(); + 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 ( + 'consumer.lag.total', + 'prod', + 'consumer_group', + 'group-a', + @bucket, + 300, + 1, + 1, + @sum_value, + 2, + 'legacy-rollup', + 'Stable', + @first_observed, + @last_observed) + """; + command.Parameters.AddWithValue( + "@bucket", + windowStart.ToString("O")); + command.Parameters.AddWithValue( + "@sum_value", + double.PositiveInfinity); + command.Parameters.AddWithValue( + "@first_observed", + windowStart.AddSeconds(1).ToString("O")); + command.Parameters.AddWithValue( + "@last_observed", + windowStart.AddSeconds(2).ToString("O")); + await command.ExecuteNonQueryAsync(); + } + + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + await maintenance.InitializeAsync(); + + var identity = + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"); + await store.AppendAsync( + [ + HistoricalMetricSample.Gauge( + identity, + windowStart.AddSeconds(3), + 10, + "consumer_observer", + "Stable"), + ]); + + 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.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: 10)); + var rollup = + Assert.Single( + Assert.Single(query.Series).Points); + + Assert.True( + double.IsFinite( + rollup.Sum)); + Assert.Equal( + double.MaxValue, + rollup.Sum); + Assert.Equal( + "Partial", + rollup.State); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_recovers_forward_progress_from_oversized_legacy_count() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-legacy-count-{Guid.NewGuid():N}.db"); + + try + { + var factory = + new SqliteHistoricalMetricsDbConnectionFactory( + path); + var store = + new AdoHistoricalMetricStore( + factory, + TestPolicy()); + + await store.InitializeAsync(); + + var now = + DateTimeOffset.UtcNow; + var windowStart = + DateTimeOffset.FromUnixTimeSeconds( + now.AddHours(-3) + .ToUnixTimeSeconds() / + 300 * + 300); + + await using (var connection = + new SqliteConnection( + $"Data Source={path}")) + { + await connection.OpenAsync(); + await using var command = + connection.CreateCommand(); + 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 ( + 'consumer.lag.total', + 'prod', + 'consumer_group', + 'group-a', + @bucket, + 300, + 1, + 1, + 10, + 9223372036854775808.0, + 'legacy-rollup', + 'Stable', + @first_observed, + @last_observed) + """; + command.Parameters.AddWithValue( + "@bucket", + windowStart.ToString("O")); + command.Parameters.AddWithValue( + "@first_observed", + windowStart.AddSeconds(1).ToString("O")); + command.Parameters.AddWithValue( + "@last_observed", + windowStart.AddSeconds(2).ToString("O")); + await command.ExecuteNonQueryAsync(); + } + + var maintenance = + new AdoHistoricalMetricMaintenanceStore( + factory); + await maintenance.InitializeAsync(); + + var identity = + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"); + await store.AppendAsync( + [ + HistoricalMetricSample.Gauge( + identity, + windowStart.AddSeconds(3), + 5, + "consumer_observer", + "Stable"), + ]); + + 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.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: 10)); + var rollup = + Assert.Single( + Assert.Single(query.Series).Points); + + Assert.Equal( + long.MaxValue, + rollup.Count); + Assert.True( + double.IsFinite( + rollup.Sum)); + Assert.Equal( + "Partial", + rollup.State); + } + finally + { + DeleteSqliteFiles(path); + } + } + + [Fact] + public async Task Maintenance_exact_late_accumulator_recovers_after_temporary_saturation() + { + var path = Path.Combine( + Path.GetTempPath(), + $"kafdeck-history-exact-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 windowStart = + DateTimeOffset.FromUnixTimeSeconds( + now.AddHours(-3) + .ToUnixTimeSeconds() / + 300 * + 300); + var identity = + new HistoricalMetricIdentity( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a"); + var policy = + TestMaintenancePolicy(); + + async Task AppendAndRollAsync( + double value, + int second) + { + await store.AppendAsync( + [ + new HistoricalMetricSample( + identity, + windowStart.AddSeconds(second), + value, + value, + value, + 1, + 0, + "consumer_observer", + "Stable"), + ]); + + var cycleNow = + DateTimeOffset.UtcNow; + var lease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + cycleNow, + policy.LeaseDuration)); + _ = await maintenance.RunCycleAsync( + lease, + cycleNow, + policy); + } + + await AppendAndRollAsync( + double.MaxValue, + 1); + await AppendAndRollAsync( + double.MaxValue, + 2); + + var saturated = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + windowStart.AddMinutes(-1), + DateTimeOffset.UtcNow.AddMinutes(1), + MaxSeries: 1, + MaxPoints: 10)); + var saturatedRollup = + Assert.Single( + Assert.Single(saturated.Series).Points); + Assert.Equal( + double.MaxValue, + saturatedRollup.Sum); + Assert.Equal( + "Partial", + saturatedRollup.State); + + await AppendAndRollAsync( + -double.MaxValue, + 3); + + var recovered = + await store.QueryAsync( + new HistoricalMetricQuery( + "consumer.lag.total", + "prod", + "consumer_group", + "group-a", + windowStart.AddMinutes(-1), + DateTimeOffset.UtcNow.AddMinutes(1), + MaxSeries: 1, + MaxPoints: 10)); + var recoveredRollup = + Assert.Single( + Assert.Single(recovered.Series).Points); + + Assert.Equal( + double.MaxValue, + recoveredRollup.Sum); + Assert.Equal( + 3, + recoveredRollup.Count); + Assert.Equal( + "Stable", + recoveredRollup.State); + } + 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( @@ -1923,6 +2906,131 @@ await store.QueryAsync( } } + [Fact] + public async Task PostgreSql_exact_writer_fence_rejects_legacy_rollup_update() + { + var baseConnectionString = + Environment.GetEnvironmentVariable( + "KAFDECK_TEST_POSTGRES"); + if (string.IsNullOrWhiteSpace( + baseConnectionString)) + { + return; + } + + var schema = + $"w63_exact_fence_{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 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"); + + await store.AppendAsync( + [ + HistoricalMetricSample.Gauge( + identity, + windowStart.AddSeconds(1), + 10, + "consumer_observer", + "Stable"), + ]); + + var policy = + TestMaintenancePolicy(); + var lease = + Assert.IsType( + await maintenance.TryAcquireLeaseAsync( + "node-a", + now, + policy.LeaseDuration)); + _ = await maintenance.RunCycleAsync( + lease, + now, + policy); + + await using var legacyWriter = + new NpgsqlConnection( + scopedBuilder.ConnectionString); + await legacyWriter.OpenAsync(); + await using var legacyUpdate = + legacyWriter.CreateCommand(); + legacyUpdate.CommandText = + """ + UPDATE kafdeck_historical_metric_samples + SET + sum_value = sum_value + 1, + sample_count = sample_count + 1 + WHERE resolution_seconds > 0 + """; + + var exception = + await Assert.ThrowsAsync( + () => legacyUpdate.ExecuteNonQueryAsync()); + + Assert.Contains( + "KAFDECK_HISTORY_EXACT_WRITER_FENCE", + exception.MessageText, + StringComparison.Ordinal); + } + finally + { + await using var drop = + admin.CreateCommand(); + drop.CommandText = + $"DROP SCHEMA IF EXISTS \"{schema}\" CASCADE"; + await drop.ExecuteNonQueryAsync(); + } + } + [Fact] public async Task PostgreSql_maintenance_uses_repeatable_read_snapshot() {