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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 73 additions & 0 deletions Darling/Darling.Tests/DarlingAnomalyBaselineTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -1066,6 +1066,79 @@ await LiveStoreCleanup.RunAsync(connectionString!, bodySucceeded, async (cleanup
}
}

/// <summary>
/// #4248: IoLatency reads the RAW file_io_stats hypertable at per-file grain over the 30-day window, so its
/// cache key is the UTC DAY, not the hour (PgBaselineProvider.IsDailyCacheMetric) — two calls hours apart on
/// the same UTC day share the one compute, and a call on the next UTC day recomputes. Proven by counting the
/// baseline reads Npgsql actually executes (CommandCapture), #3941's own live-pin technique.
/// </summary>
[Fact]
public async Task EndToEnd_IoLatencyArm_TwoCallsHoursApartOnOneDay_ShareOneCompute_NextDayRecomputes_AgainstDevPostgres()
{
var connectionString = Environment.GetEnvironmentVariable("DARLING_TEST_PG");
Assert.SkipWhen(string.IsNullOrEmpty(connectionString),
"Set DARLING_TEST_PG to a Postgres connection string to run the live IO-arm day-cache test.");

var ct = TestContext.Current.CancellationToken;
const int ioServerId = TestServerId + 5; // own id — this test cleans its own rows

using var connection = new NpgsqlConnection(connectionString);
await connection.OpenAsync(ct);
await PgMigrations.MigrateAsync(connection, ct);

await using (var cleanup = new NpgsqlCommand($"DELETE FROM file_io_stats WHERE server_id = {ioServerId};", connection))
{
await cleanup.ExecuteNonQueryAsync(ct);
}

await using var postgres = NpgsqlDataSource.Create(connectionString!);
var bodySucceeded = false;
try
{
var day = DateTime.UtcNow.Date.AddDays(-8);
while (day.DayOfWeek != DayOfWeek.Monday) day = day.AddDays(-1);
var historyStart = DateTime.SpecifyKind(day.AddHours(10), DateTimeKind.Unspecified);

for (var i = 0; i < 5; i++)
{
await InsertAsync(connection,
"INSERT INTO file_io_stats (collection_id, collection_time, server_id, server_name, delta_reads, delta_writes, delta_stall_read_ms) VALUES ($1, $2, $3, $4, $5, $6, $7)",
(long)(200 + i), historyStart.AddMinutes(5 * i), ioServerId, "IO-DAILY-CACHE",
10L, 0L, (long)(10 * (i + 1)));
}

var provider = new PgBaselineProvider(postgres);
var analysisDay = historyStart.AddDays(7).Date; // the 30-day window's end (#4248: midnight, not the hour)

var (morning, firstReads) = await CommandCapture.CountBaselineReadsAsync(
() => provider.GetBaselineAsync(ioServerId, MetricNames.IoLatency, analysisDay.AddHours(1), ct));
Assert.Equal(1, firstReads);
Assert.True(morning.SampleCount > 0, "the seed produced no baseline — the comparison would prove nothing");

/* Nineteen hours later (over CacheTtl's one hour), same UTC day: the #4248 pin — no second read. */
var (afternoon, secondReads) = await CommandCapture.CountBaselineReadsAsync(
() => provider.GetBaselineAsync(ioServerId, MetricNames.IoLatency, analysisDay.AddHours(20), ct));
Assert.Equal(0, secondReads);
Assert.Equal(morning.SampleCount, afternoon.SampleCount);
Assert.Equal(morning.Median, afternoon.Median);

/* The next UTC day is a different window end (midnight moved), so a fresh compute. */
var (_, nextDayReads) = await CommandCapture.CountBaselineReadsAsync(
() => provider.GetBaselineAsync(ioServerId, MetricNames.IoLatency, analysisDay.AddDays(1).AddHours(1), ct));
Assert.Equal(1, nextDayReads);

bodySucceeded = true;
}
finally
{
await LiveStoreCleanup.RunAsync(connectionString!, bodySucceeded, async (cleanup, cleanupCt) =>
{
using var command = new NpgsqlCommand($"DELETE FROM file_io_stats WHERE server_id = {ioServerId};", cleanup);
await command.ExecuteNonQueryAsync(cleanupCt);
});
}
}

/// <summary>
/// #3653 (A8, first slice) proven live through the PG detector: the I/O read hands the shared gate the
/// per-file-row PEAK and MEAN, and the gate fires only when both clear. History: three Mondays at 10:00,
Expand Down
2 changes: 1 addition & 1 deletion Darling/Darling.Tests/PgTargetClockTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ public void TheClockRead_IsAProtectedVirtualSeam_OverriddenOnceByThePostgresProv
var baseCode = CSharpSourceWalker.StripCommentsAndStrings(baseSource);
Assert.Single(Regex.Matches(baseCode, @"await ReadServerClockAsync\("));
Assert.Contains("await ReadServerClockAsync(connection, serverId, AsNaive(windowEnd), cancellationToken)", baseCode, StringComparison.Ordinal);
Assert.Contains("var windowEnd = RoundedHour(analysisTime);", baseCode, StringComparison.Ordinal);
Assert.Contains("var windowEnd = RoundedKeyTime(metricName, analysisTime);", baseCode, StringComparison.Ordinal);
Assert.DoesNotContain("DateTime.UtcNow", CSharpSourceWalker.StripCommentsAndStrings(
RepoFile.ReadRepoFile("Darling", "PerformanceMonitor.Darling.Analysis", "PgTargetBaselineProvider.Clock.cs")), StringComparison.Ordinal);
}
Expand Down
9 changes: 7 additions & 2 deletions Darling/PerformanceMonitor.Darling.Analysis/BaselineCache.cs
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,12 @@ namespace PerformanceMonitor.Darling.Analysis;
/// <item><b>Time.</b> <see cref="PgBaselineProvider.CacheTtl"/> after the compute an entry is dead whatever its hour —
/// the bound on anything that did move inside a settled window (a late row, a purge, the target's clock changing zone,
/// which is re-keyed on the next compute exactly as before). Dead entries are swept, so the tier holds at most one
/// TTL's worth of computes.</item>
/// TTL's worth of computes. <b>Except (#4248):</b> a successful compute of a daily-cache metric
/// (<see cref="PgBaselineProvider.IsDailyCacheMetric"/>) is never eroded by the TTL
/// (<see cref="PgBaselineProvider.CachedBaseline.FreshUntilUtc"/> is a 24-hour backstop, not the real bound) — the
/// EntryKey's analysis-day component is, exactly as the hourly case's analysis-hour component always was, so this
/// tier keeps answering every lookup for the rest of that UTC day. See <see cref="PgBaselineProvider.IsFresh"/>,
/// which this tier's live check and sweep both call, so they never evict a still-fresh daily entry early.</item>
/// <item><b>Failure.</b> Only a SUCCESSFUL compute is shared. A failed one (a timeout, a role that cannot read a
/// relation) stays in the failing provider's own tier, where it has always meant "no baseline for this pass" — one
/// caller's timeout never blanks another caller's baselines, and a success from any caller beats a local failure.</item>
Expand Down Expand Up @@ -103,7 +108,7 @@ internal void Invalidate(int serverId)
internal void Clear() => _entries.Clear();

private static bool IsLive(PgBaselineProvider.CachedBaseline entry, DateTime nowUtc)
=> nowUtc - entry.RealTime < PgBaselineProvider.CacheTtl;
=> PgBaselineProvider.IsFresh(entry, nowUtc);

/// <summary>At most once per quarter of <see cref="PgBaselineProvider.CacheTtl"/>, drops the entries no lookup can
/// take any more, so the tier holds little beyond one TTL of computes. Without it the tier would grow for the life of
Expand Down
71 changes: 64 additions & 7 deletions Darling/PerformanceMonitor.Darling.Analysis/PgBaselineProvider.cs
Original file line number Diff line number Diff line change
Expand Up @@ -354,7 +354,7 @@ private async Task<CachedBaseline> GetOrComputeBaselinesAsync(
int serverId, string metricName, DateTime analysisTime, CancellationToken cancellationToken)
{
var cacheKey = CacheKeyFor(serverId, metricName, key: null);
var roundedHour = RoundedHour(analysisTime);
var roundedHour = RoundedKeyTime(metricName, analysisTime);

if (TryGetFresh(serverId, cacheKey, roundedHour, out var cached))
{
Expand All @@ -363,11 +363,14 @@ private async Task<CachedBaseline> GetOrComputeBaselinesAsync(

var (byMember, clock, utcOffsetMinutes, timeZoneId) = await ComputeBaselinesAsync(serverId, metricName, keys: null, analysisTime, cancellationToken);

var buckets = BucketsOf(byMember, UnkeyedMember);
var realTime = DateTime.UtcNow;
var entry = new CachedBaseline
{
ComputedAt = roundedHour,
RealTime = DateTime.UtcNow,
Buckets = BucketsOf(byMember, UnkeyedMember),
RealTime = realTime,
Buckets = buckets,
FreshUntilUtc = buckets is not null && IsDailyCacheMetric(metricName) ? realTime.AddDays(1) : null,
Clock = clock,
UtcOffsetMinutes = utcOffsetMinutes,
TimeZoneId = timeZoneId
Expand Down Expand Up @@ -397,7 +400,7 @@ private async Task<CachedBaseline> GetOrComputeBaselinesAsync(
private async Task<Dictionary<string, CachedBaseline>> GetOrComputeKeyedBaselinesAsync(
int serverId, string metricName, IReadOnlyCollection<string> keys, DateTime analysisTime, CancellationToken cancellationToken)
{
var roundedHour = RoundedHour(analysisTime);
var roundedHour = RoundedKeyTime(metricName, analysisTime);
var entries = new Dictionary<string, CachedBaseline>(StringComparer.Ordinal);
var misses = new List<string>();
var asked = new HashSet<string>(StringComparer.Ordinal);
Expand Down Expand Up @@ -435,11 +438,13 @@ private async Task<Dictionary<string, CachedBaseline>> GetOrComputeKeyedBaseline
for (var member = 1; member <= set.Length; member++)
{
var key = set[member - 1];
var memberBuckets = BucketsOf(byMember, member);
var entry = new CachedBaseline
{
ComputedAt = roundedHour,
RealTime = computedAt,
Buckets = BucketsOf(byMember, member),
Buckets = memberBuckets,
FreshUntilUtc = memberBuckets is not null && IsDailyCacheMetric(metricName) ? computedAt.AddDays(1) : null,
Clock = clock,
UtcOffsetMinutes = utcOffsetMinutes,
TimeZoneId = timeZoneId,
Expand All @@ -460,6 +465,35 @@ private async Task<Dictionary<string, CachedBaseline>> GetOrComputeKeyedBaseline
internal static DateTime RoundedHour(DateTime analysisTime)
=> new(analysisTime.Year, analysisTime.Month, analysisTime.Day, analysisTime.Hour, 0, 0);

/// <summary>Midnight UTC of <paramref name="analysisTime"/>'s day (#4248) — the key time and window end for an
/// arm that reads a RAW hypertable at full grain over the 30-day window (<see cref="IsDailyCacheMetric"/>),
/// playing <see cref="RoundedHour"/>'s role at day grain instead of hour grain. The underlying rows move by
/// about 1/720 an hour, so an hourly key bought nothing but 24x the recomputes.</summary>
internal static DateTime RoundedDay(DateTime analysisTime)
=> new(analysisTime.Year, analysisTime.Month, analysisTime.Day, 0, 0, 0);

/// <summary>#4248: the two arms whose <c>clean</c> CTE reads a RAW hypertable (<c>cpu_utilization_stats</c>,
/// <c>file_io_stats</c> — the #1743 follow-up pair, see <see cref="GetBaselineQuery"/>'s remarks) rather than a
/// pre-aggregated <c>CREATE MATERIALIZED VIEW ... _baseline</c> supply. Both tables carry their own 30-day
/// service-side retention floor (<c>DarlingRetentionHorizons.BaselineServingRawCollectors</c>), so the 30-day
/// WINDOW does not change here — only the cache KEY's grain does, because a full-grain 30-day read is what made
/// an hourly recompute expensive (measured: ~50 MB of temp per <see cref="MetricNames.IoLatency"/> call). The
/// other seven <see cref="RobustTierScaffold"/> arms and the two event arms read an already-aggregated view —
/// far fewer rows for the same 30 days — and keep the hourly key. <c>PgTargetBaselineProvider</c>'s arms all
/// read raw PostgreSQL-target hypertables too (its own remarks say so), but they are NOT in this set: nobody
/// has ruled on a keyed-arm day-long lifetime against <see cref="KeyedBaselineCacheWarnCount"/>'s hourly-turnover
/// calibration, so that file is unchanged here.</summary>
internal static bool IsDailyCacheMetric(string metricName)
=> metricName is MetricNames.Cpu or MetricNames.IoLatency;

/// <summary>The cache key's time AND the compute's window end (#3941's invariant, restated for #4248): whichever
/// grain <paramref name="metricName"/> uses (<see cref="IsDailyCacheMetric"/>) — the day for a raw-hypertable
/// arm, the hour for every other one. Both call sites (the key and the window) read this ONE seam so they can
/// never diverge — a key that named a different instant than the window it was computed over would be sharing
/// an answer that is not the answer a fresh compute at that key would give.</summary>
internal static DateTime RoundedKeyTime(string metricName, DateTime analysisTime)
=> IsDailyCacheMetric(metricName) ? RoundedDay(analysisTime) : RoundedHour(analysisTime);

/// <summary>This provider's engine in the shared tier's key (#3941): a SQL Server series and a PostgreSQL-target
/// series of one server id are never each other's.</summary>
private string SharedKind => GetType().FullName ?? GetType().Name;
Expand All @@ -474,7 +508,7 @@ private bool TryGetFresh(int serverId, string cacheKey, DateTime roundedHour, [N
{
var local = _cache.TryGetValue(cacheKey, out var own)
&& own.ComputedAt == roundedHour
&& (DateTime.UtcNow - own.RealTime) < CacheTtl;
&& IsFresh(own, DateTime.UtcNow);
if (local && own!.Buckets is not null)
{
cached = own;
Expand All @@ -492,6 +526,22 @@ private bool TryGetFresh(int serverId, string cacheKey, DateTime roundedHour, [N
return local;
}

/// <summary>#4248: is <paramref name="entry"/> still good to return? A successful compute of a daily-cache
/// metric (<see cref="CachedBaseline.FreshUntilUtc"/> set, a rolling 24 hours from <see cref="CachedBaseline.RealTime"/>
/// — never eroded by <see cref="CacheTtl"/>) stays live for a full day of real time no matter when in the UTC
/// day it ran, so it is never the TTL that ends it: the CALLER'S key (<c>ComputedAt == roundedHour</c> in
/// <see cref="TryGetFresh"/> and <see cref="BaselineCache.TryGet"/>) already stops matching the instant the
/// requested day rolls over, which is what actually bounds an entry to "the rest of the UTC day it was computed
/// in" — the whole point, two calls an hour or more apart on the same day share the one compute. Every
/// hourly-cache metric, and a FAILED compute of ANY metric (<see cref="CachedBaseline.FreshUntilUtc"/> null — a
/// failure never earns the day-long trust), keep the original rolling <see cref="CachedBaseline.RealTime"/> plus
/// <see cref="CacheTtl"/> bound, so a timeout still retries within the hour. Shared by <see cref="BaselineCache"/>'s
/// live check and sweep, so the shared tier never evicts a still-fresh daily entry early.</summary>
internal static bool IsFresh(CachedBaseline entry, DateTime nowUtc)
=> entry.FreshUntilUtc is DateTime freshUntil
? nowUtc < freshUntil
: (nowUtc - entry.RealTime) < CacheTtl;

/// <summary>Files a compute in this provider's cache and, when it SUCCEEDED, in the shared tier (#3941). A failed
/// compute (null buckets) is this caller's "no baseline this pass" and nobody else's.</summary>
private void Store(int serverId, string cacheKey, CachedBaseline entry)
Expand Down Expand Up @@ -803,7 +853,7 @@ hour was answered from a window up to 59 minutes off its own. Ending the window
in the hour ask for the SAME rows — what lets the process share one compute between the scheduled pass,
analyze_server and compare_analysis (BaselineCache) without changing anyone's answer. The lookup still keys
the bucket on the analysis instant (LookUp); its hour-of-week is the hour's. */
var windowEnd = RoundedHour(analysisTime);
var windowEnd = RoundedKeyTime(metricName, analysisTime);
var windowStart = windowEnd.AddDays(-BaselineMath.BaselineWindowDays);
var clock = LocalClockWindow.Utc(windowEnd);

Expand Down Expand Up @@ -1508,6 +1558,13 @@ internal sealed class CachedBaseline
public DateTime RealTime { get; init; }
public Dictionary<(int HourOfDay, int DayOfWeek), BaselineBucket>? Buckets { get; init; }

/// <summary>#4248: for a successful compute of a daily-cache metric, <see cref="RealTime"/> plus 24 hours —
/// null for every hourly-cache metric and for a failed compute of any metric. See <see cref="IsFresh"/>,
/// the only reader: this bound alone would outlive the metric's own UTC day, but <c>ComputedAt</c>'s key
/// match already stops answering the moment that day ends, so in practice this is the "still trustworthy in
/// real time" backstop, not the day boundary itself.</summary>
public DateTime? FreshUntilUtc { get; init; }

/// <summary>The clock the buckets were keyed with (#3653 Q6) — the lookup must use the SAME one.</summary>
public LocalClockWindow Clock { get; init; } = LocalClockWindow.Utc(DateTime.MinValue);

Expand Down
Loading
Loading