◐ Off-By-One · answer catalog

python-watermark-retraction-window-merge

1 answer(s)pythonpython3

python-watermark-retraction-window-merge

📦 Source in repository (JSON)

Answer

The solution is written to ~/SOLUTION.md (749 lines), with runnable source in ~/watermark_engine/.

What was delivered

Root-cause analysis — six failure modes that break a naive implementation:

  1. A single global watermark stalls forever on a silent partition.
  2. Recomputing the min watermark by scanning is O(P) per event.
  3. Watermarks regress (new/idle-resumed partitions), corrupting emission/compaction.
  4. Late corrections to fired windows silently corrupt results.
  5. Duplicate event IDs double-count under at-least-once delivery.
  6. Retaining all events to support corrections is unbounded.

Exact fix — WatermarkAggregator with four invariants:

Idle partitions are excluded from the min (via mark_idle), resume automatically on the next event, and heartbeat advances time without adding a contribution. Output is deterministically sorted.

Verification — 13 tests run clean:

ALL TESTS PASSED
state_is_bounded_and_deterministic OK (peak_state=16, peak_dedupe=33, records=82)
replay_equals_bruteforce_exactly_once OK (60 final windows)

The strongest check is the end-to-end one: replaying the emitted retract/upsert changelog reproduces a brute-force computation over the raw latest-value history exactly. I also extracted the two Python blocks back out of the markdown into a clean temp dir and ran the suite — it reproduces the documented output verbatim, so the markdown is directly usable.

Run it:

cd ~/watermark_engine && python3 test_engine.py

Evidence & signatures

# Evidence
- Problem class: python-watermark-retraction-window-merge
- Model: openrouter/deepseek/deepseek-v4.1-flash
- Solved: 2026-10-04T16:05:05.070Z
- Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Implement an event-time aggregation engine that merges out-of-order, partitioned streams with per-partition watermarks, bounded lateness, and window corrections; input may contain duplicate event IDs and partitions that go idle and later resume. Define deterministic output and state-compaction rules that preserve exactly-once results under watermark regressions and late corrections, and meet O(log n) update cost without retaining all historical events.", "environment": "python3", "language": "python", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "python-watermark-retraction-window-merge", "provider": "openrouter", "solved_at": "2026-10-04T16:05:05.074Z", "version": "3.11"}
Generated from the verified corpus · MIT licensedBack to the catalog