feat(kpihistory): add hourly rollup fold + query + purge repository methods (plan #22 T3)

Add FoldHourlyRollupsAsync (in-memory grouped, per-metric gauge-last/rate-sum,
idempotent upsert on the series+hour key with null-ScopeKey equality, exclusive
upper hour bound so the in-progress hour is never folded), GetHourlySeriesAsync
(same KpiSeriesPoint contract as GetRawSeriesAsync), and PurgeRollupsOlderThanAsync
(one-hour-sliced batched DELETE mirroring PurgeOlderThanAsync). Adds 7 repo tests.

Claude-Session: https://claude.ai/code/session_01MtdgwpEeCUn6cUA5f1LMPj
This commit is contained in:
Joseph Doherty
2026-07-10 11:50:53 -04:00
parent ccdc649641
commit 09f67d1c65
3 changed files with 396 additions and 0 deletions
@@ -1,4 +1,5 @@
using System.Data.Common;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Diagnostics;
using ZB.MOM.WW.ScadaBridge.Commons.Entities.Kpi;
using ZB.MOM.WW.ScadaBridge.ConfigurationDatabase;
@@ -195,6 +196,200 @@ public class KpiHistoryRepositoryTests
Assert.Equal(new[] { 4d, 5d }, remaining.Select(p => p.Value).ToArray());
}
// ---- Hourly rollup fold / query / purge (plan #22 T3) ------------------------------
[Fact]
public async Task FoldHourlyRollupsAsync_GaugeMetric_FoldsToLastValuePerHour()
{
await using var ctx = NewContext();
var repo = new KpiHistoryRepository(ctx);
// "queueDepth" is a Gauge → last (latest-timestamp) value in the hour wins.
await repo.RecordSamplesAsync(new[]
{
Sample("NotificationOutbox", "queueDepth", "Global", null, value: 3, capturedAtUtc: Base.AddMinutes(5)),
Sample("NotificationOutbox", "queueDepth", "Global", null, value: 9, capturedAtUtc: Base.AddMinutes(59)),
Sample("NotificationOutbox", "queueDepth", "Global", null, value: 7, capturedAtUtc: Base.AddMinutes(30)),
});
// toHourUtc = Base + 1 h (next hour start) → the Base hour is complete and folded.
await repo.FoldHourlyRollupsAsync(Base, Base.AddHours(1));
var series = await repo.GetHourlySeriesAsync(
"NotificationOutbox", "queueDepth", "Global", scopeKey: null,
fromUtc: Base, toUtc: Base.AddHours(24));
var row = Assert.Single(series);
Assert.Equal(Base, row.BucketStartUtc);
Assert.Equal(9, row.Value); // last-value (max CapturedAtUtc)
}
[Fact]
public async Task FoldHourlyRollupsAsync_RateMetric_FoldsToSumPerHour_WithMinMaxCount()
{
await using var ctx = NewContext();
var repo = new KpiHistoryRepository(ctx);
// "deliveredLastInterval" is a Rate → sum of the hour's per-interval deltas.
await repo.RecordSamplesAsync(new[]
{
Sample("NotificationOutbox", "deliveredLastInterval", "Global", null, value: 2, capturedAtUtc: Base.AddMinutes(5)),
Sample("NotificationOutbox", "deliveredLastInterval", "Global", null, value: 5, capturedAtUtc: Base.AddMinutes(25)),
Sample("NotificationOutbox", "deliveredLastInterval", "Global", null, value: 4, capturedAtUtc: Base.AddMinutes(45)),
});
await repo.FoldHourlyRollupsAsync(Base, Base.AddHours(1));
// Read the entity directly to assert Min/Max/Count fidelity, not just Value.
var rollup = Assert.Single(await ctx.KpiRollupHourly.ToListAsync());
Assert.Equal(Base, rollup.HourStartUtc);
Assert.Equal(11, rollup.Value); // sum 2+5+4
Assert.Equal(2, rollup.MinValue);
Assert.Equal(5, rollup.MaxValue);
Assert.Equal(3, rollup.SampleCount);
}
[Fact]
public async Task FoldHourlyRollupsAsync_SplitsIntoSeparateHourBuckets_ForGlobalScope()
{
await using var ctx = NewContext();
var repo = new KpiHistoryRepository(ctx);
// Two hours of one Global (null ScopeKey) gauge series.
await repo.RecordSamplesAsync(new[]
{
Sample("SiteHealth", "connectionsUp", "Global", null, value: 4, capturedAtUtc: Base.AddMinutes(10)),
Sample("SiteHealth", "connectionsUp", "Global", null, value: 6, capturedAtUtc: Base.AddMinutes(50)),
Sample("SiteHealth", "connectionsUp", "Global", null, value: 2, capturedAtUtc: Base.AddHours(1).AddMinutes(20)),
});
// Fold two complete hours: [Base, Base+2h).
await repo.FoldHourlyRollupsAsync(Base, Base.AddHours(2));
var series = await repo.GetHourlySeriesAsync(
"SiteHealth", "connectionsUp", "Global", scopeKey: null,
fromUtc: Base, toUtc: Base.AddHours(24));
Assert.Equal(2, series.Count);
Assert.Equal(Base, series[0].BucketStartUtc);
Assert.Equal(6, series[0].Value); // hour 0 last-value
Assert.Equal(Base.AddHours(1), series[1].BucketStartUtc);
Assert.Equal(2, series[1].Value); // hour 1 last-value
}
[Fact]
public async Task FoldHourlyRollupsAsync_ExcludesIncompleteHour_AtExclusiveUpperBound()
{
await using var ctx = NewContext();
var repo = new KpiHistoryRepository(ctx);
await repo.RecordSamplesAsync(new[]
{
// Complete hour (Base) — folded.
Sample("NotificationOutbox", "queueDepth", "Global", null, value: 3, capturedAtUtc: Base.AddMinutes(30)),
// Next, in-progress hour — must NOT be folded when toHourUtc == Base+1h.
Sample("NotificationOutbox", "queueDepth", "Global", null, value: 99, capturedAtUtc: Base.AddHours(1).AddMinutes(15)),
});
await repo.FoldHourlyRollupsAsync(Base, Base.AddHours(1));
var series = await repo.GetHourlySeriesAsync(
"NotificationOutbox", "queueDepth", "Global", scopeKey: null,
fromUtc: Base, toUtc: Base.AddHours(24));
// Only the complete Base hour is present; the in-progress hour is excluded.
var row = Assert.Single(series);
Assert.Equal(Base, row.BucketStartUtc);
Assert.Equal(3, row.Value);
}
[Fact]
public async Task FoldHourlyRollupsAsync_IsIdempotent_AcrossRepeatedRunsOverSameWindow()
{
await using var ctx = NewContext();
var repo = new KpiHistoryRepository(ctx);
await repo.RecordSamplesAsync(new[]
{
Sample("NotificationOutbox", "deliveredLastInterval", "Global", null, value: 2, capturedAtUtc: Base.AddMinutes(5)),
Sample("NotificationOutbox", "deliveredLastInterval", "Global", null, value: 5, capturedAtUtc: Base.AddMinutes(45)),
});
await repo.FoldHourlyRollupsAsync(Base, Base.AddHours(1));
await repo.FoldHourlyRollupsAsync(Base, Base.AddHours(1)); // re-run same window
// Exactly one row for the (series, hour); the rate sum is not double-counted.
var rollups = await ctx.KpiRollupHourly.ToListAsync();
var row = Assert.Single(rollups);
Assert.Equal(7, row.Value); // still 2+5, not 14
Assert.Equal(2, row.SampleCount);
}
[Fact]
public async Task GetHourlySeriesAsync_ReturnsAscending_AndHonorsNullVsSiteScopeKey()
{
await using var ctx = NewContext();
var repo = new KpiHistoryRepository(ctx);
// Same metric, one Global (null key) hour and one Site-keyed hour, both foldable.
await repo.RecordSamplesAsync(new[]
{
Sample("SiteCallAudit", "buffered", "Global", null, value: 3, capturedAtUtc: Base.AddMinutes(10)),
Sample("SiteCallAudit", "buffered", "Global", null, value: 8, capturedAtUtc: Base.AddHours(1).AddMinutes(10)),
Sample("SiteCallAudit", "buffered", "Site", "plant-a", value: 42, capturedAtUtc: Base.AddMinutes(10)),
});
await repo.FoldHourlyRollupsAsync(Base, Base.AddHours(2));
var global = await repo.GetHourlySeriesAsync(
"SiteCallAudit", "buffered", "Global", scopeKey: null,
fromUtc: Base, toUtc: Base.AddHours(24));
// Two Global hours, ascending; the Site-keyed row must not leak in.
Assert.Equal(2, global.Count);
Assert.Equal(Base, global[0].BucketStartUtc);
Assert.Equal(3, global[0].Value);
Assert.Equal(Base.AddHours(1), global[1].BucketStartUtc);
Assert.Equal(8, global[1].Value);
var site = await repo.GetHourlySeriesAsync(
"SiteCallAudit", "buffered", "Site", scopeKey: "plant-a",
fromUtc: Base, toUtc: Base.AddHours(24));
Assert.Equal(42, Assert.Single(site).Value);
}
[Fact]
public async Task PurgeRollupsOlderThanAsync_DeletesOnlyRollupsOlderThanCutoff()
{
await using var ctx = NewContext();
var repo = new KpiHistoryRepository(ctx);
// Fold four distinct hours spanning several days.
await repo.RecordSamplesAsync(new[]
{
Sample("NotificationOutbox", "queueDepth", "Global", null, value: 1, capturedAtUtc: Base.AddDays(-10).AddMinutes(5)),
Sample("NotificationOutbox", "queueDepth", "Global", null, value: 2, capturedAtUtc: Base.AddDays(-8).AddMinutes(5)),
Sample("NotificationOutbox", "queueDepth", "Global", null, value: 3, capturedAtUtc: Base.AddDays(-7).AddMinutes(5)),
Sample("NotificationOutbox", "queueDepth", "Global", null, value: 4, capturedAtUtc: Base.AddDays(-1).AddMinutes(5)),
});
await repo.FoldHourlyRollupsAsync(Base.AddDays(-30), Base);
// Cutoff = the -7d hour start; strictly-older predicate keeps the ==cutoff row.
var cutoff = TruncateHour(Base.AddDays(-7));
await repo.PurgeRollupsOlderThanAsync(cutoff);
var remaining = await repo.GetHourlySeriesAsync(
"NotificationOutbox", "queueDepth", "Global", scopeKey: null,
fromUtc: Base.AddDays(-30), toUtc: Base.AddDays(30));
// The -7d (==cutoff) and -1d rollups survive; the -10d and -8d are purged.
Assert.Equal(new[] { 3d, 4d }, remaining.Select(p => p.Value).ToArray());
}
// Local mirror of the repo's private hour-truncation, for cutoff assertions.
private static DateTime TruncateHour(DateTime utc) =>
new(utc.Year, utc.Month, utc.Day, utc.Hour, 0, 0, DateTimeKind.Utc);
/// <summary>
/// EF command interceptor that counts the DELETE statements actually issued
/// to the store, so a test can prove the purge is sliced into multiple