FAILURE MAP
← Case archive

FA-45671 / Data systems / Open access

Watermark buffering releases events at the frontier prematurely · case 01

Watermark buffering releases events at the frontier prematurely.

Verified by executionVariant 1 · 7 checks per implementationDownload source bundle ↓JSON ↗

ROOT CAUSE

watermark-buffer-release: Watermark buffering releases events at the frontier prematurely.

VERIFIED REPAIR

Preserve the stated physical representation and operation order: Process event actions [event,time,id] and watermark actions [watermark,time,None]. Maintain monotone watermark; events older than it go to late output. On a watermark, release buffered events with time strictly below watermark ordered by (time,id). Return [emitted,late,sorted-pending].

Unsuccessful approach: Adding one retains the premature frontier release.

Case contract

Process event actions [event,time,id] and watermark actions [watermark,time,None]. Maintain monotone watermark; events older than it go to late output. On a watermark, release buffered events with time strictly below watermark ordered by (time,id). Return [emitted,late,sorted-pending].

Why this case matters

A bounded deterministic data engine model makes representation and changelog faults reproducible.

1 / The failure

Exit 1
"""Failure Map reference implementation. Python standard library only."""
import json

N = 1
observations = []
def solve(d):
    try:
        watermark=None; pending=[]; emitted=[]; late=[]
        for kind,time,ident in d:
            if kind=='event':
                if watermark is not None and time<watermark: late.append([time,ident])
                else: pending.append([time,ident])
            else:
                watermark=time if watermark is None else max(watermark,time)
                ready=sorted(e for e in pending if e[0]<=watermark)
                emitted.extend(ready)
                pending=[e for e in pending if e[0]>=watermark]
        return [emitted,late,sorted(pending)]
    except (IndexError, KeyError, ValueError, StopIteration) as exc:
        return {"representation_error": type(exc).__name__}
def check(label, actual, expected):
    observations.append({"check": label, "actual": actual, "expected": expected, "passed": actual == expected})
if N == 1:
    check('frontier equality', solve([['watermark', 1, None], ['event', 1, 10]]), [[], [], [[1, 10]]])
    check('regression', solve([['watermark', 4, None], ['watermark', 1, None], ['event', 2, 10]]), [[], [[2, 10]], []])
    check('frontier stays pending', solve([['event', 1, 10], ['watermark', 1, None]]), [[], [], [[1, 10]]])
    check('event time ordering', solve([['event', 3, 10], ['event', 1, 11], ['watermark', 4, None]]), [[[1, 11], [3, 10]], [], []])
    check('emit once', solve([['event', 1, 10], ['watermark', 2, None], ['watermark', 3, None]]), [[[1, 10]], [], []])
    check('late row', solve([['watermark', 3, None], ['event', 1, 10]]), [[], [[1, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 2:
    check('frontier equality', solve([['watermark', 2, None], ['event', 2, 10]]), [[], [], [[2, 10]]])
    check('regression', solve([['watermark', 5, None], ['watermark', 2, None], ['event', 3, 10]]), [[], [[3, 10]], []])
    check('frontier stays pending', solve([['event', 2, 10], ['watermark', 2, None]]), [[], [], [[2, 10]]])
    check('event time ordering', solve([['event', 4, 10], ['event', 2, 11], ['watermark', 5, None]]), [[[2, 11], [4, 10]], [], []])
    check('emit once', solve([['event', 2, 10], ['watermark', 3, None], ['watermark', 4, None]]), [[[2, 10]], [], []])
    check('late row', solve([['watermark', 4, None], ['event', 2, 10]]), [[], [[2, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 3:
    check('frontier equality', solve([['watermark', 3, None], ['event', 3, 10]]), [[], [], [[3, 10]]])
    check('regression', solve([['watermark', 6, None], ['watermark', 3, None], ['event', 4, 10]]), [[], [[4, 10]], []])
    check('frontier stays pending', solve([['event', 3, 10], ['watermark', 3, None]]), [[], [], [[3, 10]]])
    check('event time ordering', solve([['event', 5, 10], ['event', 3, 11], ['watermark', 6, None]]), [[[3, 11], [5, 10]], [], []])
    check('emit once', solve([['event', 3, 10], ['watermark', 4, None], ['watermark', 5, None]]), [[[3, 10]], [], []])
    check('late row', solve([['watermark', 5, None], ['event', 3, 10]]), [[], [[3, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 4:
    check('frontier equality', solve([['watermark', 4, None], ['event', 4, 10]]), [[], [], [[4, 10]]])
    check('regression', solve([['watermark', 7, None], ['watermark', 4, None], ['event', 5, 10]]), [[], [[5, 10]], []])
    check('frontier stays pending', solve([['event', 4, 10], ['watermark', 4, None]]), [[], [], [[4, 10]]])
    check('event time ordering', solve([['event', 6, 10], ['event', 4, 11], ['watermark', 7, None]]), [[[4, 11], [6, 10]], [], []])
    check('emit once', solve([['event', 4, 10], ['watermark', 5, None], ['watermark', 6, None]]), [[[4, 10]], [], []])
    check('late row', solve([['watermark', 6, None], ['event', 4, 10]]), [[], [[4, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 5:
    check('frontier equality', solve([['watermark', 5, None], ['event', 5, 10]]), [[], [], [[5, 10]]])
    check('regression', solve([['watermark', 8, None], ['watermark', 5, None], ['event', 6, 10]]), [[], [[6, 10]], []])
    check('frontier stays pending', solve([['event', 5, 10], ['watermark', 5, None]]), [[], [], [[5, 10]]])
    check('event time ordering', solve([['event', 7, 10], ['event', 5, 11], ['watermark', 8, None]]), [[[5, 11], [7, 10]], [], []])
    check('emit once', solve([['event', 5, 10], ['watermark', 6, None], ['watermark', 7, None]]), [[[5, 10]], [], []])
    check('late row', solve([['watermark', 7, None], ['event', 5, 10]]), [[], [[5, 10]], []])
    check('empty input', solve([]), [[], [], []])
print(json.dumps({"observations": observations, "passed": all(x["passed"] for x in observations)}, ensure_ascii=False))
raise SystemExit(0 if all(x["passed"] for x in observations) else 1)
Boundary fixtureActualExpectedOutcome
frontier equality[[], [], [[1, 10]]][[], [], [[1, 10]]]Passed
regression[[], [[2, 10]], []][[], [[2, 10]], []]Passed
frontier stays pending[[[1, 10]], [], [[1, 10]]][[], [], [[1, 10]]]Failed
event time ordering[[[1, 11], [3, 10]], [], []][[[1, 11], [3, 10]], [], []]Passed
emit once[[[1, 10]], [], []][[[1, 10]], [], []]Passed
late row[[], [[1, 10]], []][[], [[1, 10]], []]Passed
empty input[[], [], []][[], [], []]Passed

SHA-256 / 7f72b03b080d49967b0814882d4352defafd55a18c1ab7df540b4854c0d527ac

2 / The unsuccessful fix

Exit 1
"""Failure Map reference implementation. Python standard library only."""
import json

N = 1
observations = []
def solve(d):
    try:
        watermark=None; pending=[]; emitted=[]; late=[]
        for kind,time,ident in d:
            if kind=='event':
                if watermark is not None and time<watermark: late.append([time,ident])
                else: pending.append([time,ident])
            else:
                watermark=time if watermark is None else max(watermark,time)
                ready=sorted(e for e in pending if e[0]<watermark+1)
                emitted.extend(ready)
                pending=[e for e in pending if e[0]>=watermark]
        return [emitted,late,sorted(pending)]
    except (IndexError, KeyError, ValueError, StopIteration) as exc:
        return {"representation_error": type(exc).__name__}
def check(label, actual, expected):
    observations.append({"check": label, "actual": actual, "expected": expected, "passed": actual == expected})
if N == 1:
    check('frontier equality', solve([['watermark', 1, None], ['event', 1, 10]]), [[], [], [[1, 10]]])
    check('regression', solve([['watermark', 4, None], ['watermark', 1, None], ['event', 2, 10]]), [[], [[2, 10]], []])
    check('frontier stays pending', solve([['event', 1, 10], ['watermark', 1, None]]), [[], [], [[1, 10]]])
    check('event time ordering', solve([['event', 3, 10], ['event', 1, 11], ['watermark', 4, None]]), [[[1, 11], [3, 10]], [], []])
    check('emit once', solve([['event', 1, 10], ['watermark', 2, None], ['watermark', 3, None]]), [[[1, 10]], [], []])
    check('late row', solve([['watermark', 3, None], ['event', 1, 10]]), [[], [[1, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 2:
    check('frontier equality', solve([['watermark', 2, None], ['event', 2, 10]]), [[], [], [[2, 10]]])
    check('regression', solve([['watermark', 5, None], ['watermark', 2, None], ['event', 3, 10]]), [[], [[3, 10]], []])
    check('frontier stays pending', solve([['event', 2, 10], ['watermark', 2, None]]), [[], [], [[2, 10]]])
    check('event time ordering', solve([['event', 4, 10], ['event', 2, 11], ['watermark', 5, None]]), [[[2, 11], [4, 10]], [], []])
    check('emit once', solve([['event', 2, 10], ['watermark', 3, None], ['watermark', 4, None]]), [[[2, 10]], [], []])
    check('late row', solve([['watermark', 4, None], ['event', 2, 10]]), [[], [[2, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 3:
    check('frontier equality', solve([['watermark', 3, None], ['event', 3, 10]]), [[], [], [[3, 10]]])
    check('regression', solve([['watermark', 6, None], ['watermark', 3, None], ['event', 4, 10]]), [[], [[4, 10]], []])
    check('frontier stays pending', solve([['event', 3, 10], ['watermark', 3, None]]), [[], [], [[3, 10]]])
    check('event time ordering', solve([['event', 5, 10], ['event', 3, 11], ['watermark', 6, None]]), [[[3, 11], [5, 10]], [], []])
    check('emit once', solve([['event', 3, 10], ['watermark', 4, None], ['watermark', 5, None]]), [[[3, 10]], [], []])
    check('late row', solve([['watermark', 5, None], ['event', 3, 10]]), [[], [[3, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 4:
    check('frontier equality', solve([['watermark', 4, None], ['event', 4, 10]]), [[], [], [[4, 10]]])
    check('regression', solve([['watermark', 7, None], ['watermark', 4, None], ['event', 5, 10]]), [[], [[5, 10]], []])
    check('frontier stays pending', solve([['event', 4, 10], ['watermark', 4, None]]), [[], [], [[4, 10]]])
    check('event time ordering', solve([['event', 6, 10], ['event', 4, 11], ['watermark', 7, None]]), [[[4, 11], [6, 10]], [], []])
    check('emit once', solve([['event', 4, 10], ['watermark', 5, None], ['watermark', 6, None]]), [[[4, 10]], [], []])
    check('late row', solve([['watermark', 6, None], ['event', 4, 10]]), [[], [[4, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 5:
    check('frontier equality', solve([['watermark', 5, None], ['event', 5, 10]]), [[], [], [[5, 10]]])
    check('regression', solve([['watermark', 8, None], ['watermark', 5, None], ['event', 6, 10]]), [[], [[6, 10]], []])
    check('frontier stays pending', solve([['event', 5, 10], ['watermark', 5, None]]), [[], [], [[5, 10]]])
    check('event time ordering', solve([['event', 7, 10], ['event', 5, 11], ['watermark', 8, None]]), [[[5, 11], [7, 10]], [], []])
    check('emit once', solve([['event', 5, 10], ['watermark', 6, None], ['watermark', 7, None]]), [[[5, 10]], [], []])
    check('late row', solve([['watermark', 7, None], ['event', 5, 10]]), [[], [[5, 10]], []])
    check('empty input', solve([]), [[], [], []])
print(json.dumps({"observations": observations, "passed": all(x["passed"] for x in observations)}, ensure_ascii=False))
raise SystemExit(0 if all(x["passed"] for x in observations) else 1)
Boundary fixtureActualExpectedOutcome
frontier equality[[], [], [[1, 10]]][[], [], [[1, 10]]]Passed
regression[[], [[2, 10]], []][[], [[2, 10]], []]Passed
frontier stays pending[[[1, 10]], [], [[1, 10]]][[], [], [[1, 10]]]Failed
event time ordering[[[1, 11], [3, 10]], [], []][[[1, 11], [3, 10]], [], []]Passed
emit once[[[1, 10]], [], []][[[1, 10]], [], []]Passed
late row[[], [[1, 10]], []][[], [[1, 10]], []]Passed
empty input[[], [], []][[], [], []]Passed

SHA-256 / 1e179148d7197fb67a639f7819cc7740e303ad359743bcd8f5ae77e0cdfcbc36

3 / The verified repair

Exit 0
"""Failure Map reference implementation. Python standard library only."""
import json

N = 1
observations = []
def solve(d):
    try:
        watermark=None; pending=[]; emitted=[]; late=[]
        for kind,time,ident in d:
            if kind=='event':
                if watermark is not None and time<watermark: late.append([time,ident])
                else: pending.append([time,ident])
            else:
                watermark=time if watermark is None else max(watermark,time)
                ready=sorted(e for e in pending if e[0]<watermark)
                emitted.extend(ready)
                pending=[e for e in pending if e[0]>=watermark]
        return [emitted,late,sorted(pending)]
    except (IndexError, KeyError, ValueError, StopIteration) as exc:
        return {"representation_error": type(exc).__name__}
def check(label, actual, expected):
    observations.append({"check": label, "actual": actual, "expected": expected, "passed": actual == expected})
if N == 1:
    check('frontier equality', solve([['watermark', 1, None], ['event', 1, 10]]), [[], [], [[1, 10]]])
    check('regression', solve([['watermark', 4, None], ['watermark', 1, None], ['event', 2, 10]]), [[], [[2, 10]], []])
    check('frontier stays pending', solve([['event', 1, 10], ['watermark', 1, None]]), [[], [], [[1, 10]]])
    check('event time ordering', solve([['event', 3, 10], ['event', 1, 11], ['watermark', 4, None]]), [[[1, 11], [3, 10]], [], []])
    check('emit once', solve([['event', 1, 10], ['watermark', 2, None], ['watermark', 3, None]]), [[[1, 10]], [], []])
    check('late row', solve([['watermark', 3, None], ['event', 1, 10]]), [[], [[1, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 2:
    check('frontier equality', solve([['watermark', 2, None], ['event', 2, 10]]), [[], [], [[2, 10]]])
    check('regression', solve([['watermark', 5, None], ['watermark', 2, None], ['event', 3, 10]]), [[], [[3, 10]], []])
    check('frontier stays pending', solve([['event', 2, 10], ['watermark', 2, None]]), [[], [], [[2, 10]]])
    check('event time ordering', solve([['event', 4, 10], ['event', 2, 11], ['watermark', 5, None]]), [[[2, 11], [4, 10]], [], []])
    check('emit once', solve([['event', 2, 10], ['watermark', 3, None], ['watermark', 4, None]]), [[[2, 10]], [], []])
    check('late row', solve([['watermark', 4, None], ['event', 2, 10]]), [[], [[2, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 3:
    check('frontier equality', solve([['watermark', 3, None], ['event', 3, 10]]), [[], [], [[3, 10]]])
    check('regression', solve([['watermark', 6, None], ['watermark', 3, None], ['event', 4, 10]]), [[], [[4, 10]], []])
    check('frontier stays pending', solve([['event', 3, 10], ['watermark', 3, None]]), [[], [], [[3, 10]]])
    check('event time ordering', solve([['event', 5, 10], ['event', 3, 11], ['watermark', 6, None]]), [[[3, 11], [5, 10]], [], []])
    check('emit once', solve([['event', 3, 10], ['watermark', 4, None], ['watermark', 5, None]]), [[[3, 10]], [], []])
    check('late row', solve([['watermark', 5, None], ['event', 3, 10]]), [[], [[3, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 4:
    check('frontier equality', solve([['watermark', 4, None], ['event', 4, 10]]), [[], [], [[4, 10]]])
    check('regression', solve([['watermark', 7, None], ['watermark', 4, None], ['event', 5, 10]]), [[], [[5, 10]], []])
    check('frontier stays pending', solve([['event', 4, 10], ['watermark', 4, None]]), [[], [], [[4, 10]]])
    check('event time ordering', solve([['event', 6, 10], ['event', 4, 11], ['watermark', 7, None]]), [[[4, 11], [6, 10]], [], []])
    check('emit once', solve([['event', 4, 10], ['watermark', 5, None], ['watermark', 6, None]]), [[[4, 10]], [], []])
    check('late row', solve([['watermark', 6, None], ['event', 4, 10]]), [[], [[4, 10]], []])
    check('empty input', solve([]), [[], [], []])
elif N == 5:
    check('frontier equality', solve([['watermark', 5, None], ['event', 5, 10]]), [[], [], [[5, 10]]])
    check('regression', solve([['watermark', 8, None], ['watermark', 5, None], ['event', 6, 10]]), [[], [[6, 10]], []])
    check('frontier stays pending', solve([['event', 5, 10], ['watermark', 5, None]]), [[], [], [[5, 10]]])
    check('event time ordering', solve([['event', 7, 10], ['event', 5, 11], ['watermark', 8, None]]), [[[5, 11], [7, 10]], [], []])
    check('emit once', solve([['event', 5, 10], ['watermark', 6, None], ['watermark', 7, None]]), [[[5, 10]], [], []])
    check('late row', solve([['watermark', 7, None], ['event', 5, 10]]), [[], [[5, 10]], []])
    check('empty input', solve([]), [[], [], []])
print(json.dumps({"observations": observations, "passed": all(x["passed"] for x in observations)}, ensure_ascii=False))
raise SystemExit(0 if all(x["passed"] for x in observations) else 1)
Boundary fixtureActualExpectedOutcome
frontier equality[[], [], [[1, 10]]][[], [], [[1, 10]]]Passed
regression[[], [[2, 10]], []][[], [[2, 10]], []]Passed
frontier stays pending[[], [], [[1, 10]]][[], [], [[1, 10]]]Passed
event time ordering[[[1, 11], [3, 10]], [], []][[[1, 11], [3, 10]], [], []]Passed
emit once[[[1, 10]], [], []][[[1, 10]], [], []]Passed
late row[[], [[1, 10]], []][[], [[1, 10]], []]Passed
empty input[[], [], []][[], [], []]Passed

SHA-256 / 15b2eb00a8b04c69e6f07835ab64ddf0303420db98d87d21c7e99daef68f39a2

Verification & scope

Offline stipulated semantics over valid small inputs; no performance, concurrency, or production-engine conformance claim. This reproducer isolates one failure mechanism. Results cover the supplied fixtures. Variants within a family share a test contract and should remain grouped when constructing evaluation splits. Related mechanisms with a shared evaluation_group must also remain together; these controlled models are not independent production incidents.

Observations recorded using Python 3.12.14 at 2026-09-29T14:44:24.588946+00:00.

Case digest / 790b9f166f7cc591e01f32e79468b18d161883ff4a86b268eedf93bd24218ca3