Skip to content
89 changes: 89 additions & 0 deletions Darling/Darling.Tests/DarlingRetentionTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -1052,6 +1052,95 @@ the command timeout. Pinned as arithmetic so raising the cap without re-measurin
+ "the margin is deliberate");
}

/* ---------------- #4250 item 3: the unordered row-capped prune for the liveness-touched tables ---------------- */

/// <summary>
/// The sibling builder's shape: row-capped like <see cref="DarlingRetention.RowCappedDeleteSql"/>, but
/// with NO <c>ORDER BY</c> — since V149 (#4250) dropped both tables' <c>last_seen</c> btree, an ORDER BY
/// would force a sort of every under-cutoff row before the LIMIT could apply, a second full pass over
/// data the seq scan already read. No <c>min()</c> subquery and no <c>INTERVAL</c> either: this is a
/// single bound against a FIXED cutoff, not a time slice.
/// </summary>
[Fact]
public void UnorderedRowCappedDelete_HasNoOrderByAndNoMinSubquery()
{
var sql = DarlingRetention.UnorderedRowCappedDeleteSql("collect.query_store_text", "last_seen", 300_000);

Assert.Equal(
"DELETE FROM collect.query_store_text WHERE ctid IN ("
+ "SELECT ctid FROM collect.query_store_text WHERE last_seen < $1 "
+ "LIMIT 300000)",
sql);

Assert.DoesNotContain("ORDER BY", sql, StringComparison.OrdinalIgnoreCase);
Assert.DoesNotContain("min(", sql, StringComparison.OrdinalIgnoreCase);
Assert.DoesNotContain("INTERVAL", sql, StringComparison.OrdinalIgnoreCase);
Assert.Contains("ctid IN", sql, StringComparison.Ordinal);
Assert.Contains("LIMIT 300000", sql, StringComparison.Ordinal);
}

/// <summary>
/// Both liveness-touched tables' purge calls use the unordered builder, carry the shared cap as their
/// batch size (so the drain loop's "a full-cap batch means there may be more" contract applies), and
/// leave <c>adaptiveRowCapTimeColumn</c> unset — that parameter exists only for the plan dimension's
/// #4130 retry-at-half-cap behavior, which neither of these tables has been measured to need: both
/// clear their whole steady-state backlog in ONE batch at the sized cap (see
/// <see cref="DarlingRetention.LivenessTouchedTablePruneRowCap"/>'s remarks), so there is nothing here
/// for a shrinking retry to protect against. Both call sites read the cap from the LOCAL parameter
/// <c>livenessTouchedTablePruneRowCap</c> rather than the constant directly, since #4250 item 3's live
/// loop test needs a seam to run the same call sites at a small cap; the parameter itself defaults to
/// the constant (pinned separately, below), so production is unchanged.
/// </summary>
[Fact]
public void MapAndTextPurges_UseTheUnorderedCap_WithoutTheAdaptiveRetry()
{
var source = ReadRetentionSource();

var mapAt = source.IndexOf("var mapCutoff = ComputeMapCutoff(", StringComparison.Ordinal);
Assert.True(mapAt >= 0, "the map purge call moved");
var mapBody = source[mapAt..Math.Min(source.Length, mapAt + 500)];
Assert.Contains("UnorderedRowCappedDeleteSql(", mapBody, StringComparison.Ordinal);
Assert.Contains("livenessTouchedTablePruneRowCap", mapBody, StringComparison.Ordinal);
Assert.DoesNotContain("adaptiveRowCapTimeColumn", mapBody, StringComparison.Ordinal);

var textAt = source.IndexOf("var queryTextCutoff = utcNow.AddDays(", StringComparison.Ordinal);
Assert.True(textAt >= 0, "the query text purge call moved");
var textBody = source[textAt..Math.Min(source.Length, textAt + 500)];
Assert.Contains("UnorderedRowCappedDeleteSql(", textBody, StringComparison.Ordinal);
Assert.Contains("livenessTouchedTablePruneRowCap", textBody, StringComparison.Ordinal);
Assert.DoesNotContain("adaptiveRowCapTimeColumn", textBody, StringComparison.Ordinal);
}

/// <summary>
/// The #4250 item 3 test seam's default: every real caller (the daily sweep, the on-demand
/// <c>purge_now</c> command) omits <c>livenessTouchedTablePruneRowCap</c>, so production must always
/// run the shipped 300,000-row constant, never a silently different value.
/// </summary>
[Fact]
public void LivenessTouchedTablePruneRowCapSeam_DefaultsToTheShippedConstant()
{
var method = typeof(DarlingRetention).GetMethod(
nameof(DarlingRetention.PurgeAsync), System.Reflection.BindingFlags.Public | System.Reflection.BindingFlags.Static)!;
var parameter = Array.Find(method.GetParameters(), p => p.Name == "livenessTouchedTablePruneRowCap")!;
Assert.NotNull(parameter);
Assert.Equal(DarlingRetention.LivenessTouchedTablePruneRowCap, (int)parameter.DefaultValue!);
}

/// <summary>
/// The cap clears the busiest single day observed in the field's <c>last_seen</c> age histogram (map
/// 255k rows at its oldest surviving day, text 220k) in ONE batch — a steady-state run never issues a
/// second, empty-batch statement.
/// </summary>
[Fact]
public void LivenessTouchedTablePruneRowCap_ClearsTheBusiestObservedDayInOneBatch()
{
const int busiestMapDay = 255_000;
const int busiestTextDay = 220_000;

Assert.True(DarlingRetention.LivenessTouchedTablePruneRowCap > busiestMapDay);
Assert.True(DarlingRetention.LivenessTouchedTablePruneRowCap > busiestTextDay);
}

/// <summary>
/// Only the plan dimension is capped. <c>query_text_dim</c> is ~40 MB in total and drains in a single
/// slice, and the fact tables need the compressed-chunk-safe shape that the <c>ctid</c> idiom cannot
Expand Down
247 changes: 247 additions & 0 deletions Darling/Darling.Tests/QueryStoreLivenessHotTouchLiveTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,247 @@
/*
* 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 V149 (#4250): the Query Store liveness touch becomes a HOT update on
/// <c>collect.query_store_plan_map</c> and <c>collect.query_store_text</c> — the two <c>last_seen</c> btree
/// indexes are dropped and both tables get <c>fillfactor = 90</c>. This file is the RUNG (ladder, viewer
/// probe) and the HOT-eligibility proof: the schema after migrate carries neither index and the fillfactor
/// reloption, startup's <c>CREATE TABLE IF NOT EXISTS</c> convergence does not recreate the index, and two
/// touches of the same row report a HOT update via <c>pg_stat_user_tables.n_tup_hot_upd</c>.
///
/// <para>This file's "I am the top rung" claim takes over from <c>ReadLatencyFlushLiveTests</c> (V148) now
/// that V149 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 QueryStoreLivenessHotTouchLiveTests
{
private const int RungVersion = 149;
private const int PreviousVersion = 148;

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

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>ReadLatencyFlushLiveTests</c> (V148) now that V149 has landed.
/// </summary>
[Fact]
public void TheRungIsRegisteredAtTheTopOfADenseLadder()
{
var versions = PgMigrations.Scripts.Select(s => s.Version).ToList();

Assert.Equal("query-store-liveness-hot-touch", 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_query_store_plan_map_last_seen", probe, StringComparison.Ordinal);
Assert.Contains("fillfactor=90", 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("hasHotLivenessTouch", 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 (hasHotLivenessTouch)", StringComparison.Ordinal);
var previousArm = viewer.IndexOf("if (hasReadLatency)", StringComparison.Ordinal);
Assert.True(thisArm >= 0, "the viewer has no V149 sentinel arm — a fully-migrated store would map one rung short");
Assert.True(thisArm < previousArm, "the V149 arm sits below V148'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: neither <c>last_seen</c> index exists on either table, and both carry
/// the <c>fillfactor=90</c> reloption. Run against <c>origin/dev</c> (pre-V149) this is RED — the indexes
/// exist and the reloption is absent — proving the pin actually checks the rung rather than a tautology.
/// </summary>
[Fact]
public async Task AfterMigrate_NeitherLastSeenIndexExists_AndBothTablesCarryFillfactor90()
{
var baseConnectionString = ConnectionString;
Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString),
"Set DARLING_TEST_PG to a Postgres connection string to run the V149 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.False(await IndexExistsAsync(connection, "idx_query_store_plan_map_last_seen", ct),
"V149 drops idx_query_store_plan_map_last_seen");
Assert.False(await IndexExistsAsync(connection, "idx_query_store_text_last_seen", ct),
"V149 drops idx_query_store_text_last_seen");

Assert.True(await HasFillfactor90Async(connection, "collect.query_store_plan_map", ct),
"V149 sets fillfactor = 90 on query_store_plan_map");
Assert.True(await HasFillfactor90Async(connection, "collect.query_store_text", ct),
"V149 sets fillfactor = 90 on query_store_text");
}

/// <summary>
/// Startup convergence (<see cref="QueryStorePlanMap.CreateTableSql"/>, <see cref="QueryStoreTextStore.CreateTableSql"/>)
/// no longer recreates either index: run the table-ensure path a second time after migrate, on top of an
/// already-migrated store, and the index must still be absent. A real mutation (re-adding the
/// <c>CREATE INDEX IF NOT EXISTS ... last_seen</c> line back to either helper's DDL) turns this RED —
/// see the PR body for the exact revert-and-rerun.
/// </summary>
[Fact]
public async Task StartupConvergence_RunTwice_DoesNotRecreateEitherIndex()
{
var baseConnectionString = ConnectionString;
Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString),
"Set DARLING_TEST_PG to a Postgres connection string to run the V149 convergence 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);

for (var i = 0; i < 2; i++)
{
await using var mapCommand = new NpgsqlCommand(QueryStorePlanMap.CreateTableSql, connection);
await mapCommand.ExecuteNonQueryAsync(ct);

await using var textCommand = new NpgsqlCommand(QueryStoreTextStore.CreateTableSql, connection);
await textCommand.ExecuteNonQueryAsync(ct);
}

Assert.False(await IndexExistsAsync(connection, "idx_query_store_plan_map_last_seen", ct),
"convergence must not recreate idx_query_store_plan_map_last_seen");
Assert.False(await IndexExistsAsync(connection, "idx_query_store_text_last_seen", ct),
"convergence must not recreate idx_query_store_text_last_seen");
}

/// <summary>
/// The payoff: two touches of the same map row after V149 report a HOT update. The first touch (row
/// starts fresh, no prior version) may or may not be HOT depending on page layout; the SECOND touch of
/// the same row — after the guard interval has re-elapsed — is the one this asserts, because it is the
/// steady-state case the field actually runs (every guard cycle touches rows that were touched last
/// cycle too). Uses <c>pg_stat_force_next_flush()</c> plus a short settle (as
/// <c>QueryStoreIntervalWideGridLiveTests</c> and <c>StoreToastAndCheckpointerTests</c> do) because
/// <c>pg_stat_user_tables</c> counters are backend-pending and throttled to once a second.
/// </summary>
[Fact]
public async Task TwoTouchesOfTheSameRow_TheSecondTouchReportsAHotUpdate()
{
var baseConnectionString = ConnectionString;
Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString),
"Set DARLING_TEST_PG to a Postgres connection string to run the V149 HOT 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 using var seed = new NpgsqlCommand(
"INSERT INTO collect.query_store_plan_map (server_id, database_name, plan_id, digest, plan_hash, last_seen) " +
"VALUES (1, 'db', 1, '\\x01', 'h1', now() AT TIME ZONE 'UTC' - interval '1 day')", connection);
await seed.ExecuteNonQueryAsync(ct);

var touch = "UPDATE collect.query_store_plan_map SET last_seen = now() AT TIME ZONE 'UTC' " +
"WHERE server_id = 1 AND database_name = 'db' AND plan_id = 1";

// First touch: may build a fresh heap version, not the steady-state case.
await using (var first = new NpgsqlCommand(touch, connection))
{
await first.ExecuteNonQueryAsync(ct);
}

await using (var flush1 = new NpgsqlCommand("SELECT pg_stat_force_next_flush()", connection))
{
await flush1.ExecuteScalarAsync(ct);
}

var before = await ReadHotUpdatesAsync(connection, ct);

// Second touch: the steady-state case — the row already has a settled heap version with room.
await using (var second = new NpgsqlCommand(touch, connection))
{
await second.ExecuteNonQueryAsync(ct);
}

await using (var flush2 = new NpgsqlCommand("SELECT pg_stat_force_next_flush()", connection))
{
await flush2.ExecuteScalarAsync(ct);
}

var after = await ReadHotUpdatesAsync(connection, ct);

Assert.True(after > before,
$"the second touch of an already-settled row must be HOT (n_tup_hot_upd {before} -> {after})");
}

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))!;
}

private static async Task<bool> HasFillfactor90Async(NpgsqlConnection connection, string tableName, System.Threading.CancellationToken ct)
{
await using var command = new NpgsqlCommand(
"SELECT coalesce((SELECT c.reloptions FROM pg_class c WHERE c.oid = $1::regclass) @> ARRAY['fillfactor=90'], false)", connection);
command.Parameters.AddWithValue(tableName);
return (bool)(await command.ExecuteScalarAsync(ct))!;
}

private static async Task<long> ReadHotUpdatesAsync(NpgsqlConnection connection, System.Threading.CancellationToken ct)
{
await using var command = new NpgsqlCommand(
"SELECT n_tup_hot_upd FROM pg_stat_user_tables WHERE schemaname = 'collect' AND relname = 'query_store_plan_map'", connection);
var scalar = await command.ExecuteScalarAsync(ct);
return Convert.ToInt64(scalar, System.Globalization.CultureInfo.InvariantCulture);
}
}
Loading
Loading