Skip to content

Windows close on the engine's signal; lateness is decided at arrival - #399

Merged
turbolytics merged 8 commits into
mainfrom
feat/watermark-driven-close
Sep 27, 2026
Merged

turbolytics merged 8 commits into
mainfrom
feat/watermark-driven-close

Conversation

@turbolytics

@turbolytics turbolytics commented Sep 27, 2026 •

Copy link
Copy Markdown
Owner

Implements docs/superpowers/specs/2026-09-26-watermark-driven-close-design.md (spec and plan in #396; this branch carries those commits too, so its diff shrinks once #396 lands).

What changes

The manager no longer polls. The engine kicks the window's signal after the commit that moved its watermark or admitted a late row. The manager runs one pass on start, one per kick, one on the drain. No clock, no ticker, no interval. poll_interval_seconds is gone and validate names the change.

Lateness is decided at arrival, by the engine. Each placed record is classified against the window's asserted watermark, before the handler: refused past end + allowed_lateness_seconds and counted in window_late_rows_total{outcome="refused"}; written and queued for a recompute within the lateness; on time otherwise. The window table never holds a row the engine did not admit.

late_rows is replaced by allowed_lateness_seconds (default 0, Flink's allowedLateness). With a positive value a closed bucket's rows are kept that long and a late row republishes the bucket whole (emit_sql over every row it has), so the sink sees the exact value, never a delta. That needs a sink that replaces by key; validate and run refuse an appending one. window_recomputes_total counts republications. The start pass republishes every retained bucket, because the recompute set lived in memory.

time_column must be time_bucket(INTERVAL '<size>', event_time). validate refuses anything else, since the engine decides lateness from that bucket. core.BucketStart uses DuckDB's origin (2000-01-03), proven against time_bucket over a grid.

Bucket table: open / due / retained / expired, rendered to docs/windows/decisions.md.

Wire: late_rows_dropped now counts refusals, late_rows_recomputed is new, late_rows_reemitted is gone.

Configs: six examples, three bench configs and the soak config drop the removed keys, add event_time to their source, and bucket on event_time. The bluesky postgres example is the one that sets allowed_lateness_seconds: 300.

Not in this PR

  • render/pipeline.yml is unchanged. The render workflow validates it against the image it pins (v2026.09.21), whose schema still requires late_rows, so the template moves in the PR that bumps the pin after the release. That PR also decides the template's bucket clock: its payload timestamp is optional and per metric, so the candidate is the request's arrival event_time, with the metric's own timestamp kept in last_at.

Verification

  • go test -short -skip TestIntegration ./...: 33 packages ok (port-binding packages run outside the sandbox).
  • go test -race on core, managers, simulate, cli/run, conformance: ok.
  • Tooling tests 205 passed; make coverage-page regenerated; schema and CLI goldens regenerated.
  • New: model with lateness shapes and the exact-value property (11,110 sequences × 8 shapes); five simulator scenarios including a driven run where the kick, not a Pass step, publishes.

Merge order

#398 (serve flake) → #396 (spec + plan) → this.

BucketStart aligns to time_bucket's origin, 2000-01-03, rather than Go's
zero time. The two agree for any size that divides a day and for multiples
of a week, which is why the difference went unnoticed, and disagree by hours
for a 25-hour or 3-day bucket. The test runs the real time_bucket over
fourteen sizes and seven instants, the boundaries included.
…n start, kick and drain

The engine classifies each placed record against the window's asserted
watermark: refused past end + allowed_lateness_seconds, written and queued
for a recompute within it, on time otherwise. It kicks the window's signal
after the commit that moved the watermark or admitted a late row. The
manager has no clock and no interval: one pass on start, one per kick, one
on the drain. A pass publishes what is due, republishes whole every bucket
a late row landed in, and deletes what is past its lateness; the start pass
republishes every retained bucket, since the recompute set lived in memory.
late_rows, poll_interval_seconds and the manager's late sweep are gone;
allowed_lateness_seconds is the one knob. The bucket table becomes
open/due/retained/expired, and the conformance harness drives Pass.
…teness needs a replacing sink, every shipped window updated

validate refuses a windowing handler whose time_column is not
time_bucket(INTERVAL '<size>', event_time), a structured batch table without
event_time TIMESTAMPTZ, and lateness above zero paired with a sink that
appends; it names the replacement for the removed late_rows and
poll_interval_seconds keys. run refuses the lateness pairing too. The six
examples, the Render template, the three bench configs and the soak config
drop the removed keys, tell their source where the record's time is, and
bucket on event_time. The Render template cuts minutes on the request's
arrival, since its payload timestamp is optional and per metric.
…th the delta

window_late_rows carries outcome=refused and outcome=recomputed, recorded
by the engine. The bundle's late_rows_dropped counts refusals,
late_rows_recomputed is new, late_rows_reemitted is gone.
… exact value

Pass replaces Poll, StartPass is what a restarted manager does first, and
RunWindowedDriven runs the manager's own loop so the kick can be seen to
work. The sink keeps the last value per bucket; LateRefused, Unplaceable
and Recomputes come from the engine's and the manager's counters. Five
scenarios cover refusal, recompute, the lateness bound, a burst, and a
restart between the late row and the pass.
…elog and the spec say so

manager.publish.eventually is reworded for the kick, the two manager.late
claims go with the sweep, and pipeline.window.lateness_decided_at_arrival
is declared and tracked until the harness has a windowed subject.
… the new claim

The render workflow validates render/pipeline.yml against the image it
pins, whose schema still requires late_rows, so the template's move to
allowed_lateness_seconds and event_time goes with the pin bump after the
release. The two pipeline status files gain the row for
pipeline.window.lateness_decided_at_arrival, which the coverage check
demanded.
@turbolytics
turbolytics force-pushed the feat/watermark-driven-close branch from 5f7c697 to 6f5d18c Compare September 27, 2026 12:56
@turbolytics
turbolytics merged commit fa3780d into main Sep 27, 2026
6 checks passed
@github-actions github-actions Bot locked and limited conversation to collaborators Sep 27, 2026
Sign up for free to subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant