Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,151 @@
/*
* Copyright (c) 2026 Erik Darling, Darling Data LLC
*
* This file is part of the SQL Server Performance Monitor.
*
* Licensed under the MIT License. See LICENSE file in the project root for full license information.
*/

using System;
using System.Linq;
using System.Threading.Tasks;
using Npgsql;
using PerformanceMonitor.Darling.Storage;
using PerformanceMonitor.Darling.Viewer;
using Xunit;

namespace Darling.Tests;

/// <summary>
/// Pins Darling rung V150 (#4469, #4477): two supporting indexes, <c>idx_collection_log_watermark</c> on
/// <c>collect.collection_log (server_id, collector_name, collection_time DESC)</c> for the per-collector
/// watermark lookup, and <c>idx_job_history_server_run</c> on <c>collect.job_history (server_id,
/// run_datetime DESC, instance_id DESC)</c> for the Viewer's Job History tab. This file is the RUNG
/// (ladder, viewer probe) and the schema-after-migrate proof: both indexes exist, plain
/// <c>CREATE INDEX IF NOT EXISTS</c>, idempotent on rerun.
///
/// <para>This file's "I am the top rung" claim takes over from <c>QueryStoreLivenessHotTouchLiveTests</c>
/// (V149) now that V150 has landed.</para>
/// </summary>
/* #1776 own-store: each fact mints its own scratch database through ScratchPostgres and never touches the
shared store's tables, so it cannot race the live collection and serializing it would be pure slowdown. */
public sealed class CollectionLogWatermarkAndJobHistoryIndexesRungTests
{
private const int RungVersion = 150;
private const int PreviousVersion = 149;

/// <summary>This rung's sentinel ordinal in the viewer probe — the newest, so the last argument.</summary>
private const int ProbeOrdinal = 125;

private static string? ConnectionString => Environment.GetEnvironmentVariable("DARLING_TEST_PG");

/// <summary>
/// The rung is registered and is the new top of the ladder — the claim this class takes over from
/// <c>QueryStoreLivenessHotTouchLiveTests</c> (V149) now that V150 has landed.
/// </summary>
[Fact]
public void TheRungIsRegisteredAtTheTopOfADenseLadder()
{
var versions = PgMigrations.Scripts.Select(s => s.Version).ToList();

Assert.Equal("collection-log-watermark-and-job-history-indexes", PgMigrations.Scripts.Single(s => s.Version == RungVersion).Name);
Assert.Equal(StorageVersion.SchemaVersion, PgMigrations.Scripts[^1].Version);
Assert.Equal(StorageVersion.SchemaVersion, versions.Max());
Assert.Equal(RungVersion, StorageVersion.SchemaVersion);
Assert.Equal(versions.Distinct().OrderBy(v => v), versions);
}

/// <summary>
/// The viewer probe's sentinel carries this rung, and the map treats it as the TOP arm: a missing top arm
/// maps a fully-migrated store one rung short, permanently, because
/// <see cref="ViewerDataService.RequiredStoreSchemaVersion"/> is <see cref="StorageVersion.SchemaVersion"/>.
/// </summary>
[Fact]
public void TheProbeMapsAFullyMigratedStoreToThisTopRung()
{
var probe = ViewerDataService.StoreSchemaProbeSql.Replace("\r\n", "\n", StringComparison.Ordinal);
Assert.Contains("idx_collection_log_watermark", probe, StringComparison.Ordinal);
Assert.Contains("idx_job_history_server_run", probe, StringComparison.Ordinal);

var viewer = RepoFile.ReadRepoFile("Darling", "PerformanceMonitor.Darling.Viewer", "ViewerDataService.cs");
Assert.Contains($"reader.GetBoolean({ProbeOrdinal})", viewer, StringComparison.Ordinal);
Assert.DoesNotContain($"reader.GetBoolean({ProbeOrdinal + 1})", viewer, StringComparison.Ordinal);

Assert.Equal(StorageVersion.SchemaVersion, ViewerDataService.RequiredStoreSchemaVersion);

var method = typeof(ViewerDataService).GetMethod("MapProbedSchemaVersion", System.Reflection.BindingFlags.NonPublic | System.Reflection.BindingFlags.Static)!;
var arity = method.GetParameters().Length;
Assert.Equal(ProbeOrdinal, arity - 1);
Assert.Equal("hasCollectionLogWatermarkAndJobHistoryIndexes", method.GetParameters()[ProbeOrdinal].Name);

var all = Enumerable.Repeat((object)true, arity).ToArray();
Assert.Equal(StorageVersion.SchemaVersion, (int)method.Invoke(null, all)!);

var behind = (object[])all.Clone();
behind[ProbeOrdinal] = false;
Assert.Equal(PreviousVersion, (int)method.Invoke(null, behind)!);

var thisArm = viewer.IndexOf("if (hasCollectionLogWatermarkAndJobHistoryIndexes)", StringComparison.Ordinal);
var previousArm = viewer.IndexOf("if (hasHotLivenessTouch)", StringComparison.Ordinal);
Assert.True(thisArm >= 0, "the viewer has no V150 sentinel arm — a fully-migrated store would map one rung short");
Assert.True(thisArm < previousArm, "the V150 arm sits below V149's, so a current store maps one rung short");
Assert.Contains(
"return " + StorageVersion.SchemaVersion.ToString(System.Globalization.CultureInfo.InvariantCulture) + ";",
viewer[thisArm..previousArm], StringComparison.Ordinal);
}

/// <summary>
/// The LIVE schema after migrate: both indexes exist. Run against <c>origin/dev</c> (pre-V150) this is
/// RED — neither index exists — proving the pin actually checks the rung rather than a tautology.
/// </summary>
[Fact]
public async Task AfterMigrate_BothIndexesExist()
{
var baseConnectionString = ConnectionString;
Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString),
"Set DARLING_TEST_PG to a Postgres connection string to run the V150 schema pin.");

var ct = TestContext.Current.CancellationToken;

await using var scratch = await ScratchPostgres.CreateAsync(baseConnectionString!, ct);
await using var connection = new NpgsqlConnection(scratch.ConnectionString);
await connection.OpenAsync(ct);
await PgMigrations.MigrateAsync(connection, ct);

Assert.True(await IndexExistsAsync(connection, "idx_collection_log_watermark", ct),
"V150 creates idx_collection_log_watermark");
Assert.True(await IndexExistsAsync(connection, "idx_job_history_server_run", ct),
"V150 creates idx_job_history_server_run");
}

/// <summary>
/// Rerunning the migration ladder (as startup does on an already-migrated store) is idempotent: both
/// <c>CREATE INDEX IF NOT EXISTS</c> statements do not error and the indexes still exist.
/// </summary>
[Fact]
public async Task MigrateAsync_RunTwice_IsIdempotent()
{
var baseConnectionString = ConnectionString;
Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString),
"Set DARLING_TEST_PG to a Postgres connection string to run the V150 idempotency pin.");

var ct = TestContext.Current.CancellationToken;

await using var scratch = await ScratchPostgres.CreateAsync(baseConnectionString!, ct);
await using var connection = new NpgsqlConnection(scratch.ConnectionString);
await connection.OpenAsync(ct);
await PgMigrations.MigrateAsync(connection, ct);
await PgMigrations.MigrateAsync(connection, ct);

Assert.True(await IndexExistsAsync(connection, "idx_collection_log_watermark", ct));
Assert.True(await IndexExistsAsync(connection, "idx_job_history_server_run", ct));
}

private static async Task<bool> IndexExistsAsync(NpgsqlConnection connection, string indexName, System.Threading.CancellationToken ct)
{
await using var command = new NpgsqlCommand(
"SELECT EXISTS (SELECT 1 FROM pg_indexes WHERE schemaname = 'collect' AND indexname = $1)", connection);
command.Parameters.AddWithValue(indexName);
return (bool)(await command.ExecuteScalarAsync(ct))!;
}
}
27 changes: 24 additions & 3 deletions Darling/Darling.Tests/DarlingWatermarkFloorPlanShapeLiveTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -98,12 +98,32 @@ the field shape (30 of 32 chunks compressed). The newest row per collector sits
/* The plan-shape half: EXPLAIN the bounded statement and prove the chunk count is bounded by the
floor, not the 30 days of retention just seeded. */
var floor = Whole(nowUtc) - DarlingWorker.WatermarkFloorLookback;
var boundedPlan = await ExplainAsync(connection, DarlingWorker.ReadCollectorWatermarksSql, floor, ct);
var collectorNames = new[] { "wait_stats", "index_object_stats" };
var boundedPlan = await ExplainAsync(connection, DarlingWorker.ReadCollectorWatermarksSql, floor, ct, collectorNames);
var boundedChunks = PlanChunkScans.DistinctChunkCount(boundedPlan);

Assert.True(boundedChunks is >= 1 and <= 3,
$"expected the 2-day floor to touch at most a couple of chunks (touched={boundedChunks}):\n{boundedPlan}");

/* V150: one SubPlan/Limit descent per collector name, driven off unnest($3) — NOT a
GROUP BY/HashAggregate over every row in range. "Index Only Scan on idx_collection_log_watermark"
alone does not discriminate: the OLD GROUP BY statement, run against a store that already has
this index, ALSO plans as an Index Only Scan (TimescaleDB satisfies a GROUP BY's per-chunk
aggregate from the index instead of a table scan once the index exists) — Function Scan on
unnest only appears in the new per-collector LIMIT 1 shape, and HashAggregate/Finalize
HashAggregate only in the old GROUP BY shape. This is the exact regression this rung exists to
prevent: reverting the read to GROUP BY (variant C's shape too) must fail here even though the
index is still present. */
Assert.Contains("idx_collection_log_watermark", boundedPlan, StringComparison.Ordinal);
Assert.Contains("Function Scan on unnest", boundedPlan, StringComparison.Ordinal);
Assert.True(
boundedPlan.Contains("Index Only Scan", StringComparison.Ordinal)
|| boundedPlan.Contains("Index Scan", StringComparison.Ordinal),
$"expected an Index Only Scan or Index Scan on idx_collection_log_watermark:\n{boundedPlan}");
Assert.DoesNotContain("Bitmap Heap Scan", boundedPlan, StringComparison.Ordinal);
Assert.DoesNotContain("Seq Scan on _hyper", boundedPlan, StringComparison.Ordinal);
Assert.DoesNotContain("HashAggregate", boundedPlan, StringComparison.Ordinal);

/* Seed a SECOND server at 2x the chunk count (60 days) and prove the bounded plan's chunk count
does not grow with retention — the property #4469 exists to establish, not just "is small
today". */
Expand All @@ -115,7 +135,7 @@ the field shape (30 of 32 chunks compressed). The newest row per collector sits
}
await CompressOldChunksAsync(connection, ct);

var widePlan = await ExplainAsync(connection, DarlingWorker.ReadCollectorWatermarksSql, floor, ct, wideServerId);
var widePlan = await ExplainAsync(connection, DarlingWorker.ReadCollectorWatermarksSql, floor, ct, collectorNames, wideServerId);
var wideChunks = PlanChunkScans.DistinctChunkCount(widePlan);
Assert.True(wideChunks <= boundedChunks + 1,
$"expected doubling retained history to leave the bounded plan's chunk count essentially unchanged (30d={boundedChunks}, 60d={wideChunks}):\n{widePlan}");
Expand Down Expand Up @@ -188,11 +208,12 @@ failing the seed. */
}

private static async Task<string> ExplainAsync(
NpgsqlConnection connection, string sql, DateTime floor, CancellationToken ct, int? serverId = null)
NpgsqlConnection connection, string sql, DateTime floor, CancellationToken ct, string[] collectorNames, int? serverId = null)
{
using var command = new NpgsqlCommand("EXPLAIN (COSTS OFF) " + sql, connection);
command.Parameters.AddWithValue(serverId ?? LiveServerId);
command.Parameters.AddWithValue(floor);
command.Parameters.Add(new NpgsqlParameter<string[]> { TypedValue = collectorNames });
await using var reader = await command.ExecuteReaderAsync(ct);
var sb = new System.Text.StringBuilder();
while (await reader.ReadAsync(ct))
Expand Down
Loading
Loading