Skip to content

Add a live event-time barrier for Remote Write (POST /api/v1/precompute/watermark) - #814

Open
GordonYuanyc wants to merge 6 commits into
mainfrom
fix/remote-write-watermark-barrier
Open

GordonYuanyc wants to merge 6 commits into
mainfrom
fix/remote-write-watermark-barrier

Conversation

@GordonYuanyc

Copy link
Copy Markdown

Addresses #772.

Why

A pane only closes when a later sample arrives in the same series, or after window + 5 s of wall-clock idle. When a series stops sending data, its last pane stays open, and every query over that window falls back for all series. Remote Write has no way for the producer to say "everything up to T is written".

What

POST /api/v1/precompute/watermark?time_ms=T closes, in every series, the panes whose samples are all at or before T. Unlike drain, input stays open.

How

  • The router sends an AdvanceWatermark message to every worker, the same way drain does. Each message queues behind the writes that were already acknowledged.
  • Worker::close_through closes buckets up to T (T+1 for left-closed membership). It goes no further than the latest open bucket.
  • flush_all and force_close_all are unchanged.

Before this PR

A task stops after epoch 0. With 150 s panes, the next window falls back for about 155 s.

After this PR

After the barrier call returns (5–12 ms), the next query is served warm.

Evidence

Demo workload: 30 tasks starting and stopping over 10 epochs, with one barrier per epoch. Warm epochs went from 0/10 to 10/10, and p50 error stayed ≤ 0.9%.

Verification

  • Unit tests:
    • watermark_barrier_closes_complete_pane_of_stopped_group: exact boundary for both membership modes; a repeated or lower barrier does nothing.
    • watermark_barrier_closes_full_windows_and_bounds_its_scan: FullWindow layout; i64::MAX terminates; later input aggregates normally.
  • End-to-end tests:
    • watermark_barrier_serves_window_of_stopped_series: real process; not warm before the barrier, warm right after it, and input still accepted.
  • Other checks: fmt, clippy, cargo test -p data_plane --lib (899 passed).

Architectural decisions

The barrier is a single producer-side statement for the whole backend. There is no producer or partition roster.

Limitations and follow-up

  • With several producers, send the barrier only after all of them have written through T. The user guide says so.
  • In --remote-write-revisions mode the barrier is a no-op that returns success.

Human review — do not complete with an agent

  • The MVP boundary is correct.
  • New conceptual layers or public interfaces are necessary.
  • The before/after description matches the intended product behavior.
  • Human reviewer:
  • Decision and rationale:

🤖 Generated with Claude Code

GordonYuanyc and others added 6 commits October 1, 2026 00:52
A pane is published only after a strictly later sample arrives in the same
series, or after the wall-clock idle rule fires (window + 5 s). A series
that stops sending leaves its last pane open, so every query over that
window misses ("materialization population has unpublished input") and is
forwarded. Remote Write has no way to say "all data up to T is written",
and /api/v1/precompute/drain seals input for the process lifetime.

POST /api/v1/precompute/watermark?time_ms=T lets the producer declare that
every sample at or before T has been written. The receiver broadcasts a
WorkerMessage::AdvanceWatermark to every worker; Remote Write acknowledges
only after enqueueing, so the barrier queues behind every acknowledged
write. Each worker closes, in every group, the stored buckets (panes or
FullWindow windows) whose samples all lie at or before T, including groups
that stopped, and replies after publishing. Right-closed PromQL groups file
a sample at t as t - 1, so their barrier is T; other groups use T + 1. The
scan is bounded by the latest open bucket, as in force_close_all. Input
stays open; later samples for closed buckets follow the late-data policy.
A failed worker answers the barrier with its error.

flush_all and force_close_all are unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Document POST /api/v1/precompute/watermark for users of the asapquery
profile, and correct the design note that said the profile has no public
barrier endpoint.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
close_through capped its scan at the end of a group's latest open bucket
but still set the closure watermark to the full barrier. A barrier far in
the future (for example microseconds or nanoseconds sent as time_ms) then
made every later sample of every existing group late until restart: under
ForwardToStore each became a standalone one-sample correction, and counter
deltas were dropped.

The closure watermark now advances only to min(barrier, latest open bucket
end), and not at all for a group without open buckets. Published buckets
are still detected as late.

The worker barrier tests are merged into one left/right-closed boundary
test, and the FullWindow test now checks that input after an i64::MAX
barrier aggregates normally (it was dropped before this change). The
async run-loop tests are removed; the process test covers that path.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Drop the receiver test (a one-line forward plus channel FIFO; the process
test shows input stays open). In the process test, query once right after
the barrier returns instead of polling, since panes are published before
the 200, and drop the missing-time_ms check.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A window also closes at the latest window + 5 s after its first sample,
not only after being idle that long.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@GordonYuanyc
GordonYuanyc requested a review from zzylol October 1, 2026 06:31
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