Repository navigation
Cache the server-scoped watermark read across a runner's lifetime (#4197) - #4399
Merged
Merged
Conversation
…#4197 part b) Adds a per-row watermark accessor (WatermarkValueAccessor / NumericWatermarkValueAccessor, null by default) to ICollectorDefinition/CollectorDefinitionBase, declared on job_history, default_trace_events, system_health_events and memory_pressure_events. A null accessor opts a collector out of the cache entirely -- pg_cpu_utilization declares none, since Targets/RdsCpuIngestor.cs is a second writer under that collector name. RunAsync now checks ServerWatermarkCache before either store read (a hit skips both the timestamp and numeric watermark reads), seeds on a miss, and advances after WriteBatchAsync's COPY returns on the plain (non-fan-out) path from the batch's own max watermark value -- never before the write, and never on a zero-row batch. Any fault in RunAsync invalidates the (server, collector) entry via a thin try/catch wrapper (RunCoreAsync holds the prior body). OnServerReconnected drops every cached entry for that server; DarlingWorker's per-connect call site already covers a re-added server_id, since a re-add's first connection reaches the same OnConnectedAsync path as any other reconnect. cpu_utilization (the #3778 UTC-twin pair) stays uncached in this lane: batch-max wiring for the twin frame needs more read time than remained. Refs #4197.
…rt b) xUnit pins in Darling.Tests, net10.0-windows (cannot execute on this machine; each states its RED-on-dev reason in its own summary): miss returns null, seed round-trips value + numeric twin, advance is monotonic (older value leaves the cache unchanged, newer moves it forward and carries its own #3778 frame flag), per-server isolation, Invalidate drops one entry, InvalidateServer drops every entry for one server. Two more pin the definition-level opt-in/out default: pg_cpu_utilization declares no accessor (stays out of the cache), the four wired collectors do, cpu_utilization stays excluded until its twin wiring lands. Build: Darling.Tests and Lite.Tests both 0 Warning(s)/0 Error(s) (EnableWindowsTargeting=true). Refs #4197.
…4197 part b) - AdvanceServerWatermark now floors the batch max to PostgreSQL's microsecond storage resolution before caching it, matching CollectionTimeClock's own truncation. Without this the cache could hold a value with finer ticks than the store ever persisted, drifting from a fresh GetLastCollectedTimeAsync read by those sub-microsecond ticks (caught live by the equivalence pins). - Fixed pin bugs in ServerWatermarkCacheRunnerLiveTests found while running them live for the first time: pg_stat_statements SUM(calls) returns numeric, not bigint (added an explicit ::bigint cast); three pins asserted read counts against an EMPTY seed table, which doubles every cold-cache read via the probe/fallback pair regardless of the cache -- each now seeds one real row first so the assertion isolates the cache's own behaviour. - Updated TimeHonestyRungTests' and ServerWatermarkDispatchGateTests' source/IL scans for the #4197 part b extraction: the watermark-pair read block moved one indentation level deeper (into ResolveServerWatermarkAsync's own try), and the gate (ServerWatermarkIsDiscarded) now sits in RunCoreAsync while the read it guards sits one call away in ResolveServerWatermarkAsync -- the IL walk now tracks calls to the seam itself as reaching the read transitively. Refs #4197 (part b).
…part b) The IL-walk pin (ServerWatermarkDispatchGateTests) checks that every method body performing the server-scoped watermark read also calls ServerWatermarkIsDiscarded. The #4197 seam extraction moved the read into ResolveServerWatermarkAsync but left the gate call in RunCoreAsync, which passed the precomputed bool in as a parameter -- splitting the read and its gate across two bodies and failing the pin (Expected 2, Actual 1). Move ServerWatermarkIsDiscarded(definition, dispatchProbe) into ResolveServerWatermarkAsync itself, passing the dispatch probe instead of the precomputed bool. RunCoreAsync gets the discarded flag back in the tuple (it still needs it for hasCollectedBefore). Behaviour is unchanged; ServerWatermarkDispatchGateTests.cs is now unchanged against dev. ServerWatermarkCacheRunnerLiveTests' ResolveAsync helper updated for the new signature (probe context instead of bool, one extra tuple member).
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Refs #4197.
Why
On a production SQL Server store (43 servers), nightly build 470:
MAX(run_datetime)MAX(instance_id)MAX(event_time)Every collector cycle re-reads the server-scoped watermark from the store even though, within one
runner's lifetime, nothing else can have written a newer value for that (server, collector) pair
between the read and this cycle's own write. The read is repeated work.
What changes
ServerWatermarkCache(new): a per-runner, per-(server, collector) cache of the last-knownwatermark (timestamp + optional UTC-frame flag + optional numeric twin).
TryGet,Seed,Advance,Invalidate.DarlingCollectorRunner.RunAsyncinvalidates the cache entry for a (server,collector) pair on ANY fault in the run (collection, dedup, write, or cancel) — a watermark ahead
of committed data would SKIP events, the one direction this cache must never take; falling behind
only re-collects duplicates, the safe direction.
ResolveServerWatermarkAsync(the cache hit/miss/seed/read block)and
AdvanceServerWatermark(the post-write cache-advance block), extracted out ofRunCoreAsyncas pure moves so a live-Postgres-only test can exercise the cache behaviour withouta live SQL Server target. The The server-scoped watermark read is discarded entirely for collectors that declare a per-database watermark #2797 dispatch gate (
ServerWatermarkIsDiscarded) is called INSIDEResolveServerWatermarkAsync(passed the dispatch probe), not precomputed inRunCoreAsyncandpassed in as a bool — an earlier extraction had split the gate call and the read it guards across
two method bodies, failing
ServerWatermarkDispatchGateTests' IL-walk invariant (every body thatperforms the server-scoped read must also call the gate).
RunCoreAsyncstill needs the discardedflag for
hasCollectedBefore, soResolveServerWatermarkAsync's return tuple grew one member(
ServerWatermarkDiscarded) rather than the gate being called twice.AdvanceServerWatermark(a real bug found and fixed here): PostgreSQL'stimestampcolumn truncates to microsecond resolution (10 .NET ticks) on write, but a batch'sin-memory
DateTime.Tickscan carry finer ticks than what the store actually persisted. Withoutflooring, a cached value sits sub-microsecond AHEAD of what a fresh read of the store would return
— the "cached vs. fresh" equivalence pin drifts by those ticks, and in principle the cache could
then be ahead of the committed watermark and skip an event on the next cycle's bound.
AdvanceServerWatermarknow floorsbatchMaxtoTimeSpan.TicksPerMicrosecondbefore caching it.ICollectorDefinition/CollectorDefinitionBaseaccessors:WatermarkValueAccessorandNumericWatermarkValueAccessor, letting the cache pull the batch's watermark value directly from awritten row rather than a second store read.
JobHistoryCollector(timestamp + numeric twin),DefaultTraceEventsCollector,SystemHealthEventsCollector,MemoryPressureEventsCollector(timestamp only).
QueryStoreCollector/CpuUtilizationCollector/PgCpuUtilizationCollectorareexcluded — see below.
Behaviour changes
MAXwould read as afirst run (the cache remembers what was actually collected; the store's own history of it may have
aged out).
pg_cpu_utilizationis excluded from the cache: a second writer,RdsCpuIngestor, can advance thestore's watermark for it outside this runner's own write path, and the cache has no way to observe
that write.
cpu_utilizationis excluded: it declares a UTC-twin frame (sample_time_utcbesidesample_time), and caching across that frame distinction is a follow-up, not part of this change.(server, collector) pair, so the next cycle re-seeds from the store rather than risk advancing past
rows that never committed.
Test plan
Run in-process on this Mac (
Darling.Tests.dll,Microsoft.WindowsDesktop.Appstripped from theruntimeconfig, against a local
timescale/timescaledb:2.30.1-pg18container withpg_stat_statements+timescaledbpreloaded and thedarlingrole/db present):ServerWatermarkDispatchGateTestsServerWatermarkCacheRunnerLiveTests#1776). The read counts come from the data source's own Npgsql command log (a counting logger attached withUseLoggerFactory), so the tests don't needpg_stat_statementspreloaded; CI's PostgreSQL doesn't preload it.ServerWatermarkCacheTestsDarlingCollectorRunnerTestsTimeHonestyRungTestsStorageCommandTimeoutTestsDocCommentHygieneTestsRED on dev:
ServerWatermarkCacheRunnerLiveTests.csandServerWatermarkCacheTests.cscopiedinto a detached worktree of
origin/dev(bb04415cd) fail to COMPILE there —ServerWatermarkCachedoes not exist, and neither do
ResolveServerWatermarkAsync/AdvanceServerWatermark/theWatermarkValueAccessorproperties the pins reference (multipleCS0246/CS1061).The fault-invalidation live pin (
Fault_InvalidatesTheCachedEntry_SoTheNextResolveReadsExactlyOnce)drives the fault through
ServerWatermarkCache.Invalidatedirectly, standing in for the runner's owntry/catch — the real
RunAsync/RunCoreAsyncfault path needs a live SQL Server target thislive-Postgres-only rig cannot provide.
Darling.Tests/Lite.Testsboth build 0 warnings/0 errors (net10.0-windows, cannot execute onmacOS; CI decides the Windows suite).
Measurement
Before/after, measured with a temporary in-process xUnit test (deleted before this commit) against
a local Postgres/TimescaleDB rig: 43 servers seeded with 250,000 job_history rows each (packed
inside
WatermarkPolicy.RecentWatermarkWindow, 6h) so one server's boundedMAX(run_datetime)touches 5,482 shared blocks (
EXPLAIN (ANALYZE, BUFFERS), above the 5,000-block target; field wasabout 14.6k). One simulated field hour = 43 servers × 12 runs, one small batch written per server
per run.
pg_stat_statements_reset()before each side.MAX(timestamp)MAX(timestamp)AFTER's call count matches the target: 43 servers × 2 seed calls (timestamp + numeric twin) on the
FIRST resolve of a fresh runner, then zero for every one of the 11 subsequent runs in the hour — the
cache absorbs the rest. BEFORE's call count is one direct read per server per run (43 × 12 = 516),
plus job_history's numeric-twin read doubling its own row (1,032).
CHANGELOG entry
SECTION: Changed
ENTRY:
MAX()reads againstjob_history,default_trace_events,system_health_events, andmemory_pressure_eventsto one seed per (server, collector) pair instead of one per collection cycle ([Cache the server-scoped watermark read across a runner's lifetime (#4197) #4399])REF:
[Cache the server-scoped watermark read across a runner's lifetime (#4197) #4399]: Cache the server-scoped watermark read across a runner's lifetime (#4197) #4399