Skip to content

fix(store): merge late ForwardToStore corrections in exact-agg reads - #813

Open
GordonYuanyc wants to merge 2 commits into
mainfrom
fix/late-correction-merge
Open

GordonYuanyc wants to merge 2 commits into
mainfrom
fix/late-correction-merge

Conversation

@GordonYuanyc

Copy link
Copy Markdown

Why

A late sample for an already closed pane is written as a correction. Exact-agg reads kept only the last state per window, so the correction silently replaced the published value.

What

Exact-agg reads (sum/count/min/max/rate/increase) merge all states stored for a window. Sketch reads already did this.

How

query_exact_agg_range merges into the existing state (merge_with) instead of overwriting it. A merge failure fails the read, so the query falls back.

Before this PR

10 s pane: write 1..5, idle-close (16 s), write 6..10 → sum = 40.

After this PR

Same steps → 55. Retries still count once.

Evidence

Remote Write against a release build: 40 → 55 for the steps above.

Verification

  • Unit tests:
    • exact_agg_range_merges_late_correction_for_same_window: two states for one window read as their merge.
    • late_sample_merges_into_closed_pane_and_retries_count_once: receiver → worker → store → read.
    • Both fail on main.
  • End-to-end tests: Remote Write check above.
  • Other checks: fmt, clippy, cargo test -p data_plane --lib (899 passed).

Architectural decisions

Merge on read; the store stays append-only, matching sketch reads.

Limitations and follow-up

  • Distributed profile only: re-running a backfill over the same window now adds instead of overwriting.
  • The persistence disk tier still keeps one record per window. Persistence is not available in the asapquery profile.

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 2 commits October 1, 2026 01:01
With LateDataPolicy::ForwardToStore, input for a pane that the idle rule
or absolute deadline already closed is emitted as a separate correction
state and appended beside the published state for the same window. The
sketch read path merges every frame for a window, but
query_exact_agg_range kept only the last state per window end, so each
correction replaced the published sum instead of adding to it (Remote
Write repro: 15 published, then 40 more late, read 40 instead of 55).

Merge in-memory exact states that share a window, as the precompute
design doc specifies for ForwardToStore. States that cannot merge return
an error and the bound read fails closed instead of returning a partial
answer.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Replace the two clock-driven Remote Write tests with one test that closes
the pane through the worker drain, then checks the late merge and
exactly-once retries. Share the bound-read setup with the existing
queued-population test. Surface exact merge failures as Decode errors
and note that the disk tier keeps one record per window.

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