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() {