The failure mode is eager eviction driven by an over-trusted watermark. A temporal (interval) join needs both
The failure mode is eager eviction driven by an over-trusted watermark. A temporal (interval) join needs both
members of every pair. The naive design advances a single global watermark, emits/evicts an event as soon as the
watermark exceeds t + W, and drops anything at or below the watermark.
Three independent bugs fall out of that:
L > 0, an event below the watermark can still arrive and
match. Evicting at wm > t + W loses those matches. The correct safe point is t + W + L < min(wm_A, wm_B),
because the opposite side must also have moved past the match window minus the lateness slack.Two consequences the problem statement is testing:
L (or a disk spill / backpressure channel — which is
exactly what "hard cap" must mean).Events: (side ∈ {A,B}, key, timestamp, event_id). Join output: all pairs with equal key and |t_A − t_B| ≤ W.
State per key:
| state | meaning |
|---|---|
wm[k][side] |
max timestamp observed for (k, side) — a lower bound on future timestamps |
buf[k][A], buf[k][B] |
unsealed events, kept sorted by timestamp |
seen[k] |
ids of events currently in buf (dedup) |
Invariants that make the result exact:
F[k] = min(wm[k][A], wm[k][B]).x may be evicted only when t_x + W + L < F[k]. At that instant every
possible partner on the other side has already arrived (watermark honest, lateness ≤ L), and is still
buffered, so emitting x against the opposite buffer is complete.(a,b) is emitted exactly once — when the first of its two members is evicted.
Candidates are popped in increasing timestamp order, so the partner is guaranteed present.seen[k] ⇒ skip. Id not in seen[k] but t < F[k] − L − W ⇒ the event is below
the seal threshold; under the lateness bound it is a redundant copy (or a contract violation) and is dropped
without affecting the final result.size ≥ cap, the event goes to the spill queue (disk) and is replayed
when capacity frees. Dropping unsealed events is only ever safe if you accept approximate results.Bounded memory guarantee: at any moment per key only events with timestamps in (F[k]−W−L, F[k]] are retained,
so cap ≥ max_k (#events in a W+L horizon) suffices for zero spilling. A single hot/stalled key can still
consume the cap; that key’s source must be back-pressured (or spilled), not the data dropped.
Self-contained, tested implementation. Save as bounded_temporal_join.py and run.
"""
Bounded-memory temporal (interval) join over out-of-order streams with
per-key watermarks, bounded lateness, duplicate suppression, hard state cap
with disk spill, and a retraction protocol for speculative outer joins.
"""
from bisect import insort
from collections import defaultdict, deque
NEG = float('-inf')
POS = float('inf')
class IntervalJoin:
"""
Inner interval join: (a,b) emitted iff key(a)==key(b) and |t_a - t_b| <= W.
An event x is sealed once t_x + W + L < min(wm[key]['A'], wm[key]['B']).
Only sealed events are evicted; every pair is emitted exactly once.
"""
def __init__(self, W, lateness=0, max_state=POS, spill_enabled=True):
self.W, self.L = W, lateness
self.cap, self.spill_enabled = max_state, spill_enabled
self.buf = defaultdict(lambda: {'A': [], 'B': []})
self.wm = {}
self.seen = defaultdict(set) # ids of currently buffered events
self.size = 0
self.spill = deque() # disk-backed pending queue
self.emitted = []
self.n_spilled = self.n_late_dropped = self.n_dup = 0
def _ensure(self, k):
if k not in self.wm:
self.wm[k] = {'A': NEG, 'B': NEG}
return self.buf[k]
def _frontier(self, k):
return min(self.wm[k]['A'], self.wm[k]['B'])
def _seal_threshold(self, k):
return self._frontier(k) - self.L - self.W
def _insert(self, side, k, t, eid, payload):
b = self._ensure(k)
if t > self.wm[k][side]:
self.wm[k][side] = t
insort(b[side], (t, eid, payload))
self.seen[k].add(eid)
self.size += 1
def _advance(self, k):
b = self.buf[k]
while True:
thr = self._seal_threshold(k)
side = None
for s in ('A', 'B'): # pop globally smallest t
if b[s] and b[s][0][0] < thr:
if side is None or b[s][0][0] < b[side][0][0]:
side = s
if side is None:
break
t, eid, _ = b[side].pop(0)
self.size -= 1
self.seen[k].discard(eid)
other = 'B' if side == 'A' else 'A'
for (yt, yeid, _) in b[other]:
if yt < t - self.W:
continue
if yt > t + self.W:
break
self.emitted.append((k, eid, yeid) if side == 'A'
else (k, yeid, eid))
def _drain_spill(self):
while self.spill and self.size < self.cap:
side, k, t, eid, payload = self.spill.popleft()
if eid in self.seen[k]:
self.n_dup += 1
elif t < self._seal_threshold(k):
self.n_late_dropped += 1
else:
self._insert(side, k, t, eid, payload)
self._advance(k)
def push(self, side, k, t, eid, payload=None):
self._ensure(k)
if eid in self.seen[k]:
self.n_dup += 1
return 'dup'
if t < self._seal_threshold(k):
self.n_late_dropped += 1
return 'too-late'
if self.size >= self.cap:
if not self.spill_enabled:
self.n_late_dropped += 1
return 'dropped-cap'
self.spill.append((side, k, t, eid, payload))
self.n_spilled += 1
return 'spilled'
self._insert(side, k, t, eid, payload)
self._advance(k)
self._drain_spill()
return 'ok'
def close(self):
"""End of stream: drain disk ignoring cap, then flush all watermarks."""
while self.spill:
side, k, t, eid, payload = self.spill.popleft()
if eid not in self.seen[k]:
self._insert(side, k, t, eid, payload)
self._advance(k)
for k in list(self.wm):
self.wm[k]['A'] = self.wm[k]['B'] = POS
self._advance(k)
return sorted(self.emitted)
class RetractingLeftOuter:
"""
A-left outer interval join with low-latency speculative emission.
A provisional (a, None) row is emitted at the SOFT frontier t_a+W < wm_B,
retracted if a late B within L matches, and its retraction state is dropped
at the HARD frontier t_a+W+L < min(wm_A, wm_B).
"""
def __init__(self, W, L):
self.W, self.L = W, L
self.bufA, self.bufB = [], []
self.wmA = self.wmB = NEG
self.prov = {} # aid -> {'t','rows','final'}
self.seenA, self.seenB = set(), set()
self.log = [] # ('assert'|'retract', row)
self.n_retract = 0
def _matches(self, a):
return [(b, be) for (b, be) in self.bufB if abs(b - a[0]) <= self.W]
def _emit(self, kind, row):
self.log.append((kind, row))
if kind == 'retract':
self.n_retract += 1
def _soft_seal(self, a):
aid = a[1]
ms = self._matches(a)
rows = [(aid, be) for (_, be) in ms] or [(aid, None)]
for r in rows:
self._emit('assert', r)
self.prov[aid] = {'t': a[0], 'rows': rows}
def _refresh(self, aid):
a = next((x for x in self.bufA if x[1] == aid), (self.prov[aid]['t'], aid))
ms = self._matches(a)
newrows = [(aid, be) for (_, be) in ms] or [(aid, None)]
for r in self.prov[aid]['rows']:
self._emit('retract', r)
for r in newrows:
self._emit('assert', r)
self.prov[aid]['rows'] = newrows
def _flush_A(self):
F = min(self.wmA, self.wmB)
for a in list(self.bufA):
t, eid = a
hard = t < F - self.L - self.W
soft = t < self.wmB - self.W
if eid not in self.prov and (hard or soft):
self._soft_seal(a)
if hard:
self.prov.pop(eid, None) # rows final; drop retraction state
self.bufA.remove(a)
def push(self, side, t, eid):
if side == 'A':
if eid in self.seenA:
return
self.seenA.add(eid)
self.bufA.append((t, eid))
self.wmA = max(self.wmA, t)
self._flush_A()
else:
if eid in self.seenB:
return
self.seenB.add(eid)
self.bufB.append((t, eid))
self.wmB = max(self.wmB, t)
for aid, st in list(self.prov.items()):
if abs(st['t'] - t) <= self.W:
self._refresh(aid)
self._flush_A()
def close(self):
self.wmA = self.wmB = POS
self._flush_A()
for a in list(self.bufA):
if a[1] not in self.prov:
self._soft_seal(a)
final = set()
for kind, row in self.log:
final.add(row) if kind == 'assert' else final.discard(row)
return final
# ----------------------------- verification -----------------------------
def brute_force(events, W):
"""Exact batch answer, dedup by (side,key,id)."""
seen, A, B = set(), defaultdict(list), defaultdict(list)
for side, k, t, eid in events:
if (side, k, eid) in seen:
continue
seen.add((side, k, eid))
(A if side == 'A' else B)[k].append((t, eid))
return {(k, ea, eb)
for k in set(A) | set(B)
for ta, ea in A[k] for tb, eb in B[k]
if abs(ta - tb) <= W}
if __name__ == '__main__':
import random
def actual_lateness(events):
wm = defaultdict(lambda: {'A': NEG, 'B': NEG})
worst = 0.0
for s, k, t, _ in events:
f = min(wm[k]['A'], wm[k]['B'])
if f > NEG:
worst = max(worst, f - t)
wm[k][s] = max(wm[k][s], t)
return max(worst, 0.0)
# 1. exactness under arbitrary permutations
bad = 0
for seed in range(500):
rng = random.Random(seed)
W = rng.choice([0, 1, 3, 5])
evs = [(rng.choice('AB'), rng.randrange(3), rng.randrange(40), f'e{i}')
for i in range(rng.randrange(5, 60))]
order = list(range(len(evs))); rng.shuffle(order)
sh = [evs[i] for i in order]
j = IntervalJoin(W, int(actual_lateness(sh)) + 1)
for s, k, t, e in sh:
j.push(s, k, t, e)
bad += set(j.close()) != brute_force(sh, W)
assert bad == 0
print('inner exact: 500/500 random permutations match brute force')
# 2. duplicates + hard cap with spill
bad = 0
for seed in range(300):
rng = random.Random(seed)
W = rng.choice([1, 2, 4])
evs = [(rng.choice('AB'), rng.randrange(3), rng.randrange(20), f'e{i}')
for i in range(rng.randrange(5, 50))]
evs += [evs[i] for i in rng.sample(range(len(evs)), min(8, len(evs)))]
rng.shuffle(evs)
j = IntervalJoin(W, int(actual_lateness(evs)) + 1, max_state=6)
for s, k, t, e in evs:
j.push(s, k, t, e)
bad += set(j.close()) != brute_force(evs, W)
assert bad == 0
print('duplicates + hard cap/disk-spill: 300/300 match brute force')
# 3. adversarial skew: B stalls, A floods; drop-on-cap loses data
evs = [('A', 'k', t, f'A{t}') for t in range(50)] + [('B', 'k', 25, 'B25')]
random.Random(1).shuffle(evs)
L = int(actual_lateness(evs)) + 1
j = IntervalJoin(1, L, max_state=4)
for s, k, t, e in evs:
j.push(s, k, t, e)
exact = set(j.close())
j2 = IntervalJoin(1, L, max_state=4, spill_enabled=False)
for s, k, t, e in evs:
j2.push(s, k, t, e)
lossy = set(j2.close())
assert exact == brute_force(evs, 1) and lossy != exact
print(f'adversarial skew: exact={len(exact)} via {j.n_spilled} spilled; '
f'drop-on-cap loses {len(exact) - len(lossy)}')
# 4. retraction for speculative outer join
evs = [('A', 5, 'A5'), ('B', 10, 'B10'), ('B', 5, 'B5')]
r = RetractingLeftOuter(1, 3)
for s, t, e in evs:
r.push(s, t, e)
assert r.close() == {('A5', 'B5')} and r.n_retract >= 1
print('retraction log:', r.log)
print('ALL CHECKS PASSED')
Expected output:
inner exact: 500/500 random permutations match brute force
duplicates + hard cap/disk-spill: 300/300 match brute force
adversarial skew: exact=3 via 47 spilled; drop-on-cap loses 3
retraction log: [('assert', ('A5', None)), ('retract', ('A5', None)), ('assert', ('A5', 'B5'))]
ALL CHECKS PASSED
Watermark skew. With F[k] = min(wm_A, wm_B), a stalling side correctly freezes sealing: no matches can be
retired while the other side is behind. But that same rule means a single hot/stalled key accumulates state up to
cap. Per-key watermarks fix cross-key head-of-line blocking (a stall in key k1 no longer pins k2); they do
not fix within-key growth. The cap must therefore back-pressure that key’s source or spill. The skew test shows
47 events spilled while the result stayed exact; the spill_enabled=False variant silently lost all 3 pairs.
Duplicates. seen[k] holds only buffered ids. Once an event seals, any later copy has t < F − L − W, hits
the too-late branch, and is a no-op. Dedup memory is O(buffered), not O(all history). Caveat: if copies can
arrive arbitrarily late (no lateness bound), you need either an unbounded id set or an idempotent sink keyed by
pair id — the latter keeps memory bounded and is the recommended production choice.
Impossibility. Take a key, an A event at t=0, and an adversary that delivers the matching B event at an
arbitrary future time. To be exact you must keep A forever, so any finite cap without spill/backpressure must
either lose the match or drop a later event. Hence: exactness ⇔ bounded lateness L (or unbounded/disk state).
This is why the fix replaces "drop at cap" with "spill at cap".
Correctness argument for the seal rule. When t_x + W + L < F[k], the honest watermark and lateness bound
guarantee no future opposite event can fall in [t_x − W, t_x + W]; every such partner has arrived and, because
it is not yet evictable by the same bound, is still in the opposite buffer. Emitting then is complete and final.
Popping candidates in increasing t ensures each pair is emitted exactly once: at the first member’s eviction the
other member is still buffered, and after eviction it cannot be emitted again.
| check | result |
|---|---|
| 500 random permutations vs. batch brute force (inner join) | 0 mismatches |
300 cases with injected duplicates and cap=6 (spill on) |
0 mismatches |
| adversarial skew, stalled side + flooded side | exact via 47 spilled; drop-on-cap loses 3 pairs |
| forced speculative null→pair retraction (outer join) | correct final {('A5','B5')}, 1 retraction |
| duplicate handling | seen bounded to buffered ids only |
Run: python3 bounded_temporal_join.py → ALL CHECKS PASSED.
# Evidence - Problem class: temporal-join-retraction-frontier - Model: openrouter/deepseek/deepseek-v4.1-flash - Solved: 2026-10-06T04:02:07.933Z - Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Design a bounded-memory temporal join over two out-of-order event streams with per-key watermarks, late-data retractions, and a hard cap on retained state. Specify the minimal state and eviction rules that preserve exact results under arbitrary arrival order, and analyze adversarial watermark skew and duplicate events.", "environment": "python3", "language": "python", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "temporal-join-retraction-frontier", "provider": "openrouter", "solved_at": "2026-10-06T04:02:07.936Z", "version": "3.11"}