From e95166c21cc35aeb56fe02faa1485a5cbf0f5b91 Mon Sep 17 00:00:00 2001
From: Erik Darling <2136037+erikdarlingdata@users.noreply.github.com>
Date: Wed, 30 Sep 2026 06:09:10 -0400
Subject: [PATCH 1/3] Hourly and daily rollup reads stop at the window end
instead of taking the bucket that starts there
---
.../DarlingMcpTrendToolsTests.cs | 5 +-
.../Darling.Tests/HourlyWindowEdgesTests.cs | 24 +-
.../QueryStoreTrendRoutingTests.cs | 5 +-
.../RollupWindowEndBoundLiveTests.cs | 378 ++++++++++++++++++
.../RollupWindowEndBoundTests.cs | 206 ++++++++++
.../Compose/ComposeCompiler.cs | 7 +-
.../Mcp/DarlingDataReader.cs | 12 +-
.../Mcp/DarlingTrendReader.cs | 2 +-
.../DurationTrendRouting.cs | 10 +-
.../HourlyWindowEdges.cs | 10 +-
.../QueryStoreTrendRouting.cs | 7 +-
.../ViewerDataService.ProcedureStats.cs | 55 +--
.../ViewerDataService.QueryStats.cs | 49 ++-
13 files changed, 704 insertions(+), 66 deletions(-)
create mode 100644 Darling/Darling.Tests/RollupWindowEndBoundLiveTests.cs
create mode 100644 Darling/Darling.Tests/RollupWindowEndBoundTests.cs
diff --git a/Darling/Darling.Tests/DarlingMcpTrendToolsTests.cs b/Darling/Darling.Tests/DarlingMcpTrendToolsTests.cs
index 3e913c27f..2346d0b14 100644
--- a/Darling/Darling.Tests/DarlingMcpTrendToolsTests.cs
+++ b/Darling/Darling.Tests/DarlingMcpTrendToolsTests.cs
@@ -415,7 +415,10 @@ public void DurationTrendHourlySql_ReadsTheRollup_BucketsByIt_ProjectsTheSharedS
Assert.DoesNotContain("FROM " + rawTable + "\n", sql, StringComparison.Ordinal);
Assert.DoesNotContain("collection_time >=", sql, StringComparison.Ordinal);
Assert.Contains("bucket >= $2", sql, StringComparison.Ordinal);
- Assert.Contains("bucket <= $3", sql, StringComparison.Ordinal);
+ /* A bucket is stamped at its START, so the window end is exclusive: `bucket <= $3` would take the whole
+ hour that begins at an end falling on the hour (RollupWindowEndBoundTests pins the same for every rollup read). */
+ Assert.Contains("bucket < $3", sql, StringComparison.Ordinal);
+ Assert.DoesNotContain("bucket <= $3", sql, StringComparison.Ordinal);
Assert.Contains("GROUP BY bucket", sql, StringComparison.Ordinal);
/* Since #3897 the rollup's hours are gathered into $4-minute buckets on the shared origin; the columns and
diff --git a/Darling/Darling.Tests/HourlyWindowEdgesTests.cs b/Darling/Darling.Tests/HourlyWindowEdgesTests.cs
index 1cc5d2e59..d1a263e34 100644
--- a/Darling/Darling.Tests/HourlyWindowEdgesTests.cs
+++ b/Darling/Darling.Tests/HourlyWindowEdgesTests.cs
@@ -33,16 +33,30 @@ public void AlignedStart_HasNoLeadingSentence()
Assert.DoesNotContain("partial hour", note, StringComparison.Ordinal);
}
- [Theory]
- [InlineData(20)]
- [InlineData(0)]
- public void End_CountsTheWholeHourPastAsOf(int minute)
+ [Fact]
+ public void UnalignedEnd_CountsTheWholeHourPastAsOf()
{
- var end = Aligned.AddHours(4).AddMinutes(minute);
+ var end = Aligned.AddHours(4).AddMinutes(20);
var note = HourlyWindowEdges.Note(Aligned, Aligned, end, Aligned.AddHours(9));
Assert.Contains("up to " + Aligned.AddHours(5).ToString("o"), note, StringComparison.Ordinal);
}
+ /// The hourly reads stop BEFORE the window end (bucket < end), so an end exactly on the
+ /// hour no longer takes the hour that begins there: the served span ends at the window end and the note
+ /// claims nothing past it. (This case used to be pinned as "counts the whole hour", when the bound was
+ /// bucket <= end.)
+ [Fact]
+ public void AlignedEnd_StopsAtTheEnd_AndCountsNothingPastAsOf()
+ {
+ var end = Aligned.AddHours(4);
+ var note = HourlyWindowEdges.Note(Aligned, Aligned, end, Aligned.AddHours(9));
+ Assert.DoesNotContain("included whole", note, StringComparison.Ordinal);
+ Assert.DoesNotContain("counted past", note, StringComparison.Ordinal);
+
+ var (_, servedEnd) = HourlyWindowEdges.ServedSpan(Aligned, Aligned, end, Aligned.AddHours(9));
+ Assert.Equal(end, servedEnd);
+ }
+
[Fact]
public void EndCutAtTheCeiling_SaysNothingAfterItWasRead()
{
diff --git a/Darling/Darling.Tests/QueryStoreTrendRoutingTests.cs b/Darling/Darling.Tests/QueryStoreTrendRoutingTests.cs
index 2091f6a86..ea01d10cb 100644
--- a/Darling/Darling.Tests/QueryStoreTrendRoutingTests.cs
+++ b/Darling/Darling.Tests/QueryStoreTrendRoutingTests.cs
@@ -184,7 +184,10 @@ public void RollupTrendSql_ServesTheRollup_AndPartitionsAtTheBoundary(bool withD
/* The partition seam: rollup buckets strictly BELOW $4 (load-bearing against a refresh landing
between the probe and the read), raw points at or above it. */
Assert.Contains("bucket >= $2", sql, StringComparison.Ordinal);
- Assert.Contains("bucket <= $3", sql, StringComparison.Ordinal);
+ /* The window end is exclusive too: a bucket is stamped at its START, so `bucket <= $3` would take the
+ whole hour that begins at an end falling on the hour. */
+ Assert.Contains("bucket < $3", sql, StringComparison.Ordinal);
+ Assert.DoesNotContain("bucket <= $3", sql, StringComparison.Ordinal);
Assert.Contains("bucket < $4", sql, StringComparison.Ordinal);
Assert.Contains("interval_start_time_utc >= $4", sql, StringComparison.Ordinal);
diff --git a/Darling/Darling.Tests/RollupWindowEndBoundLiveTests.cs b/Darling/Darling.Tests/RollupWindowEndBoundLiveTests.cs
new file mode 100644
index 000000000..41baf1ddf
--- /dev/null
+++ b/Darling/Darling.Tests/RollupWindowEndBoundLiveTests.cs
@@ -0,0 +1,378 @@
+/*
+ * 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.Linq;
+using System.Text.Json.Nodes;
+using System.Threading;
+using System.Threading.Tasks;
+using Npgsql;
+using PerformanceMonitor.Collectors;
+using PerformanceMonitor.Darling.Service;
+using PerformanceMonitor.Darling.Service.Mcp;
+using PerformanceMonitor.Darling.Storage;
+using PerformanceMonitor.Darling.Viewer;
+using Xunit;
+
+namespace Darling.Tests;
+
+///
+/// A rollup bucket is stamped at its START, so a read that bounds the window end with bucket <= end
+/// took the whole bucket that begins AT the end whenever the end fell exactly on a bucket start: one extra hour
+/// in an hourly sum. These live pins seed one query and one procedure in three hours (H0 = 1 execution, H1 = 10,
+/// H2 = 100, so a total names its hours), refresh the hourly rollups, and read each rollup path with three
+/// ends: (i) exactly H2, which must NOT add the hour that begins there (11); (ii) H2 plus 25 minutes, with the
+/// newest bucket still open as "now" is, which counts H2 (111); (iii) H1 plus 30 minutes, inside H1 (11, the
+/// same as before the change). Every raw row is stamped 10 minutes past its hour, never exactly on it: the raw
+/// arm still reads a row stamped exactly at the end, so an exact-on-the-hour stamp would blur the two arms.
+///
+/// Covered here against real rollups: the MCP top-queries and top-procedures reads, the viewer's Queries
+/// and Procedures tab reads (with the raw arm's total for the same window as the agreement check), the hourly
+/// duration trend (plain and bucketed), the one-query hourly history, and a Custom View panel on the hourly
+/// route. The Query Store rollup trend and the daily route of a Custom View are pinned at the text level in
+/// RollupWindowEndBoundTests.
+///
+/// #1776 own-store — each test mints a scratch database (it materializes continuous aggregates the
+/// shared fixture must never inherit), so the class is deliberately NOT in the live-postgres
+/// collection.
+///
+public sealed class RollupWindowEndBoundLiveTests
+{
+ private const int ServerId = -947101;
+ private const string ServerName = "rollup-window-end-bound";
+ private const string ComposeServerName = "rollup-window-end-compose";
+ private const int ComposeServerId = -947102;
+ private const string Db = "RollupEndDb";
+ private const string QueryHash = "0xRWEQ1";
+ private const string ProcName = "usp_RollupEnd";
+
+ /// A fixed anchor, never wall-clock relative: the window is far older than raw retention, and it is
+ /// the purge of raw's rows (and the refresh of the rollup) that drives the router, not calendar time.
+ private static readonly DateTime WindowStart = new(2026, 1, 5, 0, 0, 0, DateTimeKind.Unspecified);
+
+ private static readonly DateTime H0 = WindowStart.AddHours(1);
+ private static readonly DateTime H1 = WindowStart.AddHours(2);
+ private static readonly DateTime H2 = WindowStart.AddHours(3);
+
+ /// The three window ends and what each one sums, in executions.
+ private static readonly (string Label, DateTime End, long Executions, DateTime[] Buckets)[] Ends =
+ {
+ ("end exactly at H2", H2, 11L, new[] { H0, H1 }),
+ ("end 25 minutes into H2", H2.AddMinutes(25), 111L, new[] { H0, H1, H2 }),
+ ("end 30 minutes into H1", H1.AddMinutes(30), 11L, new[] { H0, H1 }),
+ };
+
+ [Fact]
+ public async Task QueryAndProcedureReads_StopAtTheWindowEnd_OnRawAndOnTheHourlyRollup()
+ {
+ 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 rollup window-end test.");
+
+ 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);
+
+ var timescaleEnabled = await TimescaleSupport.TryEnableAsync(connection, null, ct);
+ Assert.SkipWhen(!timescaleEnabled, "The live rollup window-end test needs TimescaleDB.");
+ await TimescaleSupport.ConvertToHypertablesAsync(connection, null, ct);
+ Assert.True(await TimescaleSupport.EnsureCollectionLogHypertableAsync(connection, null, ct));
+
+ await using (var stop = new NpgsqlCommand("SELECT _timescaledb_functions.stop_background_workers()", connection))
+ {
+ await stop.ExecuteNonQueryAsync(ct);
+ }
+
+ await DarlingMcpTestData.RegisterServerAsync(connection, ServerId, ServerName, ct);
+ await TimescaleSupport.EnsureContinuousAggregatesAsync(connection, null, ct);
+
+ var windowEnd = WindowStart.AddDays(1);
+ var bodySucceeded = false;
+ try
+ {
+ /* One query and one procedure in three hours: 1, 10 and 100 executions, each stamped 10 minutes past its hour. */
+ var hours = new[] { H0, H1, H2 };
+ var executions = new[] { 1L, 10L, 100L };
+ for (var i = 0; i < hours.Length; i++)
+ {
+ await PlantQueryAsync(connection, ct, ServerId, ServerName, hours[i].AddMinutes(10), executions[i]);
+ await PlantProcedureAsync(connection, ct, hours[i].AddMinutes(10), executions[i]);
+ }
+
+ /* ── raw is still present: this is the tier the hourly reads must agree with. A viewer instance caches
+ its rollup probe for five minutes, so the raw-stage instance is never reused after the refresh. ── */
+ await using var rawSource = NpgsqlDataSource.Create(scratch.ConnectionString);
+ await using var rawViewer = new ViewerDataService(scratch.ConnectionString);
+ foreach (var (label, end, expected, _) in Ends)
+ {
+ var mcpQueries = await DarlingDataReader.GetTopQueriesByCpuRoutedAsync(
+ rawSource, ServerId, WindowStart, end, top: 10, databaseName: null, cancellationToken: ct);
+ Assert.Equal(RetentionTier.Raw, mcpQueries.Tier);
+ Assert.Equal(expected, mcpQueries.Rows.Sum(r => r.TotalExecutions));
+
+ var mcpProcedures = await DarlingDataReader.GetTopProceduresByCpuRoutedAsync(
+ rawSource, ServerId, WindowStart, end, top: 10, databaseName: null, cancellationToken: ct);
+ Assert.Equal(RetentionTier.Raw, mcpProcedures.Tier);
+ Assert.Equal(expected, mcpProcedures.Rows.Sum(r => r.TotalExecutions));
+
+ var (viewerQueries, queriesTier) = await rawViewer.GetTopQueriesByCpuTierAsync(ServerId, WindowStart, end, cancellationToken: ct);
+ Assert.Equal("raw", queriesTier);
+ Assert.True(expected == viewerQueries.Sum(r => r.TotalExecutions), $"viewer raw queries, {label}");
+
+ var (viewerProcedures, proceduresTier) = await rawViewer.GetTopProceduresByCpuTierAsync(ServerId, WindowStart, end, cancellationToken: ct);
+ Assert.Equal("raw", proceduresTier);
+ Assert.True(expected == viewerProcedures.Sum(r => r.TotalExecutions), $"viewer raw procedures, {label}");
+ }
+
+ /* Refresh the hourly successors over the window BEFORE deleting raw: the product's own refresh path,
+ not a hand-built rollup row. */
+ await RefreshAsync(connection, TimescaleSupport.QueryStatsIntervalHourlyView, WindowStart, windowEnd.AddHours(1), ct);
+ await RefreshAsync(connection, TimescaleSupport.ProcedureStatsIntervalHourlyView, WindowStart, windowEnd.AddHours(1), ct);
+
+ /* ── the trend and history builders, run over the seeded hourly relations ── */
+ var queryRollup = "collect." + TimescaleSupport.QueryStatsIntervalHourlyView;
+ var procedureRollup = "collect." + TimescaleSupport.ProcedureStatsIntervalHourlyView;
+ foreach (var (label, end, _, buckets) in Ends)
+ {
+ foreach (var rollup in new[] { queryRollup, procedureRollup })
+ {
+ var plain = await ReadFirstColumnTimesAsync(
+ connection, DurationTrendRouting.BuildHourlyTrendSql(rollup, withDatabaseFilter: false), ct, ServerId, WindowStart, end);
+ Assert.True(buckets.SequenceEqual(plain), $"hourly duration trend over {rollup}, {label}: [{string.Join(", ", plain)}]");
+
+ var bucketed = await ReadFirstColumnTimesAsync(
+ connection, DurationTrendRouting.BuildBucketedHourlyTrendSql(rollup), ct, ServerId, WindowStart, end, 60);
+ Assert.True(buckets.SequenceEqual(bucketed), $"bucketed hourly duration trend over {rollup}, {label}: [{string.Join(", ", bucketed)}]");
+ }
+
+ var history = await ReadFirstColumnTimesAsync(
+ connection, DarlingTrendReader.QueryHistoryHourlySqlFor(queryRollup), ct, ServerId, Db, QueryHash, WindowStart, end);
+ Assert.True(buckets.SequenceEqual(history), $"one-query hourly history, {label}: [{string.Join(", ", history)}]");
+ }
+
+ /* ── move raw's floor past the window: delete W's raw rows, as retention would. ── */
+ foreach (var table in new[] { "query_stats", "procedure_stats" })
+ {
+ await using var purge = new NpgsqlCommand(
+ $"DELETE FROM collect.{table} WHERE server_id = $1 AND collection_time >= $2 AND collection_time < $3", connection);
+ purge.Parameters.AddWithValue(ServerId);
+ purge.Parameters.AddWithValue(WindowStart);
+ purge.Parameters.AddWithValue(windowEnd);
+ await purge.ExecuteNonQueryAsync(ct);
+ }
+
+ /* ── the hourly tier answers. FRESH data source and viewer: coverage is cached per instance, and the
+ instances above cached the null hourly floor measured before the refresh. ── */
+ await using var hourlySource = NpgsqlDataSource.Create(scratch.ConnectionString);
+ await using var hourlyViewer = new ViewerDataService(scratch.ConnectionString);
+ foreach (var (label, end, expected, _) in Ends)
+ {
+ var mcpQueries = await DarlingDataReader.GetTopQueriesByCpuRoutedAsync(
+ hourlySource, ServerId, WindowStart, end, top: 10, databaseName: null, cancellationToken: ct);
+ Assert.Equal(RetentionTier.Hourly, mcpQueries.Tier);
+ Assert.True(expected == mcpQueries.Rows.Sum(r => r.TotalExecutions), $"MCP hourly top queries, {label}");
+
+ var mcpProcedures = await DarlingDataReader.GetTopProceduresByCpuRoutedAsync(
+ hourlySource, ServerId, WindowStart, end, top: 10, databaseName: null, cancellationToken: ct);
+ Assert.Equal(RetentionTier.Hourly, mcpProcedures.Tier);
+ Assert.True(expected == mcpProcedures.Rows.Sum(r => r.TotalExecutions), $"MCP hourly top procedures, {label}");
+
+ var (viewerQueries, queriesTier) = await hourlyViewer.GetTopQueriesByCpuTierAsync(ServerId, WindowStart, end, cancellationToken: ct);
+ Assert.Equal("hourly", queriesTier);
+ Assert.True(expected == viewerQueries.Sum(r => r.TotalExecutions), $"viewer hourly queries, {label}");
+
+ var (viewerProcedures, proceduresTier) = await hourlyViewer.GetTopProceduresByCpuTierAsync(ServerId, WindowStart, end, cancellationToken: ct);
+ Assert.Equal("hourly", proceduresTier);
+ Assert.True(expected == viewerProcedures.Sum(r => r.TotalExecutions), $"viewer hourly procedures, {label}");
+ }
+
+ bodySucceeded = true;
+ }
+ finally
+ {
+ await CleanupAsync(scratch.ConnectionString, bodySucceeded);
+ }
+ }
+
+ ///
+ /// A Custom View panel on the hourly route sums the hours up to the window end and stops before the hour that
+ /// begins at it. The same three ends, the same rows: the panel's worker time is 1,000 per execution, so the
+ /// sums name their hours (11,000 / 111,000 / 11,000).
+ ///
+ [Fact]
+ public async Task CustomViewPanel_OnTheHourlyRoute_StopsAtTheWindowEnd()
+ {
+ 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 rollup window-end test.");
+ 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);
+
+ var timescaleEnabled = await TimescaleSupport.TryEnableAsync(connection, null, ct);
+ Assert.SkipWhen(!timescaleEnabled, "The live rollup window-end test needs TimescaleDB: the rollup route reads a continuous aggregate.");
+ await TimescaleSupport.ConvertToHypertablesAsync(connection, null, ct);
+ Assert.True(await TimescaleSupport.EnsureCollectionLogHypertableAsync(connection, null, ct));
+
+ await using (var stop = new NpgsqlCommand("SELECT _timescaledb_functions.stop_background_workers()", connection))
+ {
+ await stop.ExecuteNonQueryAsync(ct);
+ }
+
+ await DarlingMcpTestData.RegisterServerAsync(connection, ComposeServerId, ComposeServerName, ct);
+ await TimescaleSupport.EnsureContinuousAggregatesAsync(connection, null, ct);
+
+ var bodySucceeded = false;
+ try
+ {
+ var hours = new[] { H0, H1, H2 };
+ var executions = new[] { 1L, 10L, 100L };
+ for (var i = 0; i < hours.Length; i++)
+ {
+ await PlantQueryAsync(connection, ct, ComposeServerId, ComposeServerName, hours[i].AddMinutes(10), executions[i]);
+ }
+
+ /* Both hourly rollups are refreshed, as the rollup-average pin does: the compiler picks between the
+ legacy hourly and its successor by coverage, and either must read the same hours. */
+ foreach (var view in new[] { TimescaleSupport.QueryStatsHourlyView, TimescaleSupport.QueryStatsIntervalHourlyView })
+ {
+ await RefreshAsync(connection, view, WindowStart, WindowStart.AddDays(1), ct);
+ }
+
+ await using var dataSource = NpgsqlDataSource.Create(scratch.ConnectionString);
+ var rollups = await TimescaleSupport.DetectRollupsAsync(dataSource, ct);
+ var coverage = await TimescaleSupport.DetectRollupCoverageAsync(dataSource, rollups, ct);
+
+ var json = (JsonObject)JsonNode.Parse(
+ "{\"source\":\"query_stats\",\"measure\":\"query_worker_us\",\"aggregate\":\"sum\",\"viz\":\"table\"}")!;
+ var (plan, parseError) = ComposeSpec.TryParsePanel(json, Array.Empty());
+ Assert.True(parseError is null, parseError);
+
+ foreach (var (label, end, expected, _) in Ends)
+ {
+ /* "Now" is days past the window, so the panel is old enough for the rollup route. */
+ var context = new ComposeRunContext(
+ new[] { ComposeServerName }, H0, end, ComposeRunContext.NoVariables, rollups, H2.AddDays(5), coverage);
+ var (compiled, error) = ComposeCompiler.Compile(plan!, context);
+ Assert.True(error is null, error);
+ Assert.Equal(ComposeSourceTier.Hourly, compiled!.Route.Tier);
+
+ await using var command = new NpgsqlCommand(compiled.Sql, connection);
+ foreach (var p in compiled.Parameters)
+ {
+ command.Parameters.Add(p);
+ }
+
+ await using var reader = await command.ExecuteReaderAsync(ct);
+ Assert.True(await reader.ReadAsync(ct));
+ var sum = Convert.ToDouble(reader.GetValue(reader.FieldCount - 1));
+ Assert.True(expected * 1000d == sum, $"Custom View sum of worker time, {label}: {sum:R}");
+ }
+
+ bodySucceeded = true;
+ }
+ finally
+ {
+ await CleanupAsync(scratch.ConnectionString, bodySucceeded);
+ }
+ }
+
+ private static async Task CleanupAsync(string connectionString, bool bodySucceeded)
+ {
+ await LiveStoreCleanup.RunAsync(connectionString, bodySucceeded, async (cleanup, cleanupCt) =>
+ {
+ await using var probe = new NpgsqlCommand(
+ "SELECT count(*) FROM pg_catalog.pg_stat_activity WHERE datname = pg_catalog.current_database() " +
+ "AND backend_type LIKE 'TimescaleDB Background Worker Scheduler%'", cleanup);
+ var schedulers = Convert.ToInt64(await probe.ExecuteScalarAsync(cleanupCt));
+ Assert.Equal(0L, schedulers);
+ });
+ }
+
+ /// Runs with positional parameters and returns the first column of every row
+ /// as a time.
+ private static async Task> ReadFirstColumnTimesAsync(
+ NpgsqlConnection connection, string sql, CancellationToken ct, params object[] parameters)
+ {
+ await using var command = new NpgsqlCommand(sql, connection);
+ foreach (var value in parameters)
+ {
+ command.Parameters.AddWithValue(value);
+ }
+
+ var times = new List();
+ await using var reader = await command.ExecuteReaderAsync(ct);
+ while (await reader.ReadAsync(ct))
+ {
+ times.Add(reader.GetDateTime(0));
+ }
+
+ return times;
+ }
+
+ /// One query_stats row: worker time is 1,000 per execution, a real (nonzero) sample interval so the
+ /// hourly rollup admits it.
+ private static async Task PlantQueryAsync(
+ NpgsqlConnection connection, CancellationToken ct, int serverId, string serverName, DateTime at, long executions)
+ {
+ await using var insert = new NpgsqlCommand(@"
+INSERT INTO collect.query_stats
+ (collection_id, collection_time, server_id, server_name, database_name, query_hash, sql_handle,
+ host_object_name, delta_worker_time, delta_elapsed_time, delta_execution_count, sample_interval_seconds)
+VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12)", connection);
+ insert.Parameters.AddWithValue(CollectionIdGenerator.Next());
+ insert.Parameters.AddWithValue(DarlingMcpTestData.TruncateToSeconds(at));
+ insert.Parameters.AddWithValue(serverId);
+ insert.Parameters.AddWithValue(serverName);
+ insert.Parameters.AddWithValue(Db);
+ insert.Parameters.AddWithValue(QueryHash);
+ insert.Parameters.AddWithValue("0xRWEQ1H");
+ insert.Parameters.AddWithValue("usp_RollupEndHost");
+ insert.Parameters.AddWithValue(executions * 1_000L);
+ insert.Parameters.AddWithValue(executions * 900L);
+ insert.Parameters.AddWithValue(executions);
+ insert.Parameters.AddWithValue(3600);
+ await insert.ExecuteNonQueryAsync(ct);
+ }
+
+ private static async Task PlantProcedureAsync(NpgsqlConnection connection, CancellationToken ct, DateTime at, long executions)
+ {
+ await using var insert = new NpgsqlCommand(@"
+INSERT INTO collect.procedure_stats
+ (collection_id, collection_time, server_id, server_name, database_name, schema_name, object_name, sql_handle,
+ delta_worker_time, delta_elapsed_time, delta_execution_count, sample_interval_seconds)
+VALUES ($1, $2, $3, $4, $5, 'dbo', $6, $7, $8, $9, $10, $11)", connection);
+ insert.Parameters.AddWithValue(CollectionIdGenerator.Next());
+ insert.Parameters.AddWithValue(DarlingMcpTestData.TruncateToSeconds(at));
+ insert.Parameters.AddWithValue(ServerId);
+ insert.Parameters.AddWithValue(ServerName);
+ insert.Parameters.AddWithValue(Db);
+ insert.Parameters.AddWithValue(ProcName);
+ insert.Parameters.AddWithValue("0x" + ProcName);
+ insert.Parameters.AddWithValue(executions * 1_000L);
+ insert.Parameters.AddWithValue(executions * 900L);
+ insert.Parameters.AddWithValue(executions);
+ insert.Parameters.AddWithValue(3600);
+ await insert.ExecuteNonQueryAsync(ct);
+ }
+
+ private static async Task RefreshAsync(NpgsqlConnection connection, string view, DateTime from, DateTime to, CancellationToken ct)
+ {
+ await using var refresh = new NpgsqlCommand($"CALL refresh_continuous_aggregate('collect.{view}'::regclass, $1::timestamp, $2::timestamp)", connection);
+ refresh.Parameters.AddWithValue(from);
+ refresh.Parameters.AddWithValue(to);
+ await refresh.ExecuteNonQueryAsync(ct);
+ }
+}
diff --git a/Darling/Darling.Tests/RollupWindowEndBoundTests.cs b/Darling/Darling.Tests/RollupWindowEndBoundTests.cs
new file mode 100644
index 000000000..92401db47
--- /dev/null
+++ b/Darling/Darling.Tests/RollupWindowEndBoundTests.cs
@@ -0,0 +1,206 @@
+/*
+ * 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.Text.Json.Nodes;
+using System.Text.RegularExpressions;
+using PerformanceMonitor.Darling.Service;
+using PerformanceMonitor.Darling.Service.Mcp;
+using PerformanceMonitor.Darling.Storage;
+using PerformanceMonitor.Darling.Viewer;
+using Xunit;
+
+namespace Darling.Tests;
+
+///
+/// A rollup bucket is stamped at its START. A read that bounds the window end with bucket <= end
+/// therefore takes the whole bucket that begins AT the end whenever the end falls exactly on a bucket start: one
+/// extra hour in an hourly sum, one extra day on a daily-tier Custom View. Nine reads over the hourly (and, for
+/// Custom Views, daily) rollups now bound the end with bucket < end. These are the text-level pins: each
+/// read's end parameter is bounded by < and no bucket <= is left in the text. The behavior
+/// itself (the hour that starts at the end is not summed) is pinned against real rollups in
+/// RollupWindowEndBoundLiveTests.
+///
+public sealed class RollupWindowEndBoundTests
+{
+ private const string ViewName = "collect.query_stats_hourly";
+
+ /// The text bounds bucket from above with < $endParam (a bare $N, not a
+ /// longer number) and holds no bucket <= anywhere.
+ private static void AssertEndIsExclusive(string sql, int endParam, string because)
+ {
+ Assert.True(
+ Regex.IsMatch(sql, @"\bbucket\s*<\s*\$" + endParam + @"(?!\d)"),
+ because + ": the window end is no longer bounded by `bucket < $" + endParam + "`.");
+ Assert.False(
+ Regex.IsMatch(sql, @"\bbucket\s*<="),
+ because + ": a `bucket <=` bound is back, and takes the bucket that starts at the window end.");
+ }
+
+ /* Path 1: the Queries tab's hourly arm (a viewer read, so its text is a builder). */
+ [Fact]
+ public void ViewerTopQueriesHourlySql_StopsBeforeTheWindowEnd()
+ {
+ var sql = ViewerDataService.BuildTopQueriesHourlySql(ViewName);
+
+ AssertEndIsExclusive(sql, 3, "Queries tab hourly arm");
+ Assert.Contains("bucket >= $2", sql, StringComparison.Ordinal);
+ Assert.Contains("FROM " + ViewName, sql, StringComparison.Ordinal);
+ }
+
+ /* Path 2: the Procedures tab's hourly arm. */
+ [Fact]
+ public void ViewerTopProceduresHourlySql_StopsBeforeTheWindowEnd()
+ {
+ var sql = ViewerDataService.BuildTopProceduresHourlySql(ViewName);
+
+ AssertEndIsExclusive(sql, 3, "Procedures tab hourly arm");
+ Assert.Contains("bucket >= $2", sql, StringComparison.Ordinal);
+ Assert.Contains("FROM " + ViewName, sql, StringComparison.Ordinal);
+ }
+
+ /* Paths 3 and 4: the MCP top-queries / top-procedures hourly reads. The materialization-ceiling placeholder
+ stays right after the bound, so a known ceiling still narrows the read. */
+ [Fact]
+ public void McpTopQueriesHourlySql_StopsBeforeTheWindowEnd_AndKeepsTheCeilingSlot()
+ {
+ var sql = DarlingDataReader.TopQueriesHourlySql;
+
+ AssertEndIsExclusive(sql, 3, "get_top_queries_by_cpu hourly read");
+ Assert.Contains("bucket < $3$CEIL$", sql, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void McpTopProceduresHourlySql_StopsBeforeTheWindowEnd_AndKeepsTheCeilingSlot()
+ {
+ var sql = DarlingDataReader.TopProceduresHourlySql;
+
+ AssertEndIsExclusive(sql, 3, "get_top_procedures_by_cpu hourly read");
+ Assert.Contains("bucket < $3$CEIL$", sql, StringComparison.Ordinal);
+ }
+
+ /* Path 6: the unbucketed hourly duration trend, both databases-filtered and not, over both rollups. */
+ [Theory]
+ [InlineData(TimescaleSupport.QueryStatsHourlyView, false)]
+ [InlineData(TimescaleSupport.QueryStatsHourlyView, true)]
+ [InlineData(TimescaleSupport.ProcedureStatsHourlyView, false)]
+ [InlineData(TimescaleSupport.ProcedureStatsHourlyView, true)]
+ public void HourlyDurationTrendSql_StopsBeforeTheWindowEnd(string view, bool withDatabaseFilter)
+ {
+ var sql = DurationTrendRouting.BuildHourlyTrendSql(view, withDatabaseFilter);
+
+ AssertEndIsExclusive(sql, 3, "hourly duration trend");
+ Assert.Contains("FROM " + view, sql, StringComparison.Ordinal);
+ }
+
+ /* Path 7: the same trend gathered into $4-minute buckets. */
+ [Theory]
+ [InlineData(TimescaleSupport.QueryStatsHourlyView)]
+ [InlineData(TimescaleSupport.ProcedureStatsHourlyView)]
+ public void BucketedHourlyDurationTrendSql_StopsBeforeTheWindowEnd(string view)
+ {
+ var sql = DurationTrendRouting.BuildBucketedHourlyTrendSql(view);
+
+ AssertEndIsExclusive(sql, 3, "bucketed hourly duration trend");
+ Assert.Contains("FROM " + view, sql, StringComparison.Ordinal);
+ }
+
+ /* Path 8: the Query Store rollup-routed trend. Its rollup arm keeps the partition seam, `bucket < $4`, next to
+ the window end, so a bucket is read only when it starts before both. */
+ [Theory]
+ [InlineData(false)]
+ [InlineData(true)]
+ public void QueryStoreRollupTrendSql_StopsBeforeTheWindowEnd_AndKeepsTheSeam(bool withDatabaseFilter)
+ {
+ var sql = QueryStoreTrendRouting.BuildRollupTrendSql(withDatabaseFilter);
+
+ AssertEndIsExclusive(sql, 3, "Query Store rollup trend");
+ Assert.Matches(@"\bbucket\s*<\s*\$4(?!\d)", sql);
+ }
+
+ /* Path 9: one query's hourly history, over the legacy rollup and its interval-honest successor. */
+ [Theory]
+ [InlineData(TimescaleSupport.QueryStatsHourlyView)]
+ [InlineData("query_stats_interval_hourly")]
+ public void QueryHistoryHourlySql_StopsBeforeTheWindowEnd(string relation)
+ {
+ var sql = DarlingTrendReader.QueryHistoryHourlySqlFor(relation);
+
+ AssertEndIsExclusive(sql, 5, "query history hourly read");
+ Assert.Contains("bucket >= $4", sql, StringComparison.Ordinal);
+ AssertEndIsExclusive(DarlingTrendReader.QueryHistoryHourlySql, 5, "query history hourly constant");
+ }
+
+ /* Path 5: a Custom View panel. The compiler picks the rollup by the window's age; a rollup route bounds its
+ `bucket` column with `<`, and a raw route keeps `<=` on the raw time column, where a sample stamped at the
+ end still belongs to the window. */
+ private static readonly DateTime WindowEnd = new(2026, 7, 18, 6, 0, 0, DateTimeKind.Utc);
+
+ private const string QueryStatsPanel =
+ "{\"source\":\"query_stats\",\"measure\":\"query_worker_us\",\"aggregate\":\"sum\",\"timeBucket\":\"hour\",\"viz\":\"line\"}";
+
+ private static ComposeCompiled CompileOverWindowAged(int daysOld)
+ {
+ var (plan, planError) = ComposeSpec.TryParsePanel((JsonObject)JsonNode.Parse(QueryStatsPanel)!, Array.Empty());
+ Assert.True(planError is null, planError);
+ var start = WindowEnd.AddDays(-daysOld);
+
+ /* EndUtc doubles as "now" for the tier decision, exactly as the compiler pins elsewhere do. */
+ var (compiled, error) = ComposeCompiler.Compile(
+ plan!, new ComposeRunContext(null, start, WindowEnd, ComposeRunContext.NoVariables, RollupAvailability.All, WindowEnd, RollupCoverage.Unknown));
+ Assert.True(error is null, error);
+ Assert.NotNull(compiled);
+ return compiled!;
+ }
+
+ [Fact]
+ public void ComposeHourlyRoute_StopsBeforeTheWindowEnd()
+ {
+ var compiled = CompileOverWindowAged(10);
+
+ Assert.Equal(ComposeSourceTier.Hourly, compiled.Route.Tier);
+ Assert.Contains("FROM collect.query_stats_hourly AS f", compiled.Sql, StringComparison.Ordinal);
+ Assert.Contains("f.bucket >= $1", compiled.Sql, StringComparison.Ordinal);
+ Assert.Contains("f.bucket < $2", compiled.Sql, StringComparison.Ordinal);
+ Assert.DoesNotContain("f.bucket <=", compiled.Sql, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void ComposeDailyRoute_StopsBeforeTheWindowEnd()
+ {
+ var compiled = CompileOverWindowAged(120);
+
+ Assert.Equal(ComposeSourceTier.Daily, compiled.Route.Tier);
+ Assert.Contains("FROM collect.query_stats_daily AS f", compiled.Sql, StringComparison.Ordinal);
+ Assert.Contains("f.bucket >= $1", compiled.Sql, StringComparison.Ordinal);
+ Assert.Contains("f.bucket < $2", compiled.Sql, StringComparison.Ordinal);
+ Assert.DoesNotContain("f.bucket <=", compiled.Sql, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void ComposeRawRoute_KeepsTheInclusiveEndOnItsRawTimeColumn()
+ {
+ var compiled = CompileOverWindowAged(1);
+
+ Assert.Equal(ComposeSourceTier.Raw, compiled.Route.Tier);
+ Assert.Contains("FROM collect.query_stats AS f", compiled.Sql, StringComparison.Ordinal);
+ Assert.Contains("f.collection_time >= $1", compiled.Sql, StringComparison.Ordinal);
+ Assert.Contains("f.collection_time <= $2", compiled.Sql, StringComparison.Ordinal);
+ Assert.DoesNotContain("f.collection_time < $2", compiled.Sql, StringComparison.Ordinal);
+ }
+
+ /* The three first-bucket probes are NOT part of the change: they only locate the first bucket the window
+ holds (where the served span starts), and a bucket at the very end can be that first bucket only when the
+ window holds no other, in which case the read returns no rows either way. */
+ [Fact]
+ public void FirstBucketProbes_AreUnchanged()
+ {
+ Assert.Contains("f.bucket <= $3$CEIL$", DarlingDataReader.HourlyFirstBucketSql, StringComparison.Ordinal);
+ Assert.Contains("f.bucket <= $3$CEIL$", DarlingDataReader.HourlyFirstBucketSingleRelationSql, StringComparison.Ordinal);
+ }
+}
diff --git a/Darling/PerformanceMonitor.Darling.Service/Compose/ComposeCompiler.cs b/Darling/PerformanceMonitor.Darling.Service/Compose/ComposeCompiler.cs
index 9ad6c755f..db0e240ad 100644
--- a/Darling/PerformanceMonitor.Darling.Service/Compose/ComposeCompiler.cs
+++ b/Darling/PerformanceMonitor.Darling.Service/Compose/ComposeCompiler.cs
@@ -339,8 +339,13 @@ object_name resolves as m.object_name either way. */
}
}
+ /* A rollup bucket is stamped at its START, so a CAGG route's window end is EXCLUSIVE: with the end
+ exactly on a bucket start, `bucket <= end` would also take the whole hour (or, on the daily tier,
+ the whole day) that begins there, which lies after the window. A raw route stamps a sample when it
+ was taken, so a sample at the end still counts and it keeps `<=`. */
+ var endOperator = route.IsCagg ? " < " : " <= ";
sql.Append(indent).Append("WHERE ").Append(FactAlias).Append('.').Append(timeColumn).Append(" >= ").Append(startParam).Append('\n');
- sql.Append(indent).Append(" AND ").Append(FactAlias).Append('.').Append(timeColumn).Append(" <= ").Append(endParam).Append('\n');
+ sql.Append(indent).Append(" AND ").Append(FactAlias).Append('.').Append(timeColumn).Append(endOperator).Append(endParam).Append('\n');
if (hasServerScope)
{
sql.Append(indent).Append(" AND ").Append(FactAlias).Append(".server_name = ANY(").Append(serverScopeParam).Append(")\n");
diff --git a/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingDataReader.cs b/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingDataReader.cs
index 2ce997714..d39bcb5c6 100644
--- a/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingDataReader.cs
+++ b/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingDataReader.cs
@@ -1215,7 +1215,9 @@ ORDER BY r.total_cpu_us DESC
/// The rollup's min/max columns are per-collection sums, not per-execution extremes, so none are selected. $FROM$ is a
/// PLACEHOLDER, substituted (string.Replace, not string.Format — the SQL text otherwise contains braces)
/// with the FROM-clause item returns for this window at
- /// call time — never a literal relation name. $1 server_id, $2/$3 window (naive UTC), $4 top, $5 database
+ /// call time — never a literal relation name. $1 server_id, $2/$3 window (naive UTC; $3 is EXCLUSIVE — a
+ /// bucket is stamped at its START, so the bucket that begins at $3 lies after the window and is not read),
+ /// $4 top, $5 database
/// filter (NULL = all), $6 the materialization ceiling (naive UTC), bound only when the ceiling is known.
/// $CEIL$ becomes AND f.bucket < $6 or nothing. min_dop and host-object grouping need columns only raw carries, so
/// a read that sets either never reaches this const (it is forced to raw) and it takes no $6.
@@ -1232,7 +1234,7 @@ WITH ranked AS (
FROM $FROM$
WHERE server_id = $1
AND bucket >= $2
- AND bucket <= $3$CEIL$
+ AND bucket < $3$CEIL$
AND ($5::text IS NULL OR database_name = $5)
GROUP BY database_name, query_hash
HAVING (SUM(execution_count_sum) > 0 OR SUM(elapsed_time_sum) > 0)
@@ -1600,7 +1602,9 @@ ORDER BY SUM(delta_worker_time) DESC
/// rollup's pre-summed bucket columns rather than per-collection deltas. $FROM$ is a PLACEHOLDER,
/// substituted (string.Replace, not string.Format) with the FROM-clause item
/// returns for this window at call time — never a literal
- /// relation name. $1 server_id, $2/$3 window (naive UTC), $4 top, $5 database filter (NULL = all).
+ /// relation name. $1 server_id, $2/$3 window (naive UTC; $3 is EXCLUSIVE — a bucket is stamped at its
+ /// START, so the bucket that begins at $3 lies after the window and is not read), $4 top, $5 database
+ /// filter (NULL = all).
///
public const string TopProceduresHourlySql = """
WITH ranked AS (
@@ -1614,7 +1618,7 @@ WITH ranked AS (
FROM $FROM$
WHERE server_id = $1
AND bucket >= $2
- AND bucket <= $3$CEIL$
+ AND bucket < $3$CEIL$
AND ($5::text IS NULL OR database_name = $5)
GROUP BY database_name, schema_name, object_name
HAVING (SUM(execution_count_sum) > 0 OR SUM(elapsed_time_sum) > 0)
diff --git a/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingTrendReader.cs b/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingTrendReader.cs
index 325e12092..7fce073b8 100644
--- a/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingTrendReader.cs
+++ b/Darling/PerformanceMonitor.Darling.Service/Mcp/DarlingTrendReader.cs
@@ -184,7 +184,7 @@ public static string QueryHistoryHourlySqlFor(string hourlyRelation)
AND database_name = $2
AND query_hash = $3
AND bucket >= $4
- AND bucket <= $5
+ AND bucket < $5
ORDER BY bucket
""";
}
diff --git a/Darling/PerformanceMonitor.Darling.Storage/DurationTrendRouting.cs b/Darling/PerformanceMonitor.Darling.Storage/DurationTrendRouting.cs
index a1c52db65..8a064ce98 100644
--- a/Darling/PerformanceMonitor.Darling.Storage/DurationTrendRouting.cs
+++ b/Darling/PerformanceMonitor.Darling.Storage/DurationTrendRouting.cs
@@ -158,7 +158,8 @@ public static (DateTime EffectiveStartUtc, bool Truncated) DescribeCoverage(Date
/// With , $4 is the viewer's guarded text[] database
/// filter (#1319) — both hourly rollups group by database_name, so the filter survives the routing.
/// Without it the text is the MCP reader's constant byte for byte, which is what its pin asserts.
- /// $1 server_id, $2/$3 window (naive UTC).
+ /// $1 server_id, $2/$3 window (naive UTC; $3 is EXCLUSIVE — a bucket is stamped at its START, so the hour
+ /// that begins at $3 lies after the window and is not read).
///
public static string BuildHourlyTrendSql(string hourlyView, bool withDatabaseFilter)
{
@@ -176,7 +177,7 @@ public static string BuildHourlyTrendSql(string hourlyView, bool withDatabaseFil
FROM {hourlyView}
WHERE server_id = $1
AND bucket >= $2
- AND bucket <= $3{filter}
+ AND bucket < $3{filter}
GROUP BY bucket
ORDER BY bucket
""";
@@ -342,7 +343,8 @@ ORDER BY 1
/// not read low for hours it does not contain. The rollup keeps hourly sums, not the collections inside them,
/// so the peak on this tier is the bucket's busiest HOUR — the finest grain the rollup holds, and at 60 minutes
/// the point's own rate; nor is any hour unrated, because its denominator is known. $1 server_id, $2/$3 window
- /// (naive UTC), $4 the bucket width in minutes.
+ /// (naive UTC; $3 is EXCLUSIVE — a bucket is stamped at its START, so the hour that begins at $3 lies after
+ /// the window and is not read), $4 the bucket width in minutes.
///
public static string BuildBucketedHourlyTrendSql(string hourlyView)
{
@@ -358,7 +360,7 @@ WITH hourly AS
FROM {hourlyView}
WHERE server_id = $1
AND bucket >= $2
- AND bucket <= $3
+ AND bucket < $3
GROUP BY bucket
)
SELECT
diff --git a/Darling/PerformanceMonitor.Darling.Storage/HourlyWindowEdges.cs b/Darling/PerformanceMonitor.Darling.Storage/HourlyWindowEdges.cs
index a4f736ff9..9806a9815 100644
--- a/Darling/PerformanceMonitor.Darling.Storage/HourlyWindowEdges.cs
+++ b/Darling/PerformanceMonitor.Darling.Storage/HourlyWindowEdges.cs
@@ -13,8 +13,9 @@ namespace PerformanceMonitor.Darling.Storage;
///
/// The hour-bucket edges of a read served from an hourly rollup. A bucket is keyed by its start hour and the
-/// read takes bucket >= start AND bucket <= end, so an unaligned window start drops the partial
-/// hour it falls inside, and the bucket holding the window end is counted whole, up to the top of the next hour.
+/// read takes bucket >= start AND bucket < end, so an unaligned window start drops the partial
+/// hour it falls inside, and an unaligned window end counts the bucket holding it whole, up to the top of the
+/// next hour. A window end exactly on the hour stops there: the hour that begins at the end is not read.
///
public static class HourlyWindowEdges
{
@@ -26,7 +27,8 @@ public static (DateTime Start, DateTime? End) ServedSpan(
{
var startHour = FloorHour(requestedStart);
var start = firstBucket ?? (startHour != requestedStart ? startHour.AddHours(1) : requestedStart);
- var wholeEnd = FloorHour(requestedEnd).AddHours(1);
+ var endHour = FloorHour(requestedEnd);
+ var wholeEnd = endHour != requestedEnd ? endHour.AddHours(1) : endHour;
DateTime? end = ceiling is null ? null : (ceiling.Value < wholeEnd ? ceiling.Value : wholeEnd);
return (start, end);
}
@@ -60,7 +62,7 @@ public static (DateTime Start, DateTime? End) ServedSpan(
? string.Create(CultureInfo.InvariantCulture, $"served from {served:o} to {ceiling.Value:o}; {cut}")
: cut);
}
- else
+ else if (endHour != requestedEnd)
{
parts.Add(string.Create(CultureInfo.InvariantCulture,
$"the bucket at {endHour:o} is included whole, so up to {servedEnd!.Value:o} is counted past as_of"));
diff --git a/Darling/PerformanceMonitor.Darling.Storage/QueryStoreTrendRouting.cs b/Darling/PerformanceMonitor.Darling.Storage/QueryStoreTrendRouting.cs
index 565f2de43..cfca491a8 100644
--- a/Darling/PerformanceMonitor.Darling.Storage/QueryStoreTrendRouting.cs
+++ b/Darling/PerformanceMonitor.Darling.Storage/QueryStoreTrendRouting.cs
@@ -210,7 +210,10 @@ public static async Task ResolveAsync(
/// The rollup-routed trend SQL. $1 server_id, $2/$3 window (naive UTC), $4 the raw boundary
/// (); with , $5 is
/// the viewer's guarded text[] database filter (#1319) on every arm — the corrected hourly
- /// carries database_name, so the filter survives the routing.
+ /// carries database_name, so the filter survives the routing. The rollup arm stops BEFORE $3
+ /// (bucket < $3): a bucket is stamped at its START, so with $3 exactly on a bucket start the
+ /// hour that begins there lies after the window. The raw arms stamp a point when it happened and keep
+ /// <=.
///
/// The partition seam. The rollup arm takes buckets strictly BELOW $4; the raw arms take
/// points at or ABOVE it. bucket < $4 is load-bearing rather than decorative: a refresh can
@@ -283,7 +286,7 @@ WITH rollup_points AS
FROM {TimescaleSupport.QueryStoreStatsCorrectedHourlyView}
WHERE server_id = $1
AND bucket >= $2
- AND bucket <= $3
+ AND bucket < $3
AND bucket < $4{rollupFilter}
GROUP BY bucket
),
diff --git a/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.ProcedureStats.cs b/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.ProcedureStats.cs
index b4e7ed683..fead53cf5 100644
--- a/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.ProcedureStats.cs
+++ b/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.ProcedureStats.cs
@@ -147,7 +147,9 @@ ORDER BY SUM(delta_elapsed_time) DESC
/// 's legacy pair, Daily clamped to Hourly (#4231).
/// An hourly-routed page carries only what the rollup has: object_type/sql_handle/
/// plan_handle/reads/writes/spills columns are unavailable and read as their defaults, exactly the
- /// same disclosure the MCP payload's tier_used/precision_note make. Use
+ /// same disclosure the MCP payload's tier_used/precision_note make. An hourly-routed page
+ /// also stops BEFORE (a bucket is stamped at its start, so an end on the hour
+ /// does not add the hour that begins there). Use
/// to also learn which tier answered.
///
public async Task> GetTopProceduresByCpuAsync(
@@ -191,28 +193,7 @@ private async Task> GetTopProceduresByCpuHourlyAsy
{
var fromClause = coverage.StitchedRelationSql(
TimescaleSupport.ProcedureStatsHourlyView, "f", startUtc, RollupCoverage.StitchTier.Hourly);
- var sql = $"""
- SELECT
- database_name,
- schema_name,
- object_name,
- CAST(SUM(execution_count_sum) AS bigint) AS total_executions,
- CAST(SUM(worker_time_sum) AS bigint) AS total_cpu_us,
- CAST(SUM(elapsed_time_sum) AS bigint) AS total_elapsed_us,
- MIN(worker_time_min) AS min_worker_time,
- MAX(worker_time_max) AS max_worker_time,
- MIN(elapsed_time_min) AS min_elapsed_time,
- MAX(elapsed_time_max) AS max_elapsed_time
- FROM {fromClause}
- WHERE server_id = $1
- AND bucket >= $2
- AND bucket <= $3
- AND ($5::text[] IS NULL OR database_name = ANY($5))
- GROUP BY database_name, schema_name, object_name
- HAVING (SUM(execution_count_sum) > 0 OR SUM(elapsed_time_sum) > 0)
- ORDER BY SUM(worker_time_sum) DESC
- LIMIT $4
- """;
+ var sql = BuildTopProceduresHourlySql(fromClause);
var rows = new List();
await using var command = _dataSource.CreateCommand(sql);
@@ -243,6 +224,34 @@ ORDER BY SUM(worker_time_sum) DESC
return rows;
}
+ /// The hourly-rollup arm's SQL over . A rollup bucket is stamped at
+ /// its START, so the window end is EXCLUSIVE (bucket < $3): a range whose To is 14:00 sums the
+ /// hours up to 13:00-14:00 and does not add the 14:00-15:00 hour that only begins at the end. (The raw arm
+ /// stamps a sample when it was taken, so it keeps <= on collection_time.) Split out so a
+ /// test can read the text.
+ internal static string BuildTopProceduresHourlySql(string fromClause) => $"""
+ SELECT
+ database_name,
+ schema_name,
+ object_name,
+ CAST(SUM(execution_count_sum) AS bigint) AS total_executions,
+ CAST(SUM(worker_time_sum) AS bigint) AS total_cpu_us,
+ CAST(SUM(elapsed_time_sum) AS bigint) AS total_elapsed_us,
+ MIN(worker_time_min) AS min_worker_time,
+ MAX(worker_time_max) AS max_worker_time,
+ MIN(elapsed_time_min) AS min_elapsed_time,
+ MAX(elapsed_time_max) AS max_elapsed_time
+ FROM {fromClause}
+ WHERE server_id = $1
+ AND bucket >= $2
+ AND bucket < $3
+ AND ($5::text[] IS NULL OR database_name = ANY($5))
+ GROUP BY database_name, schema_name, object_name
+ HAVING (SUM(execution_count_sum) > 0 OR SUM(elapsed_time_sum) > 0)
+ ORDER BY SUM(worker_time_sum) DESC
+ LIMIT $4
+ """;
+
/// The Raw-tier read, unchanged — what ran before #4231
/// stage 3b added the hourly arm.
private async Task> GetTopProceduresByCpuRawAsync(
diff --git a/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.QueryStats.cs b/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.QueryStats.cs
index dcda1a674..a4e5dc2de 100644
--- a/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.QueryStats.cs
+++ b/Darling/PerformanceMonitor.Darling.Viewer/ViewerDataService.QueryStats.cs
@@ -326,7 +326,9 @@ ORDER BY r.total_elapsed_us DESC
/// 's legacy pair, Daily clamped to Hourly (out of scope for this lane).
/// An hourly-routed page carries only what the rollup has: host_object_name/
/// module_*/DOP/grant/spill/thread columns are unavailable and read as their defaults, exactly the
- /// same disclosure the MCP payload's tier_used/precision_note make. Use
+ /// same disclosure the MCP payload's tier_used/precision_note make. An hourly-routed page
+ /// also stops BEFORE (a bucket is stamped at its start, so an end on the hour
+ /// does not add the hour that begins there). Use
/// to also learn which tier answered.
///
public async Task> GetTopQueriesByCpuAsync(
@@ -369,25 +371,7 @@ private async Task> GetTopQueriesByCpuHourlyAsync(
{
var fromClause = coverage.StitchedRelationSql(
TimescaleSupport.QueryStatsHourlyView, "f", startUtc, RollupCoverage.StitchTier.Hourly);
- var sql = $"""
- SELECT
- database_name,
- query_hash,
- CAST(SUM(execution_count_sum) AS bigint) AS total_executions,
- CAST(SUM(worker_time_sum) AS bigint) AS total_cpu_us,
- CAST(SUM(elapsed_time_sum) AS bigint) AS total_elapsed_us,
- MIN(worker_time_min) AS min_worker_time,
- MAX(worker_time_max) AS max_worker_time
- FROM {fromClause}
- WHERE server_id = $1
- AND bucket >= $2
- AND bucket <= $3
- AND ($5::text[] IS NULL OR database_name = ANY($5))
- GROUP BY database_name, query_hash
- HAVING (SUM(execution_count_sum) > 0 OR SUM(elapsed_time_sum) > 0)
- ORDER BY SUM(elapsed_time_sum) DESC
- LIMIT $4
- """;
+ var sql = BuildTopQueriesHourlySql(fromClause);
var ranked = new List<(string Database, string QueryHash, long TotalExecutions, long TotalCpuUs, long TotalElapsedUs, long MinWorkerTime, long MaxWorkerTime)>();
await using (var command = _dataSource.CreateCommand(sql))
@@ -446,6 +430,31 @@ ORDER BY SUM(elapsed_time_sum) DESC
return rows;
}
+ /// The hourly-rollup arm's SQL over . A rollup bucket is stamped at
+ /// its START, so the window end is EXCLUSIVE (bucket < $3): a range whose To is 14:00 sums the
+ /// hours up to 13:00-14:00 and does not add the 14:00-15:00 hour that only begins at the end. (The raw arm
+ /// stamps a sample when it was taken, so it keeps <= on collection_time.) Split out so a
+ /// test can read the text.
+ internal static string BuildTopQueriesHourlySql(string fromClause) => $"""
+ SELECT
+ database_name,
+ query_hash,
+ CAST(SUM(execution_count_sum) AS bigint) AS total_executions,
+ CAST(SUM(worker_time_sum) AS bigint) AS total_cpu_us,
+ CAST(SUM(elapsed_time_sum) AS bigint) AS total_elapsed_us,
+ MIN(worker_time_min) AS min_worker_time,
+ MAX(worker_time_max) AS max_worker_time
+ FROM {fromClause}
+ WHERE server_id = $1
+ AND bucket >= $2
+ AND bucket < $3
+ AND ($5::text[] IS NULL OR database_name = ANY($5))
+ GROUP BY database_name, query_hash
+ HAVING (SUM(execution_count_sum) > 0 OR SUM(elapsed_time_sum) > 0)
+ ORDER BY SUM(elapsed_time_sum) DESC
+ LIMIT $4
+ """;
+
/// The Raw-tier read, unchanged — what ran before #4231
/// stage 3 added the hourly arm.
private async Task> GetTopQueriesByCpuRawAsync(
From 0dfc7135eb96173e68e69783f1af216bf77336f6 Mon Sep 17 00:00:00 2001
From: Erik Darling <2136037+erikdarlingdata@users.noreply.github.com>
Date: Wed, 30 Sep 2026 06:35:40 -0400
Subject: [PATCH 2/3] Live rollup pins: the window's hour count and the Custom
View's millisecond unit (#4859)
---
Darling/Darling.Tests/RollupWindowEndBoundLiveTests.cs | 9 +++++----
Darling/Darling.Tests/ViewerTrendRoutingPortTests.cs | 3 ++-
2 files changed, 7 insertions(+), 5 deletions(-)
diff --git a/Darling/Darling.Tests/RollupWindowEndBoundLiveTests.cs b/Darling/Darling.Tests/RollupWindowEndBoundLiveTests.cs
index 41baf1ddf..a7196f306 100644
--- a/Darling/Darling.Tests/RollupWindowEndBoundLiveTests.cs
+++ b/Darling/Darling.Tests/RollupWindowEndBoundLiveTests.cs
@@ -205,8 +205,8 @@ instances above cached the null hourly floor measured before the refresh. ──
///
/// A Custom View panel on the hourly route sums the hours up to the window end and stops before the hour that
- /// begins at it. The same three ends, the same rows: the panel's worker time is 1,000 per execution, so the
- /// sums name their hours (11,000 / 111,000 / 11,000).
+ /// begins at it. The same three ends, the same rows: the panel's worker time is 1,000 us per execution and it
+ /// reports the measure's default unit, ms, so the sums name their hours (11 / 111 / 11).
///
[Fact]
public async Task CustomViewPanel_OnTheHourlyRoute_StopsAtTheWindowEnd()
@@ -278,7 +278,8 @@ public async Task CustomViewPanel_OnTheHourlyRoute_StopsAtTheWindowEnd()
await using var reader = await command.ExecuteReaderAsync(ct);
Assert.True(await reader.ReadAsync(ct));
var sum = Convert.ToDouble(reader.GetValue(reader.FieldCount - 1));
- Assert.True(expected * 1000d == sum, $"Custom View sum of worker time, {label}: {sum:R}");
+ /* The panel reports the measure's default unit, ms, so 1,000 us per execution reads as 1. */
+ Assert.True(expected == sum, $"Custom View sum of worker time, {label}: {sum:R}");
}
bodySucceeded = true;
@@ -322,7 +323,7 @@ private static async Task> ReadFirstColumnTimesAsync(
return times;
}
- /// One query_stats row: worker time is 1,000 per execution, a real (nonzero) sample interval so the
+ /// One query_stats row: worker time is 1,000 us per execution, a real (nonzero) sample interval so the
/// hourly rollup admits it.
private static async Task PlantQueryAsync(
NpgsqlConnection connection, CancellationToken ct, int serverId, string serverName, DateTime at, long executions)
diff --git a/Darling/Darling.Tests/ViewerTrendRoutingPortTests.cs b/Darling/Darling.Tests/ViewerTrendRoutingPortTests.cs
index db4d7628f..5b9b05cb1 100644
--- a/Darling/Darling.Tests/ViewerTrendRoutingPortTests.cs
+++ b/Darling/Darling.Tests/ViewerTrendRoutingPortTests.cs
@@ -551,7 +551,8 @@ await RollupBackfill.RunSliceAsync(
Assert.Equal(8000.0 / 3600.0, p.Value, 6);
Assert.Equal(0, p.ExecutionCount);
});
- Assert.Equal(7 * 24 + 1, routed.Points.Count);
+ /* The window's 168 hour buckets; the bucket that starts at the window end lies after it. */
+ Assert.Equal(7 * 24, routed.Points.Count);
/* The database filter survives the routing (#1319 on the rollup). */
var filtered = await viewer.GetQueryDurationTrendAsync(ServerId, start, end, new[] { "DbA" }, nowUtc: now, cancellationToken: ct);
From 51a90079b6c0319c6728ae769e8f75c2ac3d1cf2 Mon Sep 17 00:00:00 2001
From: Erik Darling <2136037+erikdarlingdata@users.noreply.github.com>
Date: Wed, 30 Sep 2026 06:55:45 -0400
Subject: [PATCH 3/3] cpu_attribution pin: an on-the-hour as_of divides by
exactly the hours before it (#4859)
---
.../HourlyAttributionSpanTests.cs | 25 +++++++++++++++++++
1 file changed, 25 insertions(+)
diff --git a/Darling/Darling.Tests/HourlyAttributionSpanTests.cs b/Darling/Darling.Tests/HourlyAttributionSpanTests.cs
index 9b1a124fe..b459ebb15 100644
--- a/Darling/Darling.Tests/HourlyAttributionSpanTests.cs
+++ b/Darling/Darling.Tests/HourlyAttributionSpanTests.cs
@@ -48,6 +48,31 @@ public void Hourly_ShareDividesByTheServedSpan_NotTheRequestedWindow()
Assert.NotEqual(served.AttributedCpuRatio, requested.AttributedCpuRatio);
}
+ [Fact]
+ public void Hourly_OnTheHourAsOf_DividesByExactlyTheHoursBeforeIt()
+ {
+ // An on-the-hour as_of divides N hours of CPU by N hours, not N+1 (#4859): the read stops before the hour
+ // that begins at as_of, so that hour is not in the denominator. The ceiling sits above as_of, so it does
+ // not cut the span.
+ var start = Day.AddHours(10);
+ var asOf = Day.AddHours(14);
+ var (from, to, _) = DarlingMcpDataTools.HourlyAttributionSpan(true, start, asOf, Day.AddHours(10), Day.AddHours(18));
+
+ Assert.Equal(Day.AddHours(10), from);
+ Assert.Equal(Day.AddHours(14), to);
+
+ // 25% of 8 cores over exactly 4 hours = 28,800 CPU-seconds; 21,600 ranked seconds is 0.75.
+ var served = CpuAttribution.Compute(21600, from, to, 240, from, to, 25, 8);
+ Assert.Equal(28800, served.SqlCpuSecondsInWindow);
+ Assert.Equal(0.75, served.AttributedCpuRatio);
+
+ // Contrast: an as_of 20 minutes into the hour still counts the bucket holding it whole, so the span ends
+ // at the top of the next hour.
+ var (_, unalignedTo, _) = DarlingMcpDataTools.HourlyAttributionSpan(
+ true, start, Day.AddHours(14).AddMinutes(20), Day.AddHours(10), Day.AddHours(18));
+ Assert.Equal(Day.AddHours(15), unalignedTo);
+ }
+
[Fact]
public void Hourly_NullCeiling_SaysUnknown_NotThatNoSpanWasServed()
{