diff --git a/Darling/Darling.Tests/NpgsqlCommandCounter.cs b/Darling/Darling.Tests/NpgsqlCommandCounter.cs
new file mode 100644
index 0000000000..bcecea5b56
--- /dev/null
+++ b/Darling/Darling.Tests/NpgsqlCommandCounter.cs
@@ -0,0 +1,70 @@
+/*
+ * 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.Diagnostics;
+using System.Linq;
+using System.Threading;
+
+namespace Darling.Tests;
+
+///
+/// Counts the commands Npgsql runs against ONE database whose text contains every one of a set of fragments.
+/// Npgsql opens an per command on its Npgsql once a
+/// listener samples it, tagged with the database name and the command text; the count is taken when the
+/// activity stops, which is when the command's reader closes. The tags are matched by VALUE rather than by
+/// name (the same way SharedBaselineCacheTests.CaptureAsync finds the command text), because the
+/// OpenTelemetry attribute names Npgsql uses have changed between major versions.
+///
+/// Scoping by the scratch database's name keeps a class running in parallel against another database
+/// invisible to the count. The text carries the statement's $n placeholders, never the parameter values,
+/// so a count cannot tell one bound value from another: keep the scenario to one value, or count deltas. A
+/// caller trusts a count only after a control command it ran itself was counted once, so a renamed tag reads as
+/// a loud failure rather than a silent 0.
+///
+internal sealed class NpgsqlCommandCounter : IDisposable
+{
+ private readonly string _databaseName;
+ private readonly string[] _fragments;
+ private readonly ActivityListener _listener;
+ private long _count;
+
+ public NpgsqlCommandCounter(string databaseName, params string[] fragments)
+ {
+ _databaseName = databaseName;
+ _fragments = fragments;
+ _listener = new ActivityListener
+ {
+ ShouldListenTo = source => source.Name == "Npgsql",
+ Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded,
+ ActivityStopped = OnStopped,
+ };
+ ActivitySource.AddActivityListener(_listener);
+ }
+
+ public long Count => Interlocked.Read(ref _count);
+
+ private void OnStopped(Activity activity)
+ {
+ var onThisDatabase = false;
+ var carriesTheText = false;
+ foreach (var tag in activity.TagObjects)
+ {
+ if (tag.Value is not string value)
+ continue;
+ if (string.Equals(value, _databaseName, StringComparison.Ordinal))
+ onThisDatabase = true;
+ else if (_fragments.All(fragment => value.Contains(fragment, StringComparison.Ordinal)))
+ carriesTheText = true;
+ }
+ if (onThisDatabase && carriesTheText)
+ Interlocked.Increment(ref _count);
+ }
+
+ public void Dispose() => _listener.Dispose();
+}
diff --git a/Darling/Darling.Tests/QueryStoreBackfillCandidateLiveTests.cs b/Darling/Darling.Tests/QueryStoreBackfillCandidateLiveTests.cs
index 899a810704..60bd9e5fe5 100644
--- a/Darling/Darling.Tests/QueryStoreBackfillCandidateLiveTests.cs
+++ b/Darling/Darling.Tests/QueryStoreBackfillCandidateLiveTests.cs
@@ -163,6 +163,20 @@ public async Task CandidateRead_BoundAtItsOwnHorizon_ScansNoCompressedChunk()
Assert.DoesNotContain(compressedChunkName, plan);
Assert.True(PlanChunkScans.Count(plan) >= 1,
"expected the hot (2026-06-15) chunk to still be scanned:\n" + plan);
+
+ /* #4662: the TimescaleDB store's own read starts at the start of the chunk that holds the floor. That chunk
+ is newer than the compression boundary by design, so the read still excludes the compressed chunk. */
+ var cutStart = await ScalarStringAsync(connection, @"
+SELECT to_char(range_start AT TIME ZONE 'UTC', 'YYYY-MM-DD HH24:MI:SS.US') FROM timescaledb_information.chunks
+WHERE hypertable_name = 'query_store_stats' AND NOT is_compressed ORDER BY range_start LIMIT 1", ct);
+ var cutPlan = await ExplainAsync(connection, "EXPLAIN (COSTS OFF) " + QueryStoreBackfill.CutChunkCandidateSql
+ .Replace("$1", TestServerId.ToString(CultureInfo.InvariantCulture))
+ .Replace("{cut_start}", cutStart), ct);
+
+ Assert.DoesNotContain("DecompressChunk", cutPlan);
+ Assert.DoesNotContain("ColumnarScan", cutPlan);
+ Assert.DoesNotContain(compressedChunkName, cutPlan);
+ Assert.True(PlanChunkScans.Count(cutPlan) >= 1, "expected the hot (2026-06-15) chunk to still be scanned:\n" + cutPlan);
}
///
diff --git a/Darling/Darling.Tests/QueryStoreBackfillCutChunkLiveTests.cs b/Darling/Darling.Tests/QueryStoreBackfillCutChunkLiveTests.cs
new file mode 100644
index 0000000000..ccf53ff9c3
--- /dev/null
+++ b/Darling/Darling.Tests/QueryStoreBackfillCutChunkLiveTests.cs
@@ -0,0 +1,402 @@
+/*
+ * 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.Collections.Generic;
+using System.Diagnostics;
+using System.Globalization;
+using System.Linq;
+using System.Net;
+using System.Net.Sockets;
+using System.Text.Json;
+using System.Threading;
+using System.Threading.Tasks;
+using Microsoft.Data.SqlClient;
+using Npgsql;
+using PerformanceMonitor.Collectors;
+using PerformanceMonitor.Darling.Service;
+using PerformanceMonitor.Darling.Storage;
+using Xunit;
+
+namespace Darling.Tests;
+
+///
+/// #4662 end to end against a REAL store: the Query Store backfill's per-tick database list reads from the start
+/// of the chunk that holds the backfill floor (TimescaleDB) or walks the index (plain PostgreSQL), instead of
+/// filtering collection_time > floor inside the newest chunk, which walked that chunk's index entry by
+/// entry.
+///
+/// #1776 own-store - every test mints its own scratch database; nothing here touches the shared
+/// live fixture. The floors are fixed dates, not the wall clock, so the chunk each one lands in is fixed too.
+///
+public sealed class QueryStoreBackfillCutChunkLiveTests
+{
+ private const int Server = -466201;
+ private const int OtherServer = -466202;
+ private static long s_id = 4662000;
+
+ private static DateTime At(int month, int day, int hour, int minute = 0)
+ => new(2026, month, day, hour, minute, 0, DateTimeKind.Unspecified);
+
+ private static readonly IReadOnlyDictionary NoState = new Dictionary(StringComparer.Ordinal);
+
+ private static async Task<(ScratchPostgres Scratch, NpgsqlConnection Connection)> CreateStoreAsync(bool timescale, CancellationToken ct)
+ {
+ var baseConnectionString = Environment.GetEnvironmentVariable("DARLING_TEST_PG");
+ Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString),
+ "Set DARLING_TEST_PG to a Postgres connection string (with TimescaleDB installed) to run the live #4662 candidate-read tests (each mints its own scratch database).");
+
+ var scratch = await ScratchPostgres.CreateAsync(baseConnectionString!, ct);
+ var connection = new NpgsqlConnection(scratch.ConnectionString);
+ await connection.OpenAsync(ct);
+ await PgMigrations.MigrateAsync(connection, ct);
+ if (timescale)
+ {
+ Assert.True(await TimescaleSupport.TryEnableAsync(connection, null, ct), "the dev fixture is expected to have TimescaleDB installed");
+ await TimescaleSupport.ConvertToHypertablesAsync(connection, null, ct);
+ }
+
+ return (scratch, connection);
+ }
+
+ private static async Task SeedAsync(NpgsqlConnection connection, int serverId, string? database, DateTime collectionTime, CancellationToken ct)
+ {
+ const string sql = @"
+INSERT INTO collect.query_store_stats
+ (collection_id, collection_time, server_id, server_name, database_name, module_name, query_hash,
+ query_id, plan_id, execution_type_desc, replica_role,
+ runtime_stats_interval_id, interval_start_time_utc, first_execution_time, last_execution_time,
+ execution_count, avg_duration_us, avg_cpu_time_us, min_duration_us, max_duration_us)
+VALUES
+ ($1, $2, $3, 'SQL01', $4, 'dbo.GetOrders', '0xABCD', 91, 111, 'Regular', 'Primary',
+ 1, $2, $2, $2, 1, 100, 100, 100, 100)";
+ await using var command = new NpgsqlCommand(sql, connection);
+ command.Parameters.AddWithValue(Interlocked.Increment(ref s_id));
+ command.Parameters.AddWithValue(collectionTime);
+ command.Parameters.AddWithValue(serverId);
+ command.Parameters.AddWithValue((object?)database ?? DBNull.Value);
+ await command.ExecuteNonQueryAsync(ct);
+ }
+
+ private static async Task ExecAsync(NpgsqlConnection connection, string sql, CancellationToken ct)
+ {
+ await using var command = new NpgsqlCommand(sql, connection);
+ await command.ExecuteNonQueryAsync(ct);
+ }
+
+ private static QueryStoreBackfill Backfill(NpgsqlDataSource postgres, bool timescale)
+ => new(postgres, new DarlingCollectorRunner(postgres, new CollectorDeltaCalculator()), new CollectorDeltaCalculator(), logger: null,
+ hasContinuousAggregates: () => timescale);
+
+ private static Dictionary HoleState(string database, DateTime floor)
+ => new(StringComparer.Ordinal)
+ {
+ [QueryStoreBackfillState.HoleKeyPrefix + database] = QueryStoreBackfillState.EncodeHole(floor.AddDays(-1), floor.AddDays(1)),
+ };
+
+ /// Pin 1. The floor is 09:00 inside the 1-day chunk [06-15 00:00, 06-16 00:00). The list is the names
+ /// with rows at or after the chunk's start - one with rows after the floor, one with rows only between the
+ /// chunk start and the floor - plus the name a hole key adds; a name whose rows are all in the previous chunk,
+ /// another server's name and a NULL name are not in it. Ordinal order.
+ [Fact]
+ public async Task TimescaleStore_ListsNamesFromTheCutChunkStart_ExactlyThoseAndTheHoleName()
+ {
+ var ct = TestContext.Current.CancellationToken;
+ var floor = At(6, 15, 9);
+ var (scratch, connection) = await CreateStoreAsync(timescale: true, ct);
+ await using var scratchOwner = scratch;
+ await using var connectionOwner = connection;
+
+ await SeedAsync(connection, Server, "busy", floor.AddHours(1), ct);
+ await SeedAsync(connection, Server, "tail_only", At(6, 15, 3), ct);
+ await SeedAsync(connection, Server, "tail_at_start", At(6, 15, 0), ct);
+ await SeedAsync(connection, Server, "before_chunk", At(6, 14, 23, 59), ct);
+ await SeedAsync(connection, OtherServer, "other_server", floor.AddHours(1), ct);
+ await SeedAsync(connection, Server, null, floor.AddHours(1), ct);
+
+ await using var postgres = NpgsqlDataSource.Create(scratch.ConnectionString);
+ var list = await Backfill(postgres, timescale: true)
+ .GetCandidateDatabasesAsync(Server, floor, HoleState("zz_hole_no_rows", floor), ct);
+
+ Assert.Equal(new[] { "busy", "tail_at_start", "tail_only", "zz_hole_no_rows" }, list);
+ }
+
+ /// Pin 2. Plain PostgreSQL has no chunks: the walk lists every non-NULL name stored for the server,
+ /// older ones included, and nothing from another server.
+ [Fact]
+ public async Task PlainStore_WalksEveryNameOfTheServer_OlderIncluded_NothingFromAnotherServer()
+ {
+ var ct = TestContext.Current.CancellationToken;
+ var floor = At(6, 15, 9);
+ var (scratch, connection) = await CreateStoreAsync(timescale: false, ct);
+ await using var scratchOwner = scratch;
+ await using var connectionOwner = connection;
+
+ await SeedAsync(connection, Server, "busy", floor.AddHours(1), ct);
+ await SeedAsync(connection, Server, "older_than_the_floor", At(6, 1, 12), ct);
+ await SeedAsync(connection, Server, "older_still", At(5, 20, 12), ct);
+ await SeedAsync(connection, OtherServer, "other_server", floor.AddHours(1), ct);
+ await SeedAsync(connection, Server, null, floor.AddHours(1), ct);
+
+ await using var postgres = NpgsqlDataSource.Create(scratch.ConnectionString);
+ var list = await Backfill(postgres, timescale: false).GetCandidateDatabasesAsync(Server, floor, NoState, ct);
+
+ Assert.Equal(new[] { "busy", "older_still", "older_than_the_floor" }, list);
+ }
+
+ /// Pin 3. The boundary is the catalog's, not a constant: set_chunk_time_interval changes only new
+ /// chunks. Old rows sit in 1-day chunks, the interval is then set to 6 hours, newer rows land in 6-hour chunks.
+ /// A floor in an old 1-day chunk lists from that chunk's start; a floor in a new 6-hour chunk lists from the
+ /// 6-hour chunk's start (a 1-day boundary would also list n_before, a 6-hour one would miss early).
+ [Fact]
+ public async Task TimescaleStore_TakesTheBoundaryFromTheCatalog_AcrossAnIntervalChange()
+ {
+ var ct = TestContext.Current.CancellationToken;
+ var (scratch, connection) = await CreateStoreAsync(timescale: true, ct);
+ await using var scratchOwner = scratch;
+ await using var connectionOwner = connection;
+
+ await SeedAsync(connection, Server, "previous_day", At(6, 9, 22), ct);
+ await SeedAsync(connection, Server, "early", At(6, 10, 2), ct);
+ await SeedAsync(connection, Server, "late", At(6, 10, 21), ct);
+ await ExecAsync(connection, "SELECT set_chunk_time_interval('collect.query_store_stats'::regclass, INTERVAL '6 hours')", ct);
+ await SeedAsync(connection, Server, "n_before", At(6, 20, 7), ct);
+ await SeedAsync(connection, Server, "n_in", At(6, 20, 13), ct);
+ await SeedAsync(connection, Server, "n_after", At(6, 20, 15), ct);
+
+ await using var postgres = NpgsqlDataSource.Create(scratch.ConnectionString);
+ var backfill = Backfill(postgres, timescale: true);
+
+ var inOldChunk = await backfill.GetCandidateDatabasesAsync(Server, At(6, 10, 20), NoState, ct);
+ Assert.Equal(new[] { "early", "late", "n_after", "n_before", "n_in" }, inOldChunk);
+
+ var inNewChunk = await backfill.GetCandidateDatabasesAsync(Server, At(6, 20, 14), NoState, ct);
+ Assert.Equal(new[] { "n_after", "n_in" }, inNewChunk);
+ }
+
+ /// Pin 4. A floor that no chunk holds (the day between two seeded chunks) has no catalog row, so the
+ /// read is the floor-bound statement it always was: names with rows after the floor, none from before it.
+ [Fact]
+ public async Task TimescaleStore_AFloorNoChunkHolds_FallsBackToTheFloorBoundList()
+ {
+ var ct = TestContext.Current.CancellationToken;
+ var (scratch, connection) = await CreateStoreAsync(timescale: true, ct);
+ await using var scratchOwner = scratch;
+ await using var connectionOwner = connection;
+
+ await SeedAsync(connection, Server, "old", At(6, 13, 12), ct);
+ await SeedAsync(connection, Server, "busy", At(6, 15, 10), ct);
+ await SeedAsync(connection, Server, "tail_of_next_chunk", At(6, 15, 1), ct);
+
+ await using var postgres = NpgsqlDataSource.Create(scratch.ConnectionString);
+ var list = await Backfill(postgres, timescale: true).GetCandidateDatabasesAsync(Server, At(6, 14, 12), NoState, ct);
+
+ Assert.Equal(new[] { "busy", "tail_of_next_chunk" }, list);
+ }
+
+ /* ---- pins 5 to 7: what a whole tick spends on the extra names, a recorded hole, and the walk's plan. Pins 5 and 6
+ drive RunServerSliceAsync itself, which computes its floor from the wall clock, so their rows are seeded
+ relative to that floor. ---- */
+
+ /// The floor RunServerSliceAsync will compute moves with the wall clock (now minus the
+ /// horizon), and 1-day chunk boundaries fall at UTC midnight. Seeding and the tick are seconds apart, so this
+ /// waits until the floor is three minutes clear of a boundary on both sides; a boundary between the seed and
+ /// the tick would put the floor in another chunk than the seeded rows. Returns the floor, Unspecified kind.
+ private static async Task WaitForFloorClearOfAChunkBoundaryAsync(CancellationToken ct)
+ {
+ var clearance = TimeSpan.FromMinutes(3);
+ while (true)
+ {
+ var floor = DateTime.UtcNow - QueryStoreBackfill.HorizonFor(hasContinuousAggregates: true);
+ if (floor.TimeOfDay >= clearance && floor.TimeOfDay <= TimeSpan.FromDays(1) - clearance)
+ {
+ return DateTime.SpecifyKind(floor, DateTimeKind.Unspecified);
+ }
+
+ await Task.Delay(TimeSpan.FromSeconds(10), ct);
+ }
+ }
+
+ /// A connection string to a loopback port nothing listens on, so a slice - which opens a SqlConnection
+ /// before it does anything else - fails at once with a connection instead of reaching
+ /// a server. Connect Timeout=2 bounds it either way.
+ private static string RefusedConnectionString()
+ {
+ var listener = new TcpListener(IPAddress.Loopback, 0);
+ listener.Start();
+ var port = ((IPEndPoint)listener.LocalEndpoint).Port;
+ listener.Stop();
+ return $"Server=127.0.0.1,{port};User ID=unused;Password=unused;Connect Timeout=2;Encrypt=False;Pooling=False";
+ }
+
+ private static ServerRuntime ServerAt(string connectionString) => new()
+ {
+ Config = new MonitoredServer { Name = "qs4662", Host = "qs4662-host" },
+ ConnectionString = connectionString,
+ Target = new CollectorTargetInfo { SqlMajorVersion = 16 },
+ StorageName = "qs4662-host",
+ ServerId = Server,
+ EngineEdition = 3,
+ };
+
+ /// Counts the two floor reads (read A: SELECT 1 ... collection_time <= $3 LIMIT 1; read B:
+ /// SELECT MIN(last_execution_time)) on one scratch database. Both are trusted only after the product's own
+ /// ran once for a name with no rows (read A misses, so read B
+ /// runs) and each counted exactly once; the counts are then read as deltas from that point.
+ private static async Task<(NpgsqlCommandCounter ReadA, NpgsqlCommandCounter ReadB, long BaseA, long BaseB)> CountFloorReadsAsync(
+ QueryStoreBackfill backfill, string databaseName, DateTime floor, CancellationToken ct)
+ {
+ var readA = new NpgsqlCommandCounter(databaseName, "SELECT 1 FROM query_store_stats", "database_name = $2", "collection_time <= $3");
+ var readB = new NpgsqlCommandCounter(databaseName, "SELECT MIN(last_execution_time) FROM query_store_stats");
+
+ var missing = await backfill.GetStoredFloorAsync(Server, "a_name_with_no_rows", floor, ct);
+ Assert.Null(missing);
+ Assert.True(readA.Count == 1 && readB.Count == 1,
+ $"The Npgsql activity listener counted read A {readA.Count} times and read B {readB.Count} times for one control call on '{databaseName}' instead of 1 and 1 - it cannot see the database name or the command text, so a count of the ticks' reads would be meaningless.");
+ return (readA, readB, readA.Count, readB.Count);
+ }
+
+ /// Pin 5. A name whose rows are all in the cut chunk but before the floor is listed by the cut-chunk read
+ /// (the read starts at the chunk's start) and has no Done key, so it reaches the floor reads. It costs exactly one
+ /// read A (which hits, the row being at or before the floor) and one Done write, and it never reaches read B or a
+ /// slice (a slice would throw here: the server's port refuses). The next tick lists it again, sees the Done key
+ /// and runs no read for it.
+ [Fact]
+ public async Task TimescaleStore_AnExtraName_CostsOneFloorReadAndOneDoneWrite_Once()
+ {
+ var ct = TestContext.Current.CancellationToken;
+ var (scratch, connection) = await CreateStoreAsync(timescale: true, ct);
+ await using var scratchOwner = scratch;
+ await using var connectionOwner = connection;
+
+ var floor = await WaitForFloorClearOfAChunkBoundaryAsync(ct);
+ await SeedAsync(connection, Server, "quiet_extra", floor.AddMinutes(-1), ct);
+
+ await using var postgres = NpgsqlDataSource.Create(scratch.ConnectionString);
+ var runner = new DarlingCollectorRunner(postgres, new CollectorDeltaCalculator());
+ var backfill = new QueryStoreBackfill(postgres, runner, new CollectorDeltaCalculator(), logger: null, hasContinuousAggregates: () => true);
+ var server = ServerAt(RefusedConnectionString());
+ var (readA, readB, baseA, baseB) = await CountFloorReadsAsync(backfill, scratch.DatabaseName, floor, ct);
+ using var readAOwner = readA;
+ using var readBOwner = readB;
+
+ Assert.False(await backfill.RunServerSliceAsync(server, ct), "the first tick has no slice to run");
+ Assert.Equal(1, readA.Count - baseA);
+ Assert.Equal(0, readB.Count - baseB);
+ var state = await runner.GetCollectorStateAsync(Server, QueryStoreBackfill.StateCollectorName, ct);
+ Assert.True(state.ContainsKey(QueryStoreBackfillState.DoneKeyPrefix + "quiet_extra"),
+ "the first tick must save the Done key for the extra name, or every later tick reads for it again");
+
+ Assert.False(await backfill.RunServerSliceAsync(server, ct), "the second tick has no slice to run");
+ Assert.Equal(1, readA.Count - baseA);
+ Assert.Equal(0, readB.Count - baseB);
+ }
+
+ /// Pin 6. A database with a recorded hole that is also marked Done is still dug: the hole check runs
+ /// before the Done check. The tick attempts the hole's slice (which throws the connection error, the port
+ /// refusing) and runs neither floor read for it. This is what keeps a recorded outage gap backfilled after its
+ /// database finished its first-contact tail.
+ [Fact]
+ public async Task TimescaleStore_ARecordedHoleOnADoneDatabase_IsStillDug_WithoutAFloorRead()
+ {
+ var ct = TestContext.Current.CancellationToken;
+ var (scratch, connection) = await CreateStoreAsync(timescale: true, ct);
+ await using var scratchOwner = scratch;
+ await using var connectionOwner = connection;
+
+ var floor = await WaitForFloorClearOfAChunkBoundaryAsync(ct);
+ await SeedAsync(connection, Server, "holed_db", floor.AddMinutes(-1), ct);
+
+ await using var postgres = NpgsqlDataSource.Create(scratch.ConnectionString);
+ var runner = new DarlingCollectorRunner(postgres, new CollectorDeltaCalculator());
+ var backfill = new QueryStoreBackfill(postgres, runner, new CollectorDeltaCalculator(), logger: null, hasContinuousAggregates: () => true);
+ await runner.SaveCollectorStateAsync(Server, QueryStoreBackfill.StateCollectorName, new Dictionary(StringComparer.Ordinal)
+ {
+ [QueryStoreBackfillState.DoneKeyPrefix + "holed_db"] = DateTime.UtcNow.ToString("o", CultureInfo.InvariantCulture),
+ [QueryStoreBackfillState.HoleKeyPrefix + "holed_db"] = QueryStoreBackfillState.EncodeHole(floor.AddHours(-2), floor.AddHours(1)),
+ }, ct);
+ var server = ServerAt(RefusedConnectionString());
+ var (readA, readB, baseA, baseB) = await CountFloorReadsAsync(backfill, scratch.DatabaseName, floor, ct);
+ using var readAOwner = readA;
+ using var readBOwner = readB;
+
+ var stopwatch = Stopwatch.StartNew();
+ await Assert.ThrowsAsync(() => backfill.RunServerSliceAsync(server, ct));
+ TestContext.Current.SendDiagnosticMessage("QS4662_HOLE_SLICE_REFUSED_SECONDS=" + stopwatch.Elapsed.TotalSeconds.ToString("F2", CultureInfo.InvariantCulture));
+
+ Assert.Equal(0, readA.Count - baseA);
+ Assert.Equal(0, readB.Count - baseB);
+ }
+
+ /// Pin 7. On plain PostgreSQL the walk costs one index seek per database, not a scan of the table: the
+ /// recursive term's inner select is a Limit over an Index (Only) Scan of query_store_stats, and nothing in
+ /// the plan is a Seq Scan. The store carries the index the worker's start step creates (PgTableTuning), the table
+ /// holds 30,000 rows of two servers and five names, and it is ANALYZEd, so the index is the planner's own choice
+ /// rather than a forced one. The plan is EXPLAIN ANALYZE, not plain EXPLAIN: on PostgreSQL 18.6 the plain form
+ /// leaves the correlated select out of the recursive term's plan (it shows only the work table scan), and the
+ /// analyzed form runs the walk, which is six index seeks.
+ [Fact]
+ public async Task PlainStore_TheWalksInnerSelect_IsAnIndexSeekUnderALimit_NotASeqScan()
+ {
+ var ct = TestContext.Current.CancellationToken;
+ var (scratch, connection) = await CreateStoreAsync(timescale: false, ct);
+ await using var scratchOwner = scratch;
+ await using var connectionOwner = connection;
+
+ await PgTableTuning.ApplyAsync(connection, null, ct);
+ await ExecAsync(connection, @"
+INSERT INTO collect.query_store_stats
+ (collection_id, collection_time, server_id, server_name, database_name, module_name, query_hash,
+ query_id, plan_id, execution_type_desc, replica_role,
+ runtime_stats_interval_id, interval_start_time_utc, first_execution_time, last_execution_time,
+ execution_count, avg_duration_us, avg_cpu_time_us, min_duration_us, max_duration_us)
+SELECT 4662500000 + row_number() OVER (), x.t, s.server_id, 'SQL01', d.name, 'dbo.GetOrders', '0xABCD',
+ g, g, 'Regular', 'Primary', 1, x.t, x.t, x.t, 1, 100, 100, 100, 100
+FROM (VALUES (-466201), (-466202)) AS s(server_id)
+CROSS JOIN (VALUES ('db_a'), ('db_b'), ('db_c'), ('db_d'), ('db_e')) AS d(name)
+CROSS JOIN generate_series(1, 3000) AS g
+CROSS JOIN LATERAL (SELECT TIMESTAMP '2026-06-15 09:00:00' + g * INTERVAL '1 second') AS x(t)", ct);
+ await ExecAsync(connection, "ANALYZE collect.query_store_stats", ct);
+
+ await using var explain = new NpgsqlCommand(
+ "EXPLAIN (ANALYZE, FORMAT JSON, COSTS OFF, TIMING OFF, SUMMARY OFF) " + QueryStoreBackfill.WalkCandidateSql.Replace("$1", Server.ToString(CultureInfo.InvariantCulture), StringComparison.Ordinal),
+ connection);
+ var planJson = (string)(await explain.ExecuteScalarAsync(ct))!;
+ using var plan = JsonDocument.Parse(planJson);
+ var root = plan.RootElement[0].GetProperty("Plan");
+
+ static string NodeType(JsonElement node) => node.GetProperty("Node Type").GetString()!;
+ static IEnumerable Nodes(JsonElement node)
+ {
+ yield return node;
+ if (node.TryGetProperty("Plans", out var children))
+ {
+ foreach (var child in children.EnumerateArray())
+ {
+ foreach (var descendant in Nodes(child))
+ {
+ yield return descendant;
+ }
+ }
+ }
+ }
+
+ var recursiveUnion = Nodes(root).Single(node => NodeType(node) == "Recursive Union");
+ var recursiveTerm = recursiveUnion.GetProperty("Plans").EnumerateArray()
+ .Single(child => child.GetProperty("Parent Relationship").GetString() == "Inner");
+ var innerLimits = Nodes(recursiveTerm).Where(node => NodeType(node) == "Limit").ToList();
+ var innerScans = innerLimits
+ .SelectMany(limit => limit.GetProperty("Plans").EnumerateArray())
+ .Where(child => NodeType(child) is "Index Scan" or "Index Only Scan")
+ .ToList();
+
+ Assert.False(Nodes(root).Any(node => NodeType(node) == "Seq Scan"), "the walk must not scan the table:\n" + planJson);
+ Assert.True(innerScans.Count >= 1,
+ "the recursive term's inner select must be a Limit over an Index (Only) Scan:\n" + planJson);
+ }
+}
diff --git a/Darling/Darling.Tests/QueryStoreBackfillSkipFailingDatabaseTests.cs b/Darling/Darling.Tests/QueryStoreBackfillSkipFailingDatabaseTests.cs
new file mode 100644
index 0000000000..6357d83160
--- /dev/null
+++ b/Darling/Darling.Tests/QueryStoreBackfillSkipFailingDatabaseTests.cs
@@ -0,0 +1,311 @@
+/*
+ * 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.Collections.Generic;
+using System.Globalization;
+using System.Linq;
+using System.Threading;
+using System.Threading.Tasks;
+using Microsoft.Extensions.Logging;
+using Npgsql;
+using PerformanceMonitor.Collectors;
+using PerformanceMonitor.Darling.Service;
+using PerformanceMonitor.Darling.Storage;
+using Xunit;
+
+namespace Darling.Tests;
+
+///
+/// A database whose backfill slices always fail must not stop the databases behind it on the same server.
+/// A slice runs at most once per server per tick and the candidate list comes back in the same order every
+/// tick, so a failing database that is first in line used to be first in line forever: the throw left the
+/// tick before any other database was looked at, and no state was saved to change the next tick's answer.
+///
+/// #1776 own-store — mints its own scratch database (via ) and seeds
+/// two or three databases with a pending first-contact tail. The slice body is replaced through
+/// so no SQL Server is needed and each test scripts
+/// which slices fail; the tick loop, the candidate read, the stored-floor read and the failure accounting are
+/// the real ones.
+///
+public sealed class QueryStoreBackfillSkipFailingDatabaseTests
+{
+ private const int TestServerId = -466201;
+ private const string Failing = "aaa_failing_db";
+ private const string Healthy = "bbb_healthy_db";
+ private const string Third = "ccc_third_db";
+
+ private static int Threshold => QueryStoreBackfillState.SkipAfterConsecutiveSliceFailures;
+
+ /// What one scripted slice call gets: the database, the window the real slice would have used,
+ /// which attempt this is for that database (1-based), and a way to record the drain finishing.
+ private sealed record SliceCall(string Database, TimeSpan Span, int Nth, Func MarkDone);
+
+ [Fact]
+ public async Task AFailingDatabase_DoesNotStopTheDatabaseBehindIt()
+ {
+ var attempts = await RunTicksAsync(
+ [Failing, Healthy], ticks: 6,
+ slice: call => call.Database == Failing ? throw new InvalidOperationException("simulated slice failure") : Task.CompletedTask);
+
+ Assert.True(attempts.Any(a => a.Database == Healthy),
+ "the healthy database was never sliced in 6 ticks; attempts: " + string.Join(", ", attempts.Select(a => a.Database)));
+ }
+
+ [Fact]
+ public async Task ASkippedDatabase_IsRetriedOnceNothingElseHasWork_AtTheServersNarrowedWindow()
+ {
+ var attempts = await RunTicksAsync(
+ [Failing, Healthy], ticks: Threshold + 3,
+ slice: async call =>
+ {
+ if (call.Database == Failing)
+ {
+ throw new InvalidOperationException("simulated slice failure");
+ }
+
+ await call.MarkDone();
+ });
+
+ /* Threshold failures, then the healthy database drains and finishes, then the skipped one is retried
+ on the ticks where nothing else has work. */
+ Assert.Equal(
+ Enumerable.Repeat(Failing, Threshold).Concat([Healthy, Failing, Failing]),
+ attempts.Select(a => a.Database));
+
+ /* The window still narrows per SERVER (#2111), not per database. From a fresh count the failing database
+ at the front is tried at 60, 30 and then 15 minutes and is then skipped. The database behind it runs at
+ the 15 minutes the server has narrowed to, and its completed slice resets the server's count, so the
+ skipped database's retries start at the full 60 minutes again and narrow from there. */
+ Assert.Equal(
+ new[] { 60, 30, 15, 15, 60, 30 }.Select(m => TimeSpan.FromMinutes(m)),
+ attempts.Select(a => a.Span));
+ }
+
+ [Fact]
+ public async Task ACompletedSlice_ClearsTheFailureCount()
+ {
+ /* Fail, fail, complete, then fail forever: the completion must reset the count, so the database needs
+ a full Threshold of fresh failures before it is skipped. Without the reset it is skipped after the
+ first Threshold failures overall and the healthy database runs a few ticks earlier. */
+ var attempts = await RunTicksAsync(
+ [Failing, Healthy], ticks: (2 * Threshold) + 3,
+ slice: async call =>
+ {
+ if (call.Database == Failing && call.Nth != Threshold)
+ {
+ throw new InvalidOperationException("simulated slice failure");
+ }
+
+ if (call.Database == Healthy)
+ {
+ await call.MarkDone();
+ }
+ });
+
+ /* Threshold attempts up to and including the completion, then Threshold fresh failures. Without the
+ reset the healthy database would come one attempt after the completion instead. */
+ var healthyAt = attempts.FindIndex(a => a.Database == Healthy);
+ Assert.Equal(2 * Threshold, healthyAt);
+ }
+
+ [Fact]
+ public async Task SeveralSkippedDatabases_TakeTurnsBeingRetried()
+ {
+ /* Two databases fail forever and a third is healthy. Once both are skipped and the third has drained,
+ each idle tick retries whichever skipped database failed longest ago, so neither starves the other. */
+ var attempts = await RunTicksAsync(
+ [Failing, Healthy, Third], ticks: (3 * Threshold) + 5,
+ slice: async call =>
+ {
+ if (call.Database == Third)
+ {
+ await call.MarkDone();
+ return;
+ }
+
+ throw new InvalidOperationException("simulated slice failure");
+ });
+
+ var expected = Enumerable.Repeat(Failing, Threshold)
+ .Concat(Enumerable.Repeat(Healthy, Threshold))
+ .Append(Third)
+ .Concat([Failing, Healthy, Failing, Healthy]);
+ Assert.Equal(expected, attempts.Select(a => a.Database).Take((2 * Threshold) + 1 + 4));
+ }
+
+ [Fact]
+ public async Task TheSkipWarning_LogsOncePerSkippedDatabase_NotOnEveryTick()
+ {
+ var log = new WarningLog();
+ var attempts = await RunTicksAsync(
+ [Failing, Healthy], ticks: Threshold + 6, logger: log,
+ slice: async call =>
+ {
+ if (call.Database == Failing)
+ {
+ throw new InvalidOperationException("simulated slice failure");
+ }
+
+ await call.MarkDone();
+ });
+
+ Assert.True(attempts.Count(a => a.Database == Failing) > Threshold, "the skipped database should keep being retried");
+ var warning = Assert.Single(log.Warnings, w => w.Contains(Failing, StringComparison.Ordinal));
+ Assert.Contains(Threshold.ToString(CultureInfo.InvariantCulture), warning, StringComparison.Ordinal);
+ Assert.Contains("backfill-skip-test", warning, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void TheSkipThreshold_LeavesOneAttemptAtTheNarrowestAdaptiveSpan_AndNoEarlierOne()
+ {
+ var full = QueryStoreBackfillState.MaxSliceSpan;
+
+ /* From a fresh per-server count, the failing database at the front of the line is tried at
+ AdaptiveSpan(full, k - 1) on attempt k, so the Nth failure is the first one at the floor. The
+ ASkippedDatabase test above pins the same relation through the slice hook's span argument. */
+ Assert.Equal(QueryStoreBackfillState.MinAdaptiveSpan, QueryStoreBackfillState.AdaptiveSpan(full, Threshold - 1));
+ Assert.True(QueryStoreBackfillState.AdaptiveSpan(full, Threshold - 2) > QueryStoreBackfillState.MinAdaptiveSpan,
+ "a smaller threshold would skip a database before it ever tried the narrowest window");
+ }
+
+ [Fact]
+ public void TheLedger_SkipsAtTheThreshold_ResetsOnCompletion_AndOrdersRetriesLeastRecentFirst()
+ {
+ var ledger = new QueryStoreBackfillFailureLedger();
+ for (var failures = 1; failures < Threshold; failures++)
+ {
+ Assert.Equal(failures, ledger.RecordFailure(1, "a"));
+ Assert.False(ledger.IsSkipped(1, "a"));
+ }
+
+ Assert.Equal(Threshold, ledger.RecordFailure(1, "a"));
+ Assert.True(ledger.IsSkipped(1, "a"));
+ Assert.False(ledger.IsSkipped(2, "a"));
+ Assert.False(ledger.IsSkipped(1, "b"));
+
+ ledger.RecordCompletion(1, "a");
+ Assert.Equal(0, ledger.Failures(1, "a"));
+ Assert.False(ledger.IsSkipped(1, "a"));
+
+ ledger.RecordFailure(1, "a");
+ ledger.RecordFailure(1, "b");
+ Assert.True(ledger.LastFailureTicket(1, "a") < ledger.LastFailureTicket(1, "b"));
+ Assert.Equal(0, ledger.LastFailureTicket(1, "never-failed"));
+ }
+
+ /// Seeds one pending first-contact tail per database, replaces the slice body with
+ /// , runs real ticks (a slice that throws is caught the way
+ /// the worker's outer catch does) and returns every slice attempt in order.
+ private static async Task> RunTicksAsync(
+ string[] databases, int ticks, Func slice, ILogger? logger = null)
+ {
+ var baseConnectionString = Environment.GetEnvironmentVariable("DARLING_TEST_PG");
+ Assert.SkipWhen(string.IsNullOrEmpty(baseConnectionString),
+ "Set DARLING_TEST_PG to a Postgres connection string to run the live Query Store backfill skip tests (they mint their own scratch database).");
+
+ 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 < databases.Length; i++)
+ {
+ await SeedPendingTailAsync(connection, 466201L + i, databases[i], ct);
+ }
+ }
+
+ await using var postgres = NpgsqlDataSource.Create(scratch.ConnectionString);
+ var runner = new DarlingCollectorRunner(postgres, new CollectorDeltaCalculator());
+ var backfill = new QueryStoreBackfill(postgres, runner, new CollectorDeltaCalculator(), logger);
+ var attempts = new List();
+ backfill.SliceOverrideForTests = async (database, span) =>
+ {
+ attempts.Add(new SliceAttempt(database, span));
+ var nth = attempts.Count(a => a.Database == database);
+ await slice(new SliceCall(database, span, nth, () => runner.SaveCollectorStateAsync(
+ TestServerId,
+ QueryStoreBackfill.StateCollectorName,
+ new Dictionary(StringComparer.Ordinal)
+ {
+ [QueryStoreBackfillState.DoneKeyPrefix + database] = DateTime.UtcNow.ToString("o", CultureInfo.InvariantCulture),
+ },
+ ct)));
+ };
+
+ var server = new ServerRuntime
+ {
+ Config = new MonitoredServer { Name = "backfill-skip-test", Host = "backfill-skip-test" },
+ ConnectionString = "Server=backfill-skip-test",
+ Target = new CollectorTargetInfo { SqlMajorVersion = 16 },
+ StorageName = "backfill-skip-test",
+ ServerId = TestServerId,
+ EngineEdition = 3,
+ };
+
+ for (var tick = 0; tick < ticks; tick++)
+ {
+ try
+ {
+ await backfill.RunServerSliceAsync(server, ct);
+ }
+ catch (InvalidOperationException ex) when (ex.Message.StartsWith("simulated", StringComparison.Ordinal))
+ {
+ /* The worker's outer catch logs a failed slice and carries on to the next tick. */
+ }
+ }
+
+ return attempts;
+ }
+
+ private sealed record SliceAttempt(string Database, TimeSpan Span);
+
+ /// One row an hour old with an even older last_execution_time: inside the horizon, so the database
+ /// is a candidate, and newer than the floor, so its stored ceiling is above it and a tail slice is pending.
+ /// Relative to now because the tick computes its floor from the wall clock.
+ private static async Task SeedPendingTailAsync(NpgsqlConnection connection, long collectionId, string databaseName, CancellationToken ct)
+ {
+ const string sql = @"
+INSERT INTO collect.query_store_stats
+ (collection_id, collection_time, server_id, server_name, database_name, module_name, query_hash,
+ query_id, plan_id, execution_type_desc, replica_role,
+ runtime_stats_interval_id, interval_start_time_utc, first_execution_time, last_execution_time,
+ execution_count, avg_duration_us, avg_cpu_time_us, min_duration_us, max_duration_us)
+VALUES
+ ($1, $2, $3, 'SQL01', $4, 'dbo.GetOrders', '0xABCD', 91, 111, 'Regular', 'Primary',
+ 1, $2, $5, $5, 1, 100, 100, 100, 100)";
+
+ var now = DateTime.UtcNow;
+ await using var command = new NpgsqlCommand(sql, connection);
+ command.Parameters.AddWithValue(collectionId);
+ command.Parameters.AddWithValue(DateTime.SpecifyKind(now.AddHours(-1), DateTimeKind.Unspecified));
+ command.Parameters.AddWithValue(TestServerId);
+ command.Parameters.AddWithValue(databaseName);
+ command.Parameters.AddWithValue(DateTime.SpecifyKind(now.AddHours(-2), DateTimeKind.Unspecified));
+ await command.ExecuteNonQueryAsync(ct);
+ }
+
+ /// Keeps the Warning-and-above messages, formatted, so a test can count them.
+ private sealed class WarningLog : ILogger
+ {
+ public List Warnings { get; } = [];
+
+ public IDisposable? BeginScope(TState state) where TState : notnull => null;
+
+ public bool IsEnabled(LogLevel logLevel) => true;
+
+ public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func formatter)
+ {
+ if (logLevel >= LogLevel.Warning)
+ {
+ Warnings.Add(formatter(state, exception));
+ }
+ }
+ }
+}
diff --git a/Darling/Darling.Tests/QueryStoreBackfillTests.cs b/Darling/Darling.Tests/QueryStoreBackfillTests.cs
index 2cc95b5cf2..c35e845a53 100644
--- a/Darling/Darling.Tests/QueryStoreBackfillTests.cs
+++ b/Darling/Darling.Tests/QueryStoreBackfillTests.cs
@@ -290,6 +290,33 @@ public void Sql_LiteTwinCarriesTheSameBoundShapes()
Assert.DoesNotContain(
"$\"SELECT MIN({columnName}) FROM {tableName} WHERE server_id = $1 AND {databaseColumnName} = $2\"",
source, StringComparison.Ordinal);
+
+ /* #4662: Darling's store switch (a cut-chunk read on TimescaleDB, a walk on plain PostgreSQL) does not
+ apply to DuckDB, which has no chunks. Lite's candidate statement is Darling's fallback statement,
+ character for character - referenced, not copied. */
+ Assert.Contains("\"" + global::PerformanceMonitor.Darling.Service.QueryStoreBackfill.CandidateSql + "\"", source, StringComparison.Ordinal);
+ }
+
+ /// #4662: the three statements the store switch runs. The cut-chunk read is bound at the chunk START
+ /// (a >= against the catalog's value) rather than > floor, which inside the chunk that holds
+ /// the floor can only filter; the catalog read is qualified to the hypertable and applies UTC on both sides; the
+ /// walk carries no time bound at all and never lists a NULL name.
+ [Fact]
+ public void Sql_StoreSwitchStatements_AreBoundAtTheChunkStart_ReadTheCatalogInUtc_AndWalkWithoutATimeBound()
+ {
+ const string cutChunk = global::PerformanceMonitor.Darling.Service.QueryStoreBackfill.CutChunkCandidateSql;
+ const string catalog = global::PerformanceMonitor.Darling.Service.QueryStoreBackfill.CutChunkCatalogSql;
+ const string walk = global::PerformanceMonitor.Darling.Service.QueryStoreBackfill.WalkCandidateSql;
+
+ Assert.Contains("collection_time >= TIMESTAMP '{cut_start}'", cutChunk, StringComparison.Ordinal);
+ Assert.DoesNotContain("collection_time >", cutChunk.Replace("collection_time >=", "", StringComparison.Ordinal), StringComparison.Ordinal);
+ Assert.Contains("hypertable_schema = 'collect' AND hypertable_name = 'query_store_stats'", catalog, StringComparison.Ordinal);
+ Assert.Contains("range_start AT TIME ZONE 'UTC' AS cut_start", catalog, StringComparison.Ordinal);
+ Assert.Contains("range_start <= TIMESTAMP '{floor}' AT TIME ZONE 'UTC'", catalog, StringComparison.Ordinal);
+ Assert.Contains("range_end > TIMESTAMP '{floor}' AT TIME ZONE 'UTC'", catalog, StringComparison.Ordinal);
+ Assert.StartsWith("WITH RECURSIVE walk AS", walk, StringComparison.Ordinal);
+ Assert.DoesNotContain("collection_time", walk, StringComparison.Ordinal);
+ Assert.Contains("database_name IS NOT NULL", walk, StringComparison.Ordinal);
}
}
diff --git a/Darling/Darling.Tests/WaitRateTileReadCountLiveTests.cs b/Darling/Darling.Tests/WaitRateTileReadCountLiveTests.cs
index 9adfb9c1f6..3705215b87 100644
--- a/Darling/Darling.Tests/WaitRateTileReadCountLiveTests.cs
+++ b/Darling/Darling.Tests/WaitRateTileReadCountLiveTests.cs
@@ -8,7 +8,6 @@
using System;
using System.Diagnostics;
-using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Npgsql;
@@ -152,54 +151,4 @@ private static async Task InsertAsync(NpgsqlConnection connection, string sql, p
command.Parameters.AddWithValue(value);
await command.ExecuteNonQueryAsync(TestContext.Current.CancellationToken);
}
-
- ///
- /// Counts the commands Npgsql runs against ONE database whose text contains every one of a set of fragments.
- /// Npgsql opens an per command on its Npgsql once a
- /// listener samples it, tagged with the database name and the command text; the count is taken when the
- /// activity stops, which is when the command's reader closes. The tags are matched by VALUE rather than by
- /// name (the same way SharedBaselineCacheTests.CaptureAsync finds the command text), because the
- /// OpenTelemetry attribute names Npgsql uses have changed between major versions.
- ///
- private sealed class NpgsqlCommandCounter : IDisposable
- {
- private readonly string _databaseName;
- private readonly string[] _fragments;
- private readonly ActivityListener _listener;
- private long _count;
-
- public NpgsqlCommandCounter(string databaseName, params string[] fragments)
- {
- _databaseName = databaseName;
- _fragments = fragments;
- _listener = new ActivityListener
- {
- ShouldListenTo = source => source.Name == "Npgsql",
- Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded,
- ActivityStopped = OnStopped,
- };
- ActivitySource.AddActivityListener(_listener);
- }
-
- public long Count => Interlocked.Read(ref _count);
-
- private void OnStopped(Activity activity)
- {
- var onThisDatabase = false;
- var carriesTheText = false;
- foreach (var tag in activity.TagObjects)
- {
- if (tag.Value is not string value)
- continue;
- if (string.Equals(value, _databaseName, StringComparison.Ordinal))
- onThisDatabase = true;
- else if (_fragments.All(fragment => value.Contains(fragment, StringComparison.Ordinal)))
- carriesTheText = true;
- }
- if (onThisDatabase && carriesTheText)
- Interlocked.Increment(ref _count);
- }
-
- public void Dispose() => _listener.Dispose();
- }
}
diff --git a/Darling/PerformanceMonitor.Darling.Service/QueryStoreBackfill.cs b/Darling/PerformanceMonitor.Darling.Service/QueryStoreBackfill.cs
index 8c967cbb54..c7d945fb37 100644
--- a/Darling/PerformanceMonitor.Darling.Service/QueryStoreBackfill.cs
+++ b/Darling/PerformanceMonitor.Darling.Service/QueryStoreBackfill.cs
@@ -174,7 +174,9 @@ public QueryStoreBackfill(
///
/// Runs AT MOST one backfill slice for one server: the first database found with a pending
/// hole or an undrained first-contact tail gets one byte-budgeted slice; everything else waits
- /// for a later tick. Returns true when a slice (or an exhaustion probe) ran, false when the
+ /// for a later tick. A database that has failed
+ /// slices in a row is skipped
+ /// while any other database has work, then retried on a tick where none does. Returns true when a slice (or an exhaustion probe) ran, false when the
/// server had no backfill work — the common steady state, costing one candidate query and a
/// few MIN() lookups.
///
@@ -216,6 +218,10 @@ and this is the designed response to it. */
QueryStoreBackfillState.MergeHoleDatabases for why the union is required, not just cheaper. */
var databases = await GetCandidateDatabasesAsync(server.ServerId, floorLimit, state, cancellationToken);
+ /* Databases whose slices keep failing: their slice is held back while any other database has work
+ (see QueryStoreBackfillFailureLedger), then one of them is retried after the walk. */
+ List<(string Database, DateTime Floor, DateTime Ceiling, bool IsHole)>? skipped = null;
+
foreach (var databaseName in databases)
{
cancellationToken.ThrowIfCancellationRequested();
@@ -233,6 +239,12 @@ and this is the designed response to it. */
}
var holeFloor = holeFrom > floorLimit ? holeFrom : floorLimit;
+ if (_sliceFailures.IsSkipped(server.ServerId, databaseName))
+ {
+ (skipped ??= []).Add((databaseName, holeFloor, holeTo, true));
+ continue;
+ }
+
await RunCountedSliceAsync(server, databaseName, holeFloor, holeTo, isHole: true, cancellationToken);
return true;
}
@@ -259,10 +271,35 @@ without shipping a row so the steady state never re-probes it. */
continue;
}
+ if (_sliceFailures.IsSkipped(server.ServerId, databaseName))
+ {
+ (skipped ??= []).Add((databaseName, floorLimit, storedFloor.Value, false));
+ continue;
+ }
+
await RunCountedSliceAsync(server, databaseName, floorLimit, storedFloor.Value, isHole: false, cancellationToken);
return true;
}
+ /* No other database had work, so retry a skipped one: the one whose last failure is the oldest, so
+ several skipped databases take turns instead of the first in the list starving the rest. This costs
+ at most one failed slice per tick on an otherwise idle server, exactly what the stall cost before. */
+ if (skipped is { Count: > 0 })
+ {
+ var retry = skipped[0];
+ for (var i = 1; i < skipped.Count; i++)
+ {
+ if (_sliceFailures.LastFailureTicket(server.ServerId, skipped[i].Database)
+ < _sliceFailures.LastFailureTicket(server.ServerId, retry.Database))
+ {
+ retry = skipped[i];
+ }
+ }
+
+ await RunCountedSliceAsync(server, retry.Database, retry.Floor, retry.Ceiling, retry.IsHole, cancellationToken);
+ return true;
+ }
+
return false;
}
@@ -276,6 +313,22 @@ without shipping a row so the steady state never re-probes it. */
///
private readonly ConcurrentDictionary _consecutiveSliceFailures = new();
+ ///
+ /// Consecutive failed slices per (server, database), used ONLY to decide which database to skip: one that
+ /// fails slices in a row is served
+ /// after the databases behind it instead of ahead of them, so it can no longer stall them. Reset by that
+ /// database's completed slice. It does not size the slice window: that stays the per-server count
+ /// above, because a command timeout usually means the whole server is loaded, and narrowing per database
+ /// would add timed-out queries per database against a server that is already struggling. In memory on
+ /// purpose, like the live counters.
+ ///
+ private readonly QueryStoreBackfillFailureLedger _sliceFailures = new();
+
+ /// Test-only seam: when set, replaces the slice body (called with the database and the
+ /// window span the slice would have used). A throw counts as a failed slice and a normal return as a
+ /// completed one, through the same accounting. Null in production, where it changes nothing.
+ internal Func? SliceOverrideForTests { get; set; }
+
/// Runs one slice with the failure accounting wrapped around it — the worker's outer
/// catch still logs the throw exactly as before.
private async Task RunCountedSliceAsync(
@@ -285,10 +338,22 @@ private async Task RunCountedSliceAsync(
{
await RunSliceAsync(server, databaseName, floorUtc, ceilingUtc, isHole, cancellationToken);
_consecutiveSliceFailures.TryRemove(server.ServerId, out _);
+ _sliceFailures.RecordCompletion(server.ServerId, databaseName);
}
catch (Exception ex) when (ex is not OperationCanceledException)
{
_consecutiveSliceFailures.AddOrUpdate(server.ServerId, 1, static (_, count) => count + 1);
+ var failures = _sliceFailures.RecordFailure(server.ServerId, databaseName);
+
+ /* Logged at the failure that crosses the threshold, so it is once per stretch of failures and
+ needs no extra state: the count only grows until a completed slice clears it. */
+ if (failures == QueryStoreBackfillState.SkipAfterConsecutiveSliceFailures)
+ {
+ _logger?.LogWarning(
+ "query_store backfill on '{Server}' [{Database}]: {Failures} consecutive slice failures; serving the other databases first and retrying this one only when none has work.",
+ server.Config.DisplayName, databaseName, failures);
+ }
+
throw;
}
}
@@ -315,6 +380,12 @@ window on a big database times out at the command timeout every tick and the ran
_consecutiveSliceFailures.TryGetValue(server.ServerId, out var recentFailures) ? recentFailures : 0);
var sliceFloor = QueryStoreBackfillState.BoundSliceFloor(floorUtc, ceilingUtc, sliceSpan);
+ if (SliceOverrideForTests is { } sliceOverride)
+ {
+ await sliceOverride(databaseName, sliceSpan);
+ return;
+ }
+
var definition = QueryStoreCollector.Instance;
var context = new CollectorContext
{
@@ -463,6 +534,59 @@ private Task SaveStateAsync(int serverId, string key, string value, Cancellation
internal const string CandidateSql =
"SELECT DISTINCT database_name FROM query_store_stats WHERE server_id = $1 AND collection_time > $2 ORDER BY database_name";
+ /// #4662, TimescaleDB store, step 1: the chunk of collect.query_store_stats that holds the
+ /// floor, read from the catalog at run time. NEVER derived from a constant: set_chunk_time_interval
+ /// changes only chunks created after it, so a chunk width is wrong for every older chunk. {floor} is
+ /// replaced by the floor formatted yyyy-MM-dd HH:mm:ss.ffffff (invariant culture) - a literal, like
+ /// the statement below, so chunk exclusion happens at plan time. No row means no chunk holds the floor.
+ internal const string CutChunkCatalogSql =
+ "SELECT range_start AT TIME ZONE 'UTC' AS cut_start, range_end AT TIME ZONE 'UTC' AS cut_end FROM timescaledb_information.chunks WHERE hypertable_schema = 'collect' AND hypertable_name = 'query_store_stats' AND range_start <= TIMESTAMP '{floor}' AT TIME ZONE 'UTC' AND range_end > TIMESTAMP '{floor}' AT TIME ZONE 'UTC'";
+
+ /// #4662, TimescaleDB store, step 2: the database list from the cut chunk's start. Inside the chunk
+ /// that holds the floor a collection_time > floor bound can only FILTER (it is the fifth key column
+ /// of the index), so the old read walked the chunk's index entry by entry; a bound at the chunk start lets
+ /// the skip scan seek once per name. {cut_start} is the catalog's cut_start, formatted as above.
+ internal const string CutChunkCandidateSql =
+ "SELECT DISTINCT database_name FROM query_store_stats WHERE server_id = $1 AND collection_time >= TIMESTAMP '{cut_start}' ORDER BY database_name";
+
+ /// #4662, plain PostgreSQL store: no chunks, so no time bound - a walk down the index, one seek per
+ /// database. Every name stored for the server is listed, and each is marked done once.
+ internal const string WalkCandidateSql =
+ "WITH RECURSIVE walk AS ((SELECT s.database_name FROM query_store_stats AS s WHERE s.server_id = $1 AND s.database_name IS NOT NULL ORDER BY s.database_name LIMIT 1) UNION ALL SELECT (SELECT s.database_name FROM query_store_stats AS s WHERE s.server_id = $1 AND s.database_name > w.database_name ORDER BY s.database_name LIMIT 1) FROM walk AS w WHERE w.database_name IS NOT NULL) SELECT database_name FROM walk WHERE database_name IS NOT NULL";
+
+ private static string TimestampLiteral(DateTime value)
+ => value.ToString("yyyy-MM-dd HH:mm:ss.ffffff", CultureInfo.InvariantCulture);
+
+ /// Picks the candidate statement for this store (#4662). Plain PostgreSQL: the walk. TimescaleDB:
+ /// the cut-chunk read when the catalog names a chunk that holds the floor; otherwise (no such chunk, or the
+ /// catalog read failed) , unchanged. A catalog failure never reaches the caller's
+ /// catch, which would drop the whole store list.
+ private async Task ChooseCandidateSqlAsync(NpgsqlConnection connection, DateTime floorLimit, CancellationToken cancellationToken)
+ {
+ if (!_hasContinuousAggregates())
+ {
+ return WalkCandidateSql;
+ }
+
+ try
+ {
+ using var catalog = new NpgsqlCommand(
+ CutChunkCatalogSql.Replace("{floor}", TimestampLiteral(floorLimit), StringComparison.Ordinal), connection);
+ catalog.CommandTimeout = ServiceCommandDeadlines.QueryStoreBackfillReadSeconds;
+ await using var catalogReader = await catalog.ExecuteReaderAsync(cancellationToken);
+ if (await catalogReader.ReadAsync(cancellationToken) && !catalogReader.IsDBNull(0))
+ {
+ return CutChunkCandidateSql.Replace("{cut_start}", TimestampLiteral(catalogReader.GetDateTime(0)), StringComparison.Ordinal);
+ }
+ }
+ catch (Exception ex) when (ex is not OperationCanceledException)
+ {
+ _logger?.LogDebug(ex, "query_store backfill cut-chunk catalog read failed; using the floor-bound candidate read");
+ }
+
+ return CandidateSql;
+ }
+
/// Databases that shipped query_store rows since , unioned
/// with every database a hole key already names — the backfill universe.
///
@@ -485,12 +609,18 @@ internal async Task> GetCandidateDatabasesAsync(
try
{
await using var connection = await _postgres.OpenConnectionAsync(cancellationToken);
- using var command = new NpgsqlCommand(CandidateSql, connection);
+ var sql = await ChooseCandidateSqlAsync(connection, floorLimit, cancellationToken);
+ using var command = new NpgsqlCommand(sql, connection);
/* #2874: the enclosing BackfillSliceDeadline ABANDONS rather than cancels, so this is the only
bound that reaches the statement. */
command.CommandTimeout = ServiceCommandDeadlines.QueryStoreBackfillReadSeconds;
command.Parameters.AddWithValue(serverId);
- command.Parameters.AddWithValue(DateTime.SpecifyKind(floorLimit, DateTimeKind.Unspecified));
+ if (string.Equals(sql, CandidateSql, StringComparison.Ordinal))
+ {
+ /* Only the floor-bound statement takes the floor as a parameter; the cut-chunk read and the
+ walk carry no second placeholder, and a surplus parameter is a bind error. */
+ command.Parameters.AddWithValue(DateTime.SpecifyKind(floorLimit, DateTimeKind.Unspecified));
+ }
await using var reader = await command.ExecuteReaderAsync(cancellationToken);
while (await reader.ReadAsync(cancellationToken))
{
diff --git a/Lite.Tests/QueryStoreBackfillSkipFailingDatabaseTests.cs b/Lite.Tests/QueryStoreBackfillSkipFailingDatabaseTests.cs
new file mode 100644
index 0000000000..4ef8c20d04
--- /dev/null
+++ b/Lite.Tests/QueryStoreBackfillSkipFailingDatabaseTests.cs
@@ -0,0 +1,301 @@
+/*
+ * Copyright (c) 2026 Erik Darling, Darling Data LLC
+ *
+ * This file is part of the SQL Server Performance Monitor Lite.
+ *
+ * Licensed under the MIT License. See LICENSE file in the project root for full license information.
+ */
+
+using System;
+using System.Collections.Generic;
+using System.Globalization;
+using System.IO;
+using System.Linq;
+using System.Threading;
+using System.Threading.Tasks;
+using DuckDB.NET.Data;
+using Microsoft.Extensions.Logging;
+using PerformanceMonitor.Collectors;
+using PerformanceMonitorLite.Database;
+using PerformanceMonitorLite.Models;
+using PerformanceMonitorLite.Services;
+using PerformanceMonitorLite.Tests;
+using Xunit;
+
+namespace Lite.Tests;
+
+///
+/// The Lite twin of Darling's skip-a-failing-database pins. A database whose backfill slices always fail must
+/// not stop the databases behind it on the same server: a slice runs at most once per server per tick and the
+/// candidate list comes back in the same order every tick, so a failing database first in line used to be first
+/// in line forever. The slice body is replaced through
+/// so no SQL Server is needed; the tick loop, the DuckDB candidate and stored-floor reads and the failure
+/// accounting are the real ones.
+///
+public sealed class QueryStoreBackfillSkipFailingDatabaseTests : IClassFixture
+{
+ private const string Failing = "aaa_failing_db";
+ private const string Healthy = "bbb_healthy_db";
+ private const string Third = "ccc_third_db";
+ private const string ServerLabel = "backfill-skip-test";
+
+ private readonly DuckDbInitializer _duckDb;
+
+ public QueryStoreBackfillSkipFailingDatabaseTests(SharedDuckDbFixture fixture)
+ {
+ fixture.ResetData();
+ _duckDb = fixture.DuckDb;
+ }
+
+ private static int Threshold => QueryStoreBackfillState.SkipAfterConsecutiveSliceFailures;
+
+ private sealed record SliceCall(string Database, TimeSpan Span, int Nth, Func MarkDone);
+
+ private sealed record SliceAttempt(string Database, TimeSpan Span);
+
+ [Fact]
+ public async Task AFailingDatabase_DoesNotStopTheDatabaseBehindIt()
+ {
+ var (attempts, _) = await RunTicksAsync(
+ [Failing, Healthy], ticks: 6,
+ slice: call => call.Database == Failing ? throw new InvalidOperationException("simulated slice failure") : Task.CompletedTask);
+
+ Assert.True(attempts.Any(a => a.Database == Healthy),
+ "the healthy database was never sliced in 6 ticks; attempts: " + string.Join(", ", attempts.Select(a => a.Database)));
+ }
+
+ [Fact]
+ public async Task ASkippedDatabase_IsRetriedOnceNothingElseHasWork_AtTheServersNarrowedWindow()
+ {
+ var (attempts, _) = await RunTicksAsync(
+ [Failing, Healthy], ticks: Threshold + 3,
+ slice: async call =>
+ {
+ if (call.Database == Failing)
+ {
+ throw new InvalidOperationException("simulated slice failure");
+ }
+
+ await call.MarkDone();
+ });
+
+ Assert.Equal(
+ Enumerable.Repeat(Failing, Threshold).Concat([Healthy, Failing, Failing]),
+ attempts.Select(a => a.Database));
+
+ /* The window still narrows per SERVER (#2111), not per database. From a fresh count the failing database
+ at the front is tried at 60, 30 and then 15 minutes and is then skipped. The database behind it runs at
+ the 15 minutes the server has narrowed to, and its completed slice resets the server's count, so the
+ skipped database's retries start at the full 60 minutes again and narrow from there. */
+ Assert.Equal(
+ new[] { 60, 30, 15, 15, 60, 30 }.Select(m => TimeSpan.FromMinutes(m)),
+ attempts.Select(a => a.Span));
+ }
+
+ [Fact]
+ public async Task ACompletedSlice_ClearsTheFailureCount()
+ {
+ /* Fail, fail, complete, then fail forever: the completion must reset the count, so the database needs
+ a full Threshold of fresh failures before it is skipped. Without the reset the healthy database
+ would come one attempt after the completion instead. */
+ var (attempts, _) = await RunTicksAsync(
+ [Failing, Healthy], ticks: (2 * Threshold) + 3,
+ slice: async call =>
+ {
+ if (call.Database == Failing && call.Nth != Threshold)
+ {
+ throw new InvalidOperationException("simulated slice failure");
+ }
+
+ if (call.Database == Healthy)
+ {
+ await call.MarkDone();
+ }
+ });
+
+ Assert.Equal(2 * Threshold, attempts.FindIndex(a => a.Database == Healthy));
+ }
+
+ [Fact]
+ public async Task SeveralSkippedDatabases_TakeTurnsBeingRetried()
+ {
+ var (attempts, _) = await RunTicksAsync(
+ [Failing, Healthy, Third], ticks: (3 * Threshold) + 5,
+ slice: async call =>
+ {
+ if (call.Database == Third)
+ {
+ await call.MarkDone();
+ return;
+ }
+
+ throw new InvalidOperationException("simulated slice failure");
+ });
+
+ var expected = Enumerable.Repeat(Failing, Threshold)
+ .Concat(Enumerable.Repeat(Healthy, Threshold))
+ .Append(Third)
+ .Concat([Failing, Healthy, Failing, Healthy]);
+ Assert.Equal(expected, attempts.Select(a => a.Database).Take((2 * Threshold) + 1 + 4));
+ }
+
+ [Fact]
+ public async Task TheSkipWarning_LogsOncePerSkippedDatabase_NotOnEveryTick()
+ {
+ var (attempts, log) = await RunTicksAsync(
+ [Failing, Healthy], ticks: Threshold + 6,
+ slice: async call =>
+ {
+ if (call.Database == Failing)
+ {
+ throw new InvalidOperationException("simulated slice failure");
+ }
+
+ await call.MarkDone();
+ });
+
+ Assert.True(attempts.Count(a => a.Database == Failing) > Threshold, "the skipped database should keep being retried");
+ var warning = Assert.Single(log.Warnings, w => w.Contains(Failing, StringComparison.Ordinal));
+ Assert.Contains(Threshold.ToString(CultureInfo.InvariantCulture), warning, StringComparison.Ordinal);
+ Assert.Contains(ServerLabel, warning, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void TheSkipThreshold_LeavesOneAttemptAtTheNarrowestAdaptiveSpan_AndNoEarlierOne()
+ {
+ var full = QueryStoreBackfillState.MaxSliceSpan;
+
+ /* From a fresh per-server count, the failing database at the front of the line is tried at
+ AdaptiveSpan(full, k - 1) on attempt k. The ASkippedDatabase test above pins the same relation
+ through the slice hook's span argument. */
+ Assert.Equal(QueryStoreBackfillState.MinAdaptiveSpan, QueryStoreBackfillState.AdaptiveSpan(full, Threshold - 1));
+ Assert.True(QueryStoreBackfillState.AdaptiveSpan(full, Threshold - 2) > QueryStoreBackfillState.MinAdaptiveSpan,
+ "a smaller threshold would skip a database before it ever tried the narrowest window");
+ }
+
+ private sealed class Harness(DuckDbInitializer duckDb, ServerManager servers, ScheduleManager schedules, ILogger logger)
+ : RemoteCollectorService(duckDb, servers, schedules, logger)
+ {
+ public Task TickAsync(ServerConnection server) => RunQueryStoreBackfillSliceAsync(server, CancellationToken.None);
+
+ public Task MarkDoneAsync(int serverId, string database) => SaveCollectorStateAsync(
+ serverId,
+ QueryStoreBackfillState.StateCollectorName,
+ new Dictionary(StringComparer.Ordinal)
+ {
+ [QueryStoreBackfillState.DoneKeyPrefix + database] = DateTime.UtcNow.ToString("o", CultureInfo.InvariantCulture),
+ },
+ CancellationToken.None);
+ }
+
+ /// Seeds one pending first-contact tail per database, replaces the slice body with
+ /// , runs real ticks (a slice that throws is caught the way
+ /// the collection loop's per-server catch does) and returns every slice attempt in order.
+ private async Task<(List Attempts, WarningLog Log)> RunTicksAsync(
+ string[] databases, int ticks, Func slice)
+ {
+ var configDirectory = Path.Combine(Path.GetTempPath(), "qs-backfill-skip-" + Guid.NewGuid().ToString("N"));
+ Directory.CreateDirectory(configDirectory);
+ try
+ {
+ var log = new WarningLog();
+ var harness = new Harness(_duckDb, new ServerManager(configDirectory), new ScheduleManager(configDirectory), log);
+ var server = new ServerConnection { ServerName = ServerLabel, DisplayName = ServerLabel };
+ var serverId = RemoteCollectorService.GetDeterministicHashCode(RemoteCollectorService.GetServerNameForStorage(server));
+
+ for (var i = 0; i < databases.Length; i++)
+ {
+ await SeedPendingTailAsync(_duckDb, -466201L - i, serverId, databases[i]);
+ }
+
+ var attempts = new List();
+ harness.SliceOverrideForTests = async (database, span) =>
+ {
+ attempts.Add(new SliceAttempt(database, span));
+ var nth = attempts.Count(a => a.Database == database);
+ await slice(new SliceCall(database, span, nth, () => harness.MarkDoneAsync(serverId, database)));
+ };
+
+ for (var tick = 0; tick < ticks; tick++)
+ {
+ try
+ {
+ await harness.TickAsync(server);
+ }
+ catch (InvalidOperationException ex) when (ex.Message.StartsWith("simulated", StringComparison.Ordinal))
+ {
+ /* The collection loop's per-server catch logs a failed slice and carries on to the next tick. */
+ }
+ }
+
+ return (attempts, log);
+ }
+ finally
+ {
+ try
+ {
+ Directory.Delete(configDirectory, recursive: true);
+ }
+ catch (IOException)
+ {
+ /* Best effort: a leftover temp directory is harmless. */
+ }
+ }
+ }
+
+ /// One row an hour old with an even older last_execution_time: inside the horizon, so the database is
+ /// a candidate, and newer than the floor, so its stored ceiling is above it and a tail slice is pending.
+ /// Relative to now because the tick computes its floor from the wall clock.
+ private static async Task SeedPendingTailAsync(DuckDbInitializer duckDb, long collectionId, int serverId, string databaseName)
+ {
+ var now = DateTime.UtcNow;
+ using var readLock = duckDb.AcquireReadLock();
+ using var connection = duckDb.CreateConnection();
+ await connection.OpenAsync();
+ using var cmd = connection.CreateCommand();
+ cmd.CommandText = @"
+INSERT INTO query_store_stats
+ (collection_id, collection_time, server_id, server_name, database_name,
+ query_id, plan_id, execution_type_desc, first_execution_time, last_execution_time,
+ query_text, query_hash, execution_count, avg_cpu_time_us, avg_duration_us,
+ query_plan_hash, is_forced_plan, force_failure_count)
+VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18)";
+ cmd.Parameters.Add(new DuckDBParameter { Value = collectionId });
+ cmd.Parameters.Add(new DuckDBParameter { Value = now.AddHours(-1) });
+ cmd.Parameters.Add(new DuckDBParameter { Value = serverId });
+ cmd.Parameters.Add(new DuckDBParameter { Value = ServerLabel });
+ cmd.Parameters.Add(new DuckDBParameter { Value = databaseName });
+ cmd.Parameters.Add(new DuckDBParameter { Value = collectionId });
+ cmd.Parameters.Add(new DuckDBParameter { Value = 1L });
+ cmd.Parameters.Add(new DuckDBParameter { Value = "Regular" });
+ cmd.Parameters.Add(new DuckDBParameter { Value = now.AddHours(-2) });
+ cmd.Parameters.Add(new DuckDBParameter { Value = now.AddHours(-2) });
+ cmd.Parameters.Add(new DuckDBParameter { Value = "SELECT 1" });
+ cmd.Parameters.Add(new DuckDBParameter { Value = "0xTESTHASH" });
+ cmd.Parameters.Add(new DuckDBParameter { Value = 10L });
+ cmd.Parameters.Add(new DuckDBParameter { Value = 1000L });
+ cmd.Parameters.Add(new DuckDBParameter { Value = 2000L });
+ cmd.Parameters.Add(new DuckDBParameter { Value = "0xTESTPLANHASH" });
+ cmd.Parameters.Add(new DuckDBParameter { Value = false });
+ cmd.Parameters.Add(new DuckDBParameter { Value = 0L });
+ await cmd.ExecuteNonQueryAsync();
+ }
+
+ /// Keeps the Warning-and-above messages, formatted, so a test can count them.
+ private sealed class WarningLog : ILogger
+ {
+ public List Warnings { get; } = [];
+
+ public IDisposable? BeginScope(TState state) where TState : notnull => null;
+
+ public bool IsEnabled(LogLevel logLevel) => true;
+
+ public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func formatter)
+ {
+ if (logLevel >= LogLevel.Warning)
+ {
+ Warnings.Add(formatter(state, exception));
+ }
+ }
+ }
+}
diff --git a/Lite/Services/RemoteCollectorService.QueryStoreBackfill.cs b/Lite/Services/RemoteCollectorService.QueryStoreBackfill.cs
index 3adb365f89..15015fc624 100644
--- a/Lite/Services/RemoteCollectorService.QueryStoreBackfill.cs
+++ b/Lite/Services/RemoteCollectorService.QueryStoreBackfill.cs
@@ -68,7 +68,9 @@ collection paths.
///
/// Runs AT MOST one backfill slice per enabled server: the first database found with a pending
/// hole or an undrained first-contact tail gets one byte-budgeted slice; everything else waits
- /// for a later tick. Per-server failures log and skip, and per-server WEDGES are abandoned and
+ /// for a later tick. A database that has failed
+ /// slices in a row is skipped
+ /// while any other database has work, then retried on a tick where none does. Per-server failures log and skip, and per-server WEDGES are abandoned and
/// quarantined (#2148) — one stuck server never stalls the sweep in either failure mode. Called
/// from CollectionBackgroundService on its own due-cadence.
///
@@ -168,6 +170,22 @@ private void OnQueryStoreItemSucceeded(int serverId, string database)
/// any completed slice resets it.
private readonly ConcurrentDictionary _consecutiveSliceFailures = new();
+ ///
+ /// Consecutive failed backfill slices per (server, database), the twin of Darling's, used ONLY to decide
+ /// which database to skip: one that fails
+ /// slices in a row is served after the databases behind it instead of ahead of them, so it can no longer
+ /// stall them; that database's completed slice resets it. It does not size the slice window: that stays the
+ /// per-server count above, because a command timeout usually means the whole server is loaded, and
+ /// narrowing per database would add timed-out queries per database against a server that is already
+ /// struggling.
+ ///
+ private readonly QueryStoreBackfillFailureLedger _sliceFailures = new();
+
+ /// Test-only seam: when set, replaces the slice body (called with the database and the window
+ /// span the slice would have used). A throw counts as a failed slice and a normal return as a completed
+ /// one, through the same accounting. Null in production, where it changes nothing.
+ internal Func? SliceOverrideForTests { get; set; }
+
/// Runs one slice with the failure accounting wrapped around it — the caller's outer
/// catch still logs the throw exactly as before.
private async Task RunCountedBackfillSliceAsync(
@@ -178,10 +196,22 @@ private async Task RunCountedBackfillSliceAsync(
{
await RunBackfillSliceAsync(server, serverId, target, databaseName, floorUtc, ceilingUtc, isHole, cancellationToken);
_consecutiveSliceFailures.TryRemove(serverId, out _);
+ _sliceFailures.RecordCompletion(serverId, databaseName);
}
catch (Exception ex) when (ex is not OperationCanceledException)
{
_consecutiveSliceFailures.AddOrUpdate(serverId, 1, static (_, current) => current + 1);
+ var failures = _sliceFailures.RecordFailure(serverId, databaseName);
+
+ /* Logged at the failure that crosses the threshold, so it is once per stretch of failures and
+ needs no extra state: the count only grows until a completed slice clears it. */
+ if (failures == QueryStoreBackfillState.SkipAfterConsecutiveSliceFailures)
+ {
+ _logger?.LogWarning(
+ "query_store backfill on '{Server}' [{Database}]: {Failures} consecutive slice failures; serving the other databases first and retrying this one only when none has work.",
+ server.DisplayName, databaseName, failures);
+ }
+
throw;
}
}
@@ -230,6 +260,10 @@ internal async Task RunQueryStoreBackfillSliceAsync(ServerConnection serve
a hole key already names (state is loaded above, for free). */
var databases = await GetBackfillCandidateDatabasesAsync(serverId, floorLimit, state, cancellationToken);
+ /* Databases whose slices keep failing: their slice is held back while any other database has work
+ (see QueryStoreBackfillFailureLedger), then one of them is retried after the walk. */
+ List<(string Database, DateTime Floor, DateTime Ceiling, bool IsHole)>? skipped = null;
+
foreach (var databaseName in databases)
{
cancellationToken.ThrowIfCancellationRequested();
@@ -245,6 +279,12 @@ internal async Task RunQueryStoreBackfillSliceAsync(ServerConnection serve
}
var holeFloor = holeFrom > floorLimit ? holeFrom : floorLimit;
+ if (_sliceFailures.IsSkipped(serverId, databaseName))
+ {
+ (skipped ??= []).Add((databaseName, holeFloor, holeTo, true));
+ continue;
+ }
+
await RunCountedBackfillSliceAsync(server, serverId, target, databaseName, holeFloor, holeTo, isHole: true, cancellationToken);
return true;
}
@@ -273,10 +313,35 @@ await SaveCollectorStateAsync(serverId, QueryStoreBackfillState.StateCollectorNa
continue;
}
+ if (_sliceFailures.IsSkipped(serverId, databaseName))
+ {
+ (skipped ??= []).Add((databaseName, floorLimit, storedFloor.Value, false));
+ continue;
+ }
+
await RunCountedBackfillSliceAsync(server, serverId, target, databaseName, floorLimit, storedFloor.Value, isHole: false, cancellationToken);
return true;
}
+ /* No other database had work, so retry a skipped one: the one whose last failure is the oldest, so
+ several skipped databases take turns instead of the first in the list starving the rest. This costs
+ at most one failed slice per tick on an otherwise idle server, exactly what the stall cost before. */
+ if (skipped is { Count: > 0 })
+ {
+ var retry = skipped[0];
+ for (var i = 1; i < skipped.Count; i++)
+ {
+ if (_sliceFailures.LastFailureTicket(serverId, skipped[i].Database)
+ < _sliceFailures.LastFailureTicket(serverId, retry.Database))
+ {
+ retry = skipped[i];
+ }
+ }
+
+ await RunCountedBackfillSliceAsync(server, serverId, target, retry.Database, retry.Floor, retry.Ceiling, retry.IsHole, cancellationToken);
+ return true;
+ }
+
return false;
}
@@ -305,6 +370,12 @@ window on a big database times out at the command timeout every tick and the ran
_consecutiveSliceFailures.TryGetValue(serverId, out var recentFailures) ? recentFailures : 0);
var sliceFloor = QueryStoreBackfillState.BoundSliceFloor(floorUtc, ceilingUtc, sliceSpan);
+ if (SliceOverrideForTests is { } sliceOverride)
+ {
+ await sliceOverride(databaseName, sliceSpan);
+ return;
+ }
+
var definition = QueryStoreCollector.Instance;
var context = new CollectorContext
{
diff --git a/PerformanceMonitor.Collectors/QueryStoreBackfillFailureLedger.cs b/PerformanceMonitor.Collectors/QueryStoreBackfillFailureLedger.cs
new file mode 100644
index 0000000000..79e4c9aa65
--- /dev/null
+++ b/PerformanceMonitor.Collectors/QueryStoreBackfillFailureLedger.cs
@@ -0,0 +1,60 @@
+/*
+ * 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.Collections.Concurrent;
+using System.Threading;
+
+namespace PerformanceMonitor.Collectors;
+
+///
+/// The in-memory count of consecutive failed Query Store backfill slices per (server, database), shared by
+/// both products' workers so they skip a failing database by the same rule. In memory on purpose, like the
+/// live path's counters: a restart forgets it and costs at most
+/// more failed slices. Entries are removed by a completed slice, so the map holds only databases that are
+/// failing right now.
+///
+/// The count has one job: once it reaches
+/// the loop serves the databases behind that database first. It does NOT size the slice window. That is each
+/// worker's own per-server failure count (), kept as it was,
+/// because a command timeout usually means the whole server is loaded and narrowing per database would only add
+/// timed-out queries against it. Each failure also takes a ticket from a running sequence so the skipped
+/// databases can be retried in turn, least recently failed first; without it the first skipped database in the
+/// list would starve every other skipped one, the same stall one level down.
+///
+public sealed class QueryStoreBackfillFailureLedger
+{
+ private readonly ConcurrentDictionary<(int ServerId, string Database), (int Failures, long Ticket)> _entries = new();
+ private long _ticket;
+
+ /// Consecutive failed slices for the database; 0 when none (or it last completed). Feeds the skip
+ /// decision only, never the window size.
+ public int Failures(int serverId, string database)
+ => _entries.TryGetValue((serverId, database), out var entry) ? entry.Failures : 0;
+
+ /// True once the database has failed slices in a row.
+ public bool IsSkipped(int serverId, string database)
+ => Failures(serverId, database) >= QueryStoreBackfillState.SkipAfterConsecutiveSliceFailures;
+
+ /// Counts one failed slice and returns the new consecutive count.
+ public int RecordFailure(int serverId, string database)
+ {
+ var ticket = Interlocked.Increment(ref _ticket);
+ return _entries.AddOrUpdate(
+ (serverId, database),
+ (1, ticket),
+ (_, current) => (current.Failures + 1, ticket)).Failures;
+ }
+
+ /// A completed slice clears the database's count.
+ public void RecordCompletion(int serverId, string database)
+ => _entries.TryRemove((serverId, database), out _);
+
+ /// The ticket of the database's latest failure; smaller = failed longer ago, 0 = never.
+ public long LastFailureTicket(int serverId, string database)
+ => _entries.TryGetValue((serverId, database), out var entry) ? entry.Ticket : 0;
+}
diff --git a/PerformanceMonitor.Collectors/QueryStoreBackfillState.cs b/PerformanceMonitor.Collectors/QueryStoreBackfillState.cs
index 2cf6127261..d99692366a 100644
--- a/PerformanceMonitor.Collectors/QueryStoreBackfillState.cs
+++ b/PerformanceMonitor.Collectors/QueryStoreBackfillState.cs
@@ -76,6 +76,26 @@ public static bool ShouldYieldToLive(DateTime? lastLiveFailureUtc, DateTime nowU
///
public static readonly TimeSpan MinAdaptiveSpan = TimeSpan.FromMinutes(15);
+ ///
+ /// How many consecutive failed slices one database may have before the backfill loop skips it and
+ /// moves on to the databases behind it. A slice runs at most once per server per tick, and the
+ /// candidate list is ordered the same way every tick, so a database whose slices always fail used to be
+ /// first in line forever and no database after it on that server ever got a slice.
+ ///
+ /// Why 3. The window still narrows per server, not per database
+ /// ( of the server's consecutive failures). The stall this guards against is a
+ /// database that is first in line on a server with a fresh failure count, so its attempts run at the full
+ /// span (0 failures so far), half of it (1) and (2). The third failure is
+ /// therefore the first one at the narrowest span, so a database whose slices merely time out gets one
+ /// attempt at the narrowest span before it is skipped; any smaller number would skip it while a narrower
+ /// window could still fit. Pinned against so a change to either side shows up
+ /// as a failing test.
+ ///
+ /// A skipped database goes to the back of the line, not away: it is retried on any tick where no
+ /// other database has work, and a completed slice clears its count.
+ ///
+ public const int SkipAfterConsecutiveSliceFailures = 3;
+
///
/// The window a member gets after straight failures:
/// the full span halved per failure, floored at (the exponent is