FAILURE MAP
← Case archive

FA-1556 / Replication / Open access

Complete a streaming checkpoint barrier: Barriers from different checkpoint rounds are combined · case 01

The checkpoint barrier operation is admitted even though barriers from different checkpoint rounds are combined.

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

ROOT CAUSE

The admission path omits the barrier id invariant while validating the other operation preconditions.

VERIFIED REPAIR

Require len(set(r['barrier_id'])) == 1 together with every other stated precondition before accepting the operation.

Unsuccessful approach: Adding the barrier id check repairs the reported defect, but replacing the adjacent state flushed 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 (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 fixtureActualExpectedOutcome
valid operationTrueTruePassed
A checkpoint completes before every input channel reaches its barrierFalseFalsePassed
Barriers from different checkpoint rounds are combinedTrueFalseFailed
A checkpoint is acknowledged while state exists only in memoryFalseFalsePassed
Unaligned checkpoint omits messages already in transitFalseFalsePassed
Upstream progress commits before the transactional sink preparesFalseFalsePassed
unrelated tracing metadataTrueTruePassed
repeat validation is pureTrueTruePassed
two independent violations in variantFalseFalsePassed

SHA-256 / b5b569f1d1b829ca0fe12444e59c4406809dcb6fe543b34243b2952e01eb1996

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 (len(set(r['barrier_id'])) == 1) 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 fixtureActualExpectedOutcome
valid operationTrueTruePassed
A checkpoint completes before every input channel reaches its barrierFalseFalsePassed
Barriers from different checkpoint rounds are combinedFalseFalsePassed
A checkpoint is acknowledged while state exists only in memoryTrueFalseFailed
Unaligned checkpoint omits messages already in transitFalseFalsePassed
Upstream progress commits before the transactional sink preparesFalseFalsePassed
unrelated tracing metadataTrueTruePassed
repeat validation is pureTrueTruePassed
two independent violations in variantFalseFalsePassed

SHA-256 / 4623e3eef7d45fd15bcf6d06281e66a9f5b2548d27520dfa66fc574b3dba444a

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 fixtureActualExpectedOutcome
valid operationTrueTruePassed
A checkpoint completes before every input channel reaches its barrierFalseFalsePassed
Barriers from different checkpoint rounds are combinedFalseFalsePassed
A checkpoint is acknowledged while state exists only in memoryFalseFalsePassed
Unaligned checkpoint omits messages already in transitFalseFalsePassed
Upstream progress commits before the transactional sink preparesFalseFalsePassed
unrelated tracing metadataTrueTruePassed
repeat validation is pureTrueTruePassed
two independent violations in variantFalseFalsePassed

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.635653+00:00.

Case digest / 47a46b6868576c701bab3d95034e874eeaa9323cc401eeaef95c27e0dd80d056