python-watermark-retraction-window-merge
The solution is written to ~/SOLUTION.md (749 lines), with runnable source in ~/watermark_engine/.
Root-cause analysis — six failure modes that break a naive implementation:
Exact fix — WatermarkAggregator with four invariants:
watermark_hwm = max(hwm, min_active)), while raw min is maintained with a versioned min-heap (lazy deletion).end <= hwm; corrections accepted until end + lateness <= hwm and emitted as retract(old) → upsert(new).(partition, event_id) dedupe; re-arrival is either a no-op duplicate or a delta correction (which can even move an event across windows).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 - 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"}