Skip to content

Cache the server-scoped watermark read across a runner's lifetime (#4197) - #4399

Merged
erikdarlingdata merged 10 commits into
devfrom
fix/4197-cached-watermark
Sep 26, 2026
Merged

erikdarlingdata merged 10 commits into
devfrom
fix/4197-cached-watermark

Conversation

@erikdarlingdata

@erikdarlingdata erikdarlingdata commented Sep 26, 2026 •

Copy link
Copy Markdown
Owner

Refs #4197.

Why

On a production SQL Server store (43 servers), nightly build 470:

Statement Calls/hour Mean Bytes/hour
job_history MAX(run_datetime) 506 ~63 ms ~56 GB
job_history MAX(instance_id) 506 ~63 ms ~56 GB
default_trace_events MAX(event_time) 503 ~35 ms ~10.1 GB
unbounded fallback 0 — —

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-known
    watermark (timestamp + optional UTC-frame flag + optional numeric twin). TryGet, Seed,
    Advance, Invalidate.
  • Runner wiring: DarlingCollectorRunner.RunAsync invalidates 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.
  • The resolve/advance seam: ResolveServerWatermarkAsync (the cache hit/miss/seed/read block)
    and AdvanceServerWatermark (the post-write cache-advance block), extracted out of
    RunCoreAsync as pure moves so a live-Postgres-only test can exercise the cache behaviour without
    a 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 INSIDE
    ResolveServerWatermarkAsync (passed the dispatch probe), not precomputed in RunCoreAsync and
    passed 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 that
    performs the server-scoped read must also call the gate). RunCoreAsync still needs the discarded
    flag for hasCollectedBefore, so ResolveServerWatermarkAsync's return tuple grew one member
    (ServerWatermarkDiscarded) rather than the gate being called twice.
  • Microsecond flooring in AdvanceServerWatermark (a real bug found and fixed here): PostgreSQL's
    timestamp column truncates to microsecond resolution (10 .NET ticks) on write, but a batch's
    in-memory DateTime.Ticks can carry finer ticks than what the store actually persisted. Without
    flooring, 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.
    AdvanceServerWatermark now floors batchMax to TimeSpan.TicksPerMicrosecond before caching it.
  • ICollectorDefinition/CollectorDefinitionBase accessors: WatermarkValueAccessor and
    NumericWatermarkValueAccessor, letting the cache pull the batch's watermark value directly from a
    written row rather than a second store read.
  • Four collectors wired to the new accessors: JobHistoryCollector (timestamp + numeric twin),
    DefaultTraceEventsCollector, SystemHealthEventsCollector, MemoryPressureEventsCollector
    (timestamp only). QueryStoreCollector/CpuUtilizationCollector/PgCpuUtilizationCollector are
    excluded — see below.

Behaviour changes

  • After a retention purge, the cache keeps the real watermark where a fresh MAX would read as a
    first run (the cache remembers what was actually collected; the store's own history of it may have
    aged out).
  • pg_cpu_utilization is excluded from the cache: a second writer, RdsCpuIngestor, can advance the
    store's watermark for it outside this runner's own write path, and the cache has no way to observe
    that write.
  • cpu_utilization is excluded: it declares a UTC-twin frame (sample_time_utc beside
    sample_time), and caching across that frame distinction is a follow-up, not part of this change.
  • A fault (collection, dedup, write, or cancel) unconditionally drops the cached entry for that
    (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.App stripped from the
runtimeconfig, against a local timescale/timescaledb:2.30.1-pg18 container with
pg_stat_statements + timescaledb preloaded and the darling role/db present):

Class Result
ServerWatermarkDispatchGateTests 7/7 pass — the pin for the gate-placement fix above; unchanged against dev after the fix (gate call moved into the seam, not the pin edited)
ServerWatermarkCacheRunnerLiveTests 5/5 pass (live, own scratch DB per test, #1776). The read counts come from the data source's own Npgsql command log (a counting logger attached with UseLoggerFactory), so the tests don't need pg_stat_statements preloaded; CI's PostgreSQL doesn't preload it.
ServerWatermarkCacheTests 11/11 pass
DarlingCollectorRunnerTests 9/9 pass, 3 skipped (need a live SQL Server target this rig cannot provide)
TimeHonestyRungTests 8/8 pass
StorageCommandTimeoutTests 19/19 pass
DocCommentHygieneTests 77/77 pass

RED on dev: ServerWatermarkCacheRunnerLiveTests.cs and ServerWatermarkCacheTests.cs copied
into a detached worktree of origin/dev (bb04415cd) fail to COMPILE there — ServerWatermarkCache
does not exist, and neither do ResolveServerWatermarkAsync/AdvanceServerWatermark/the
WatermarkValueAccessor properties the pins reference (multiple CS0246/CS1061).

The fault-invalidation live pin (Fault_InvalidatesTheCachedEntry_SoTheNextResolveReadsExactlyOnce)
drives the fault through ServerWatermarkCache.Invalidate directly, standing in for the runner's own
try/catch — the real RunAsync/RunCoreAsync fault path needs a live SQL Server target this
live-Postgres-only rig cannot provide.

Darling.Tests/Lite.Tests both build 0 warnings/0 errors (net10.0-windows, cannot execute on
macOS; 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 bounded MAX(run_datetime)
touches 5,482 shared blocks (EXPLAIN (ANALYZE, BUFFERS), above the 5,000-block target; field was
about 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.

Statement BEFORE (dev shape: direct reads) AFTER (this branch: cache seam)
job_history MAX (timestamp) 1,032 calls, 15.97 ms mean, 43,266 MB/hour 86 calls, 12.92 ms mean, 3,611 MB/hour
default_trace_events MAX (timestamp) 516 calls, 0.015 ms mean, 8.06 MB/hour 43 calls, 0.013 ms mean, 0.67 MB/hour

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:

…#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).
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant