FA-851 / Queues / Open access
Assign a partition during consumer rebalance: The new consumer starts while the old owner remains active · case 01
The partition rebalance operation is admitted even though the new consumer starts while the old owner remains active.
ROOT CAUSE
The admission path omits the revocation complete invariant while validating the other operation preconditions.
VERIFIED REPAIR
Require r['revocation_complete'] is True together with every other stated precondition before accepting the operation.
Unsuccessful approach: Adding the revocation complete check repairs the reported defect, but replacing the adjacent assignment epoch check loses that independent invariant.
Case contract
Return a Boolean admission decision for assign a partition during consumer rebalance. The record r must satisfy all of: r['revocation_complete'] is True; r['assignment_epoch'][0] > r['assignment_epoch'][1]; r['checkpoint_available'] is not None; r['subscription_match'][0] in r['subscription_match'][1]; r['capacity_limit'][0]+r['capacity_limit'][1] <= r['capacity_limit'][2]. Extra tracing fields are ignored; validation does not mutate the record.
Why this case matters
A deterministic local contract for queues. 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 (r['assignment_epoch'][0] > r['assignment_epoch'][1]) and (r['checkpoint_available'] is not None) and (r['subscription_match'][0] in r['subscription_match'][1]) and (r['capacity_limit'][0]+r['capacity_limit'][1] <= r['capacity_limit'][2])
def check(label, actual, expected):
observations.append({"check": label, "actual": actual, "expected": expected, "passed": actual == expected})
r = {'revocation_complete': True, 'assignment_epoch': [5, 4], 'checkpoint_available': 0, 'subscription_match': ['orders', ['orders', 'audit']], 'capacity_limit': [2, 1, 3]}
check('valid operation', solve(r), True)
check('The new consumer starts while the old owner remains active', solve(dict(r, **{'revocation_complete': False})), False)
check('A stale rebalance overwrites a newer assignment', solve(dict(r, **{'assignment_epoch': [4, 4]})), False)
check('A new owner starts without the last committed checkpoint', solve(dict(r, **{'checkpoint_available': None})), False)
check('A partition is assigned to a consumer not subscribed to its topic', solve(dict(r, **{'subscription_match': ['events', ['orders', 'audit']]})), False)
check('Rebalance exceeds the new owner partition capacity', solve(dict(r, **{'capacity_limit': [3, 1, 3]})), False)
check('unrelated tracing metadata', solve(dict(r, trace='run-'+str(N))), True)
check('repeat validation is pure', solve(r), True)
invalid = {'revocation_complete': False, 'assignment_epoch': [4, 4], 'checkpoint_available': None, 'subscription_match': ['events', ['orders', 'audit']], 'capacity_limit': [3, 1, 3]}
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 |
| The new consumer starts while the old owner remains active | True | False | Failed |
| A stale rebalance overwrites a newer assignment | False | False | Passed |
| A new owner starts without the last committed checkpoint | False | False | Passed |
| A partition is assigned to a consumer not subscribed to its topic | False | False | Passed |
| Rebalance exceeds the new owner partition capacity | 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 / 7024e84fc1e6c03d36999f747c0bb663274350bd75c08020a94d0c7dbc824d64
2 / The unsuccessful fix
Exit 1"""Failure Map reference implementation. Python standard library only."""
import json
N = 1
observations = []
def solve(r):
return (r['revocation_complete'] is True) and (r['checkpoint_available'] is not None) and (r['subscription_match'][0] in r['subscription_match'][1]) and (r['capacity_limit'][0]+r['capacity_limit'][1] <= r['capacity_limit'][2])
def check(label, actual, expected):
observations.append({"check": label, "actual": actual, "expected": expected, "passed": actual == expected})
r = {'revocation_complete': True, 'assignment_epoch': [5, 4], 'checkpoint_available': 0, 'subscription_match': ['orders', ['orders', 'audit']], 'capacity_limit': [2, 1, 3]}
check('valid operation', solve(r), True)
check('The new consumer starts while the old owner remains active', solve(dict(r, **{'revocation_complete': False})), False)
check('A stale rebalance overwrites a newer assignment', solve(dict(r, **{'assignment_epoch': [4, 4]})), False)
check('A new owner starts without the last committed checkpoint', solve(dict(r, **{'checkpoint_available': None})), False)
check('A partition is assigned to a consumer not subscribed to its topic', solve(dict(r, **{'subscription_match': ['events', ['orders', 'audit']]})), False)
check('Rebalance exceeds the new owner partition capacity', solve(dict(r, **{'capacity_limit': [3, 1, 3]})), False)
check('unrelated tracing metadata', solve(dict(r, trace='run-'+str(N))), True)
check('repeat validation is pure', solve(r), True)
invalid = {'revocation_complete': False, 'assignment_epoch': [4, 4], 'checkpoint_available': None, 'subscription_match': ['events', ['orders', 'audit']], 'capacity_limit': [3, 1, 3]}
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 |
| The new consumer starts while the old owner remains active | False | False | Passed |
| A stale rebalance overwrites a newer assignment | True | False | Failed |
| A new owner starts without the last committed checkpoint | False | False | Passed |
| A partition is assigned to a consumer not subscribed to its topic | False | False | Passed |
| Rebalance exceeds the new owner partition capacity | 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 / 5b56c679ffa8ab66bc62d7b3824c2350c875e81447643c17e3878e2c64b6b9ce
3 / The verified repair
Exit 0"""Failure Map reference implementation. Python standard library only."""
import json
N = 1
observations = []
def solve(r):
return (r['revocation_complete'] is True) and (r['assignment_epoch'][0] > r['assignment_epoch'][1]) and (r['checkpoint_available'] is not None) and (r['subscription_match'][0] in r['subscription_match'][1]) and (r['capacity_limit'][0]+r['capacity_limit'][1] <= r['capacity_limit'][2])
def check(label, actual, expected):
observations.append({"check": label, "actual": actual, "expected": expected, "passed": actual == expected})
r = {'revocation_complete': True, 'assignment_epoch': [5, 4], 'checkpoint_available': 0, 'subscription_match': ['orders', ['orders', 'audit']], 'capacity_limit': [2, 1, 3]}
check('valid operation', solve(r), True)
check('The new consumer starts while the old owner remains active', solve(dict(r, **{'revocation_complete': False})), False)
check('A stale rebalance overwrites a newer assignment', solve(dict(r, **{'assignment_epoch': [4, 4]})), False)
check('A new owner starts without the last committed checkpoint', solve(dict(r, **{'checkpoint_available': None})), False)
check('A partition is assigned to a consumer not subscribed to its topic', solve(dict(r, **{'subscription_match': ['events', ['orders', 'audit']]})), False)
check('Rebalance exceeds the new owner partition capacity', solve(dict(r, **{'capacity_limit': [3, 1, 3]})), False)
check('unrelated tracing metadata', solve(dict(r, trace='run-'+str(N))), True)
check('repeat validation is pure', solve(r), True)
invalid = {'revocation_complete': False, 'assignment_epoch': [4, 4], 'checkpoint_available': None, 'subscription_match': ['events', ['orders', 'audit']], 'capacity_limit': [3, 1, 3]}
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 |
| The new consumer starts while the old owner remains active | False | False | Passed |
| A stale rebalance overwrites a newer assignment | False | False | Passed |
| A new owner starts without the last committed checkpoint | False | False | Passed |
| A partition is assigned to a consumer not subscribed to its topic | False | False | Passed |
| Rebalance exceeds the new owner partition capacity | 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 / 7d049ff4b1891d50172e768c799b4830d80ca93d772ad441aa2e21045ded88b5
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:36:56.526732+00:00.
Case digest / fe40017ed5c411e0aef2b05f592c642b8da54f493affb60deb741a74cf334ea5