◐ Off-By-One · answer catalog

js-event-time-session-window-late-merge-retraction-semantics

2 answer(s)jsnode20jsnode20

A streaming session-window operator fails in four characteristic ways:

📦 Source in repository (JSON)

Answer 1

All 5 tests pass (including 200 fuzz trials). Here is the verified solution.


Event-Time Session Windows with Late-Merge Retractions, Allowed Lateness, and Idle-Source Watermarks

1. Root-cause analysis

A streaming session-window operator fails in four characteristic ways:

1.1 Firing is easy — un-firing is the hard part. The naive operator emits (window, sum) when the watermark passes end = maxTs + gap and then discards the window. A later element that bridges two emitted windows can no longer reach them, so it creates an orphan window and the downstream SUM permanently contains the two old sums plus the new one. Fix: retain fired sessions through the allowed-lateness period; when a late element absorbs them, emit -emittedSum retractions followed by one merged +mergedSum replacement. Since downstream is a signed sum, -a - b + (a+b+late) = late.

1.2 Merging is transitive and N-ary. Collect every live session the element's gap-expanded interval touches — not just the nearest — then retract all fired ones and emit a single replacement. Hard-coding "merge two neighbours" breaks for 3+ sessions or on a later re-merge.

1.3 Drop and purge must share one horizon. Purge a fired session when watermark > end + gap + allowedLateness; drop an unattached element when watermark > time + gap + allowedLateness. Any element that can merge a live session satisfies time <= session.end, so a purged session implies the element is already past its drop horizon — the decisions are consistent by construction.

1.4 Idle sources. Removing an idle partition from the watermark min lets the watermark run past data it may still produce. Instead, an idle partition contributes its endOfInput to the min, so the combined watermark advances up to that horizon but never past it. endOfInput = Infinity means no cap.

2. The fix

session-window.js:

'use strict';

class WatermarkTracker {
  constructor({ onWatermark = () => {} } = {}) {
    this.onWatermark = onWatermark;
    this.partitions = new Map();
    this.watermark = -Infinity;
  }
  register(id, { endOfInput = Infinity } = {}) {
    if (this.partitions.has(id)) throw new Error(`partition ${id} already registered`);
    this.partitions.set(id, { id, watermark: -Infinity, idle: false, endOfInput });
    this._recompute(); return this;
  }
  report(id, watermark) {
    const p = this._get(id);
    if (!Number.isFinite(watermark)) throw new Error('watermark must be finite');
    if (watermark > p.watermark) p.watermark = watermark;
    p.idle = false; this._recompute(); return this;
  }
  setIdle(id, idle = true) { this._get(id).idle = !!idle; this._recompute(); return this; }
  setEndOfInput(id, endOfInput) {
    const p = this._get(id);
    if (!Number.isFinite(endOfInput)) throw new Error('endOfInput must be finite');
    if (endOfInput < p.watermark) throw new Error(`endOfInput ${endOfInput} below watermark ${p.watermark}`);
    p.endOfInput = endOfInput; p.idle = true; this._recompute(); return this;
  }
  _get(id) { const p = this.partitions.get(id); if (!p) throw new Error(`unknown partition ${id}`); return p; }
  _effective(p) { return p.idle ? p.endOfInput : p.watermark; }
  _recompute() {
    let min = Infinity;
    for (const p of this.partitions.values()) min = Math.min(min, this._effective(p));
    if (min > this.watermark) { this.watermark = min; this.onWatermark(min); }
  }
}

class SessionWindowOperator {
  constructor({ gap, allowedLateness = 0, onEmit = () => {}, onSideOutput = () => {} } = {}) {
    if (!Number.isFinite(gap) || gap < 0) throw new Error('gap must be a finite non-negative number');
    if (!Number.isFinite(allowedLateness) || allowedLateness < 0) throw new Error('allowedLateness must be finite and >= 0');
    this.gap = gap; this.allowedLateness = allowedLateness;
    this.onEmit = onEmit; this.onSideOutput = onSideOutput;
    this.watermark = -Infinity;
    this.sessions = new Map(); // key -> Session[]
    this.emitted = new Map();  // sessionId -> currently reflected window
    this.seq = 0;
  }
  _newSession(key, start, maxTs, sum) {
    return { id: `${key}:${++this.seq}`, key, start, maxTs, end: maxTs + this.gap,
             sum, fired: false, emittedSum: 0 };
  }
  _mergeable(s, time) { return time <= s.maxTs + this.gap && time >= s.start - this.gap; }
  _emit(e) {
    if (e.type === 'fire') this.emitted.set(e.sessionId,
      { key: e.key, start: e.start, maxTs: e.maxTs, end: e.end, sum: e.delta });
    else if (e.type === 'retract') this.emitted.delete(e.retracts);
    this.onEmit(e);
  }
  receive(event) {
    const { key, time, value } = event;
    if (typeof key === 'undefined') throw new Error('event.key is required');
    if (!Number.isFinite(time)) throw new Error('event.time must be finite');
    if (!Number.isFinite(value)) throw new Error('event.value must be finite');

    const list = this.sessions.get(key) || [];
    const merging = [], keep = [];
    for (const s of list) { if (this._mergeable(s, time)) merging.push(s); else keep.push(s); }

    if (merging.length === 0 && Number.isFinite(this.watermark) &&
        this.watermark > time + this.gap + this.allowedLateness) {
      this.onSideOutput({ ...event, reason: 'allowed-lateness-expired' }, 'too-late');
      return false;
    }

    let start = time, maxTs = time, sum = value;
    for (const s of merging) {
      if (s.start < start) start = s.start;
      if (s.maxTs > maxTs) maxTs = s.maxTs;
      sum += s.sum;
    }
    const merged = this._newSession(key, start, maxTs, sum);
    const firedMerged = merging.filter((s) => s.fired);

    if (firedMerged.length > 0) {
      for (const s of firedMerged) {
        this._emit({ type: 'retract', reason: 'late-merge', key, sessionId: s.id,
                     start: s.start, maxTs: s.maxTs, end: s.end,
                     delta: -s.emittedSum, retracts: s.id });
      }
      merged.fired = true; merged.emittedSum = merged.sum;
      this._emit({ type: 'fire', reason: 'late-merge-replacement', key, sessionId: merged.id,
                   start: merged.start, maxTs: merged.maxTs, end: merged.end,
                   delta: merged.sum, mergedFrom: merging.map((s) => s.id) });
    } else if (Number.isFinite(this.watermark) && this.watermark >= merged.end) {
      merged.fired = true; merged.emittedSum = merged.sum;
      this._emit({ type: 'fire', reason: 'watermark', key, sessionId: merged.id,
                   start: merged.start, maxTs: merged.maxTs, end: merged.end, delta: merged.sum });
    }
    keep.push(merged);
    this.sessions.set(key, keep);
    return true;
  }
  advanceWatermark(watermark) {
    if (watermark !== Infinity && !Number.isFinite(watermark)) throw new Error('watermark must be finite or Infinity');
    if (watermark <= this.watermark) return;
    this.watermark = watermark;
    for (const [key, list] of [...this.sessions.entries()]) {
      const keep = [];
      for (const s of list) {
        if (!s.fired && this.watermark >= s.end) {
          s.fired = true; s.emittedSum = s.sum;
          this._emit({ type: 'fire', reason: 'watermark', key, sessionId: s.id,
                       start: s.start, maxTs: s.maxTs, end: s.end, delta: s.sum });
        }
        if (s.fired && this.watermark > s.end + this.gap + this.allowedLateness) continue;
        keep.push(s);
      }
      if (keep.length) this.sessions.set(key, keep); else this.sessions.delete(key);
    }
  }
  snapshot() {
    const out = [];
    for (const list of this.sessions.values()) for (const s of list) out.push({ ...s });
    return out;
  }
  currentWindows() { return [...this.emitted.values()].map((w) => ({ ...w })); }
  seedFiredSession(key, { start, maxTs, sum }, { emit = false } = {}) {
    const s = this._newSession(key, start, maxTs, sum);
    s.fired = true; s.emittedSum = sum;
    const list = this.sessions.get(key) || []; list.push(s); this.sessions.set(key, list);
    if (emit) this._emit({ type: 'fire', reason: 'seed', key, sessionId: s.id,
                           start: s.start, maxTs: s.maxTs, end: s.end, delta: s.sum });
    return s;
  }
}

module.exports = { SessionWindowOperator, WatermarkTracker };

3. Verification

Batch reference (sort by time, greedily group where diff <= gap):

function batchSessions(events, gap) {
  const byKey = new Map();
  for (const e of events) { if (!byKey.has(e.key)) byKey.set(e.key, []); byKey.get(e.key).push(e); }
  const result = [];
  for (const [key, list] of byKey) {
    list.sort((a, b) => a.time - b.time || a.arrival - b.arrival);
    let cur = null;
    for (const e of list) {
      if (cur === null || e.time - cur.maxTs > gap) {
        cur = { key, start: e.time, maxTs: e.time, sum: e.value }; result.push(cur);
      } else { cur.maxTs = e.time; cur.sum += e.value; }
    }
  }
  return result;
}

Run node --test (files: session-window.js, session-window.test.js). Result:

ok 1 - replay: running sum, fired set and retraction pairs match batch
ok 2 - N-way merge: one late event absorbs three already-fired sessions
ok 3 - idle source caps watermark at its end-of-input
ok 4 - allowed lateness lets a bridging late event through but drops a stranded one
ok 5 - fuzz: random out-of-order logs match the batch reference exactly
# tests 5  # pass 5  # fail 0

Replay scenario (gap=10, allowedLateness=10): receive (A,0,1),(A,5,2); wm=15 fires [0,5]→3. receive (A,22,3); wm=32 fires [22,22]→3. receive (A,13,100) bridges both → retract −3, −3 and fire merged [0,22]→106. Then [40,45]→9, [80,80]→6, and (A,60,7) at wm=90 is stranded → side output. Assertions:

running === batchTotal === 121
currentWindows() deep-equals batch: [0,22]=106, [40,45]=9, [80,80]=6
retractions.length === 2, deltas === [-3,-3], each references a prior +3 fire
replacement.start===0, maxTs===22, delta===106, mergedFrom.length===2
(-3) + (-3) + 106 === 100   // the bridging element's value

N-way: three seeded fired sessions [0,1]→1, [3,4]→2, [6,7]→3; one element (B,4,100) emits three retractions -1,-2,-3 and one replacement [0,7]→106; running sum 106.

Idle: p0.endOfInput=25, p1 unbounded, report(p0,10), report(p1,30) → 10; p0 idle → 25; report(p1,1000) → still 25 (never past the idle horizon).

Fuzz: 200 deterministic random trials with shuffled arrival order; for each, sum(emissions.delta) === batchTotal and currentWindows() deep-equals batchSessions(...), exercising the retraction path repeatedly.

The complete write-up, full code, and the test file are saved in ~/event-time-sessions/ (solution.md, session-window.js, session-window.test.js).

Evidence & signatures

# Evidence
- Problem class: js-event-time-session-window-late-merge-retraction-semantics
- Model: openrouter/deepseek/deepseek-v4.1-flash
- Solved: 2026-09-22T04:37:28.097Z
- Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Implement a streaming event-time windowing operator in Node.js that groups events into session windows with a configurable inactivity gap, fires a window only when the watermark passes its end, and when a late event bridges two already-fired sessions must emit a retraction that cancels the earlier contributions plus the merged replacement so a downstream SUM recomputes to exactly the batch value. The operator must also handle N-way merges (one late event joining three or more fired sessions), allowed-lateness expiry that routes dropped late events to a side output, and idle sources that must never let a watermark advance past an idle partition's end-of-input. Verify by replaying a fixed out-of-order event log and asserting the running sum, the fired-window set, and the retraction pairs all match a batch reference computed on the same log sorted by event time.", "environment": "node20", "language": "js", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "js-event-time-session-window-late-merge-retraction-semantics", "provider": "openrouter", "solved_at": "2026-09-22T04:37:28.097Z", "version": "20"}

Answer 2

All 5 tests pass (including 200 fuzz trials). Here is the verified solution.


Event-Time Session Windows with Late-Merge Retractions, Allowed Lateness, and Idle-Source Watermarks

1. Root-cause analysis

A streaming session-window operator fails in four characteristic ways:

1.1 Firing is easy — un-firing is the hard part. The naive operator emits (window, sum) when the watermark passes end = maxTs + gap and then discards the window. A later element that bridges two emitted windows can no longer reach them, so it creates an orphan window and the downstream SUM permanently contains the two old sums plus the new one. Fix: retain fired sessions through the allowed-lateness period; when a late element absorbs them, emit -emittedSum retractions followed by one merged +mergedSum replacement. Since downstream is a signed sum, -a - b + (a+b+late) = late.

1.2 Merging is transitive and N-ary. Collect every live session the element's gap-expanded interval touches — not just the nearest — then retract all fired ones and emit a single replacement. Hard-coding "merge two neighbours" breaks for 3+ sessions or on a later re-merge.

1.3 Drop and purge must share one horizon. Purge a fired session when watermark > end + gap + allowedLateness; drop an unattached element when watermark > time + gap + allowedLateness. Any element that can merge a live session satisfies time <= session.end, so a purged session implies the element is already past its drop horizon — the decisions are consistent by construction.

1.4 Idle sources. Removing an idle partition from the watermark min lets the watermark run past data it may still produce. Instead, an idle partition contributes its endOfInput to the min, so the combined watermark advances up to that horizon but never past it. endOfInput = Infinity means no cap.

2. The fix

session-window.js:

'use strict';

class WatermarkTracker {
  constructor({ onWatermark = () => {} } = {}) {
    this.onWatermark = onWatermark;
    this.partitions = new Map();
    this.watermark = -Infinity;
  }
  register(id, { endOfInput = Infinity } = {}) {
    if (this.partitions.has(id)) throw new Error(`partition ${id} already registered`);
    this.partitions.set(id, { id, watermark: -Infinity, idle: false, endOfInput });
    this._recompute(); return this;
  }
  report(id, watermark) {
    const p = this._get(id);
    if (!Number.isFinite(watermark)) throw new Error('watermark must be finite');
    if (watermark > p.watermark) p.watermark = watermark;
    p.idle = false; this._recompute(); return this;
  }
  setIdle(id, idle = true) { this._get(id).idle = !!idle; this._recompute(); return this; }
  setEndOfInput(id, endOfInput) {
    const p = this._get(id);
    if (!Number.isFinite(endOfInput)) throw new Error('endOfInput must be finite');
    if (endOfInput < p.watermark) throw new Error(`endOfInput ${endOfInput} below watermark ${p.watermark}`);
    p.endOfInput = endOfInput; p.idle = true; this._recompute(); return this;
  }
  _get(id) { const p = this.partitions.get(id); if (!p) throw new Error(`unknown partition ${id}`); return p; }
  _effective(p) { return p.idle ? p.endOfInput : p.watermark; }
  _recompute() {
    let min = Infinity;
    for (const p of this.partitions.values()) min = Math.min(min, this._effective(p));
    if (min > this.watermark) { this.watermark = min; this.onWatermark(min); }
  }
}

class SessionWindowOperator {
  constructor({ gap, allowedLateness = 0, onEmit = () => {}, onSideOutput = () => {} } = {}) {
    if (!Number.isFinite(gap) || gap < 0) throw new Error('gap must be a finite non-negative number');
    if (!Number.isFinite(allowedLateness) || allowedLateness < 0) throw new Error('allowedLateness must be finite and >= 0');
    this.gap = gap; this.allowedLateness = allowedLateness;
    this.onEmit = onEmit; this.onSideOutput = onSideOutput;
    this.watermark = -Infinity;
    this.sessions = new Map(); // key -> Session[]
    this.emitted = new Map();  // sessionId -> currently reflected window
    this.seq = 0;
  }
  _newSession(key, start, maxTs, sum) {
    return { id: `${key}:${++this.seq}`, key, start, maxTs, end: maxTs + this.gap,
             sum, fired: false, emittedSum: 0 };
  }
  _mergeable(s, time) { return time <= s.maxTs + this.gap && time >= s.start - this.gap; }
  _emit(e) {
    if (e.type === 'fire') this.emitted.set(e.sessionId,
      { key: e.key, start: e.start, maxTs: e.maxTs, end: e.end, sum: e.delta });
    else if (e.type === 'retract') this.emitted.delete(e.retracts);
    this.onEmit(e);
  }
  receive(event) {
    const { key, time, value } = event;
    if (typeof key === 'undefined') throw new Error('event.key is required');
    if (!Number.isFinite(time)) throw new Error('event.time must be finite');
    if (!Number.isFinite(value)) throw new Error('event.value must be finite');

    const list = this.sessions.get(key) || [];
    const merging = [], keep = [];
    for (const s of list) { if (this._mergeable(s, time)) merging.push(s); else keep.push(s); }

    if (merging.length === 0 && Number.isFinite(this.watermark) &&
        this.watermark > time + this.gap + this.allowedLateness) {
      this.onSideOutput({ ...event, reason: 'allowed-lateness-expired' }, 'too-late');
      return false;
    }

    let start = time, maxTs = time, sum = value;
    for (const s of merging) {
      if (s.start < start) start = s.start;
      if (s.maxTs > maxTs) maxTs = s.maxTs;
      sum += s.sum;
    }
    const merged = this._newSession(key, start, maxTs, sum);
    const firedMerged = merging.filter((s) => s.fired);

    if (firedMerged.length > 0) {
      for (const s of firedMerged) {
        this._emit({ type: 'retract', reason: 'late-merge', key, sessionId: s.id,
                     start: s.start, maxTs: s.maxTs, end: s.end,
                     delta: -s.emittedSum, retracts: s.id });
      }
      merged.fired = true; merged.emittedSum = merged.sum;
      this._emit({ type: 'fire', reason: 'late-merge-replacement', key, sessionId: merged.id,
                   start: merged.start, maxTs: merged.maxTs, end: merged.end,
                   delta: merged.sum, mergedFrom: merging.map((s) => s.id) });
    } else if (Number.isFinite(this.watermark) && this.watermark >= merged.end) {
      merged.fired = true; merged.emittedSum = merged.sum;
      this._emit({ type: 'fire', reason: 'watermark', key, sessionId: merged.id,
                   start: merged.start, maxTs: merged.maxTs, end: merged.end, delta: merged.sum });
    }
    keep.push(merged);
    this.sessions.set(key, keep);
    return true;
  }
  advanceWatermark(watermark) {
    if (watermark !== Infinity && !Number.isFinite(watermark)) throw new Error('watermark must be finite or Infinity');
    if (watermark <= this.watermark) return;
    this.watermark = watermark;
    for (const [key, list] of [...this.sessions.entries()]) {
      const keep = [];
      for (const s of list) {
        if (!s.fired && this.watermark >= s.end) {
          s.fired = true; s.emittedSum = s.sum;
          this._emit({ type: 'fire', reason: 'watermark', key, sessionId: s.id,
                       start: s.start, maxTs: s.maxTs, end: s.end, delta: s.sum });
        }
        if (s.fired && this.watermark > s.end + this.gap + this.allowedLateness) continue;
        keep.push(s);
      }
      if (keep.length) this.sessions.set(key, keep); else this.sessions.delete(key);
    }
  }
  snapshot() {
    const out = [];
    for (const list of this.sessions.values()) for (const s of list) out.push({ ...s });
    return out;
  }
  currentWindows() { return [...this.emitted.values()].map((w) => ({ ...w })); }
  seedFiredSession(key, { start, maxTs, sum }, { emit = false } = {}) {
    const s = this._newSession(key, start, maxTs, sum);
    s.fired = true; s.emittedSum = sum;
    const list = this.sessions.get(key) || []; list.push(s); this.sessions.set(key, list);
    if (emit) this._emit({ type: 'fire', reason: 'seed', key, sessionId: s.id,
                           start: s.start, maxTs: s.maxTs, end: s.end, delta: s.sum });
    return s;
  }
}

module.exports = { SessionWindowOperator, WatermarkTracker };

3. Verification

Batch reference (sort by time, greedily group where diff <= gap):

function batchSessions(events, gap) {
  const byKey = new Map();
  for (const e of events) { if (!byKey.has(e.key)) byKey.set(e.key, []); byKey.get(e.key).push(e); }
  const result = [];
  for (const [key, list] of byKey) {
    list.sort((a, b) => a.time - b.time || a.arrival - b.arrival);
    let cur = null;
    for (const e of list) {
      if (cur === null || e.time - cur.maxTs > gap) {
        cur = { key, start: e.time, maxTs: e.time, sum: e.value }; result.push(cur);
      } else { cur.maxTs = e.time; cur.sum += e.value; }
    }
  }
  return result;
}

Run node --test (files: session-window.js, session-window.test.js). Result:

ok 1 - replay: running sum, fired set and retraction pairs match batch
ok 2 - N-way merge: one late event absorbs three already-fired sessions
ok 3 - idle source caps watermark at its end-of-input
ok 4 - allowed lateness lets a bridging late event through but drops a stranded one
ok 5 - fuzz: random out-of-order logs match the batch reference exactly
# tests 5  # pass 5  # fail 0

Replay scenario (gap=10, allowedLateness=10): receive (A,0,1),(A,5,2); wm=15 fires [0,5]→3. receive (A,22,3); wm=32 fires [22,22]→3. receive (A,13,100) bridges both → retract −3, −3 and fire merged [0,22]→106. Then [40,45]→9, [80,80]→6, and (A,60,7) at wm=90 is stranded → side output. Assertions:

running === batchTotal === 121
currentWindows() deep-equals batch: [0,22]=106, [40,45]=9, [80,80]=6
retractions.length === 2, deltas === [-3,-3], each references a prior +3 fire
replacement.start===0, maxTs===22, delta===106, mergedFrom.length===2
(-3) + (-3) + 106 === 100   // the bridging element's value

N-way: three seeded fired sessions [0,1]→1, [3,4]→2, [6,7]→3; one element (B,4,100) emits three retractions -1,-2,-3 and one replacement [0,7]→106; running sum 106.

Idle: p0.endOfInput=25, p1 unbounded, report(p0,10), report(p1,30) → 10; p0 idle → 25; report(p1,1000) → still 25 (never past the idle horizon).

Fuzz: 200 deterministic random trials with shuffled arrival order; for each, sum(emissions.delta) === batchTotal and currentWindows() deep-equals batchSessions(...), exercising the retraction path repeatedly.

The complete write-up, full code, and the test file are saved in ~/event-time-sessions/ (solution.md, session-window.js, session-window.test.js).

Evidence & signatures

# Evidence
- Problem class: js-event-time-session-window-late-merge-retraction-semantics
- Model: openrouter/deepseek/deepseek-v4.1-flash
- Solved: 2026-09-22T04:37:28.097Z
- Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Implement a streaming event-time windowing operator in Node.js that groups events into session windows with a configurable inactivity gap, fires a window only when the watermark passes its end, and when a late event bridges two already-fired sessions must emit a retraction that cancels the earlier contributions plus the merged replacement so a downstream SUM recomputes to exactly the batch value. The operator must also handle N-way merges (one late event joining three or more fired sessions), allowed-lateness expiry that routes dropped late events to a side output, and idle sources that must never let a watermark advance past an idle partition's end-of-input. Verify by replaying a fixed out-of-order event log and asserting the running sum, the fired-window set, and the retraction pairs all match a batch reference computed on the same log sorted by event time.", "environment": "node20", "language": "js", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "js-event-time-session-window-late-merge-retraction-semantics", "provider": "openrouter", "solved_at": "2026-09-22T04:37:28.097Z", "version": "20"}
Generated from the verified corpus · MIT licensedBack to the catalog