FA-1551 / Replication / Open access
Complete a streaming checkpoint barrier: A checkpoint completes before every input channel reaches its barrier · case 01
The checkpoint barrier operation is admitted even though a checkpoint completes before every input channel reaches its barrier.
ROOT CAUSE
The admission path omits the all inputs invariant while validating the other operation preconditions.
VERIFIED REPAIR
Require set(r['all_inputs'][0]) == set(r['all_inputs'][1]) together with every other stated precondition before accepting the operation.
Unsuccessful approach: Adding the all inputs check repairs the reported defect, but replacing the adjacent barrier id check loses that independent invariant.
Case contract
Return a Boolean admission decision for complete a streaming checkpoint barrier. The record r must satisfy all of: set(r['all_inputs'][0]) == set(r['all_inputs'][1]); len(set(r['barrier_id'])) == 1; r['state_flushed'] is True; r['inflight_captured'][0] == r['inflight_captured'][1]; r['sink_prepared'] == 'prepared'. Extra tracing fields are ignored; validation does not mutate the record.
Why this case matters
A deterministic local contract for replication. Each negative fixture violates exactly one invariant. No transport timing, persistence, cryptographic verification, or full protocol implementation is claimed.
1 / The failure
Exit 1"""Failure Map reference implementation. Python standard library only."""
import json
N = 1
observations = []
def solve(r):
return (len(set(r['barrier_id'])) == 1) and (r['state_flushed'] is True) and (r['inflight_captured'][0] == r['inflight_captured'][1]) and (r['sink_prepared'] == 'prepared')
def check(label, actual, expected):
observations.append({"check": label, "actual": actual, "expected": expected, "passed": actual == expected})
r = {'all_inputs': [['a', 'b'], ['a', 'b']], 'barrier_id': [4, 4], 'state_flushed': True, 'inflight_captured': [8, 8], 'sink_prepared': 'prepared'}
check('valid operation', solve(r), True)
check('A checkpoint completes before every input channel reaches its barrier', solve(dict(r, **{'all_inputs': [['a', 'b'], ['a']]})), False)
check('Barriers from different checkpoint rounds are combined', solve(dict(r, **{'barrier_id': [4, 5]})), False)
check('A checkpoint is acknowledged while state exists only in memory', solve(dict(r, **{'state_flushed': False})), False)
check('Unaligned checkpoint omits messages already in transit', solve(dict(r, **{'inflight_captured': [5, 8]})), False)
check('Upstream progress commits before the transactional sink prepares', solve(dict(r, **{'sink_prepared': 'open'})), False)
check('unrelated tracing metadata', solve(dict(r, trace='run-'+str(N))), True)
check('repeat validation is pure', solve(r), True)
invalid = {'all_inputs': [['a', 'b'], ['a']], 'barrier_id': [4, 5], 'state_flushed': False, 'inflight_captured': [5, 8], 'sink_prepared': 'open'}
keys = list(invalid)
pair = {keys[N % len(keys)]: invalid[keys[N % len(keys)]], keys[(N+1) % len(keys)]: invalid[keys[(N+1) % len(keys)]]}
check('two independent violations in variant', solve(dict(r, **pair)), False)
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 fixture | Actual | Expected | Outcome |
|---|---|---|---|
| valid operation | True | True | Passed |
| A checkpoint completes before every input channel reaches its barrier | True | False | Failed |
| Barriers from different checkpoint rounds are combined | False | False | Passed |
| A checkpoint is acknowledged while state exists only in memory | False | False | Passed |
| Unaligned checkpoint omits messages already in transit | False | False | Passed |
| Upstream progress commits before the transactional sink prepares | False | False | Passed |
| unrelated tracing metadata | True | True | Passed |
| repeat validation is pure | True | True | Passed |
| two independent violations in variant | False | False | Passed |
SHA-256 / 7e2a7a6be96d79df922a6555205d9a43e6e90962cd1aa3c9f63575a799e83b57
2 / The unsuccessful fix
Exit 1"""Failure Map reference implementation. Python standard library only."""
import json
N = 1
observations = []
def solve(r):
return (set(r['all_inputs'][0]) == set(r['all_inputs'][1])) and (r['state_flushed'] is True) and (r['inflight_captured'][0] == r['inflight_captured'][1]) and (r['sink_prepared'] == 'prepared')
def check(label, actual, expected):
observations.append({"check": label, "actual": actual, "expected": expected, "passed": actual == expected})
r = {'all_inputs': [['a', 'b'], ['a', 'b']], 'barrier_id': [4, 4], 'state_flushed': True, 'inflight_captured': [8, 8], 'sink_prepared': 'prepared'}
check('valid operation', solve(r), True)
check('A checkpoint completes before every input channel reaches its barrier', solve(dict(r, **{'all_inputs': [['a', 'b'], ['a']]})), False)
check('Barriers from different checkpoint rounds are combined', solve(dict(r, **{'barrier_id': [4, 5]})), False)
check('A checkpoint is acknowledged while state exists only in memory', solve(dict(r, **{'state_flushed': False})), False)
check('Unaligned checkpoint omits messages already in transit', solve(dict(r, **{'inflight_captured': [5, 8]})), False)
check('Upstream progress commits before the transactional sink prepares', solve(dict(r, **{'sink_prepared': 'open'})), False)
check('unrelated tracing metadata', solve(dict(r, trace='run-'+str(N))), True)
check('repeat validation is pure', solve(r), True)
invalid = {'all_inputs': [['a', 'b'], ['a']], 'barrier_id': [4, 5], 'state_flushed': False, 'inflight_captured': [5, 8], 'sink_prepared': 'open'}
keys = list(invalid)
pair = {keys[N % len(keys)]: invalid[keys[N % len(keys)]], keys[(N+1) % len(keys)]: invalid[keys[(N+1) % len(keys)]]}
check('two independent violations in variant', solve(dict(r, **pair)), False)
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 fixture | Actual | Expected | Outcome |
|---|---|---|---|
| valid operation | True | True | Passed |
| A checkpoint completes before every input channel reaches its barrier | False | False | Passed |
| Barriers from different checkpoint rounds are combined | True | False | Failed |
| A checkpoint is acknowledged while state exists only in memory | False | False | Passed |
| Unaligned checkpoint omits messages already in transit | False | False | Passed |
| Upstream progress commits before the transactional sink prepares | False | False | Passed |
| unrelated tracing metadata | True | True | Passed |
| repeat validation is pure | True | True | Passed |
| two independent violations in variant | False | False | Passed |
SHA-256 / b5b569f1d1b829ca0fe12444e59c4406809dcb6fe543b34243b2952e01eb1996
3 / The verified repair
Exit 0"""Failure Map reference implementation. Python standard library only."""
import json
N = 1
observations = []
def solve(r):
return (set(r['all_inputs'][0]) == set(r['all_inputs'][1])) and (len(set(r['barrier_id'])) == 1) and (r['state_flushed'] is True) and (r['inflight_captured'][0] == r['inflight_captured'][1]) and (r['sink_prepared'] == 'prepared')
def check(label, actual, expected):
observations.append({"check": label, "actual": actual, "expected": expected, "passed": actual == expected})
r = {'all_inputs': [['a', 'b'], ['a', 'b']], 'barrier_id': [4, 4], 'state_flushed': True, 'inflight_captured': [8, 8], 'sink_prepared': 'prepared'}
check('valid operation', solve(r), True)
check('A checkpoint completes before every input channel reaches its barrier', solve(dict(r, **{'all_inputs': [['a', 'b'], ['a']]})), False)
check('Barriers from different checkpoint rounds are combined', solve(dict(r, **{'barrier_id': [4, 5]})), False)
check('A checkpoint is acknowledged while state exists only in memory', solve(dict(r, **{'state_flushed': False})), False)
check('Unaligned checkpoint omits messages already in transit', solve(dict(r, **{'inflight_captured': [5, 8]})), False)
check('Upstream progress commits before the transactional sink prepares', solve(dict(r, **{'sink_prepared': 'open'})), False)
check('unrelated tracing metadata', solve(dict(r, trace='run-'+str(N))), True)
check('repeat validation is pure', solve(r), True)
invalid = {'all_inputs': [['a', 'b'], ['a']], 'barrier_id': [4, 5], 'state_flushed': False, 'inflight_captured': [5, 8], 'sink_prepared': 'open'}
keys = list(invalid)
pair = {keys[N % len(keys)]: invalid[keys[N % len(keys)]], keys[(N+1) % len(keys)]: invalid[keys[(N+1) % len(keys)]]}
check('two independent violations in variant', solve(dict(r, **pair)), False)
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 fixture | Actual | Expected | Outcome |
|---|---|---|---|
| valid operation | True | True | Passed |
| A checkpoint completes before every input channel reaches its barrier | False | False | Passed |
| Barriers from different checkpoint rounds are combined | False | False | Passed |
| A checkpoint is acknowledged while state exists only in memory | False | False | Passed |
| Unaligned checkpoint omits messages already in transit | False | False | Passed |
| Upstream progress commits before the transactional sink prepares | False | False | Passed |
| unrelated tracing metadata | True | True | Passed |
| repeat validation is pure | True | True | Passed |
| two independent violations in variant | False | False | Passed |
SHA-256 / 48c9c7fe09f33d1f5b31b9b62dc3af923959683f650e8242ce128d6996236799
Verification & scope
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:37:03.637764+00:00.
Case digest / a507612e18fffbfac6abb3acc213066a7a5e3edb133dcf3c95b249bc5d18f5a6