FA-866 / Queues / Open access
Assign a partition during consumer rebalance: A partition is assigned to a consumer not subscribed to its topic · case 01
The partition rebalance operation is admitted even though a partition is assigned to a consumer not subscribed to its topic.
ROOT CAUSE
The admission path omits the subscription match invariant while validating the other operation preconditions.
VERIFIED REPAIR
Require r['subscription_match'][0] in r['subscription_match'][1] together with every other stated precondition before accepting the operation.
Unsuccessful approach: Adding the subscription match check repairs the reported defect, but replacing the adjacent capacity limit 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['revocation_complete'] is True) and (r['assignment_epoch'][0] > r['assignment_epoch'][1]) and (r['checkpoint_available'] is not None) 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 | True | False | Failed |
| 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 / 7e10b209168d847978a595d5167b3965d613ed8d9f71cf77a5f3f58f655dbed7
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['assignment_epoch'][0] > r['assignment_epoch'][1]) and (r['checkpoint_available'] is not None) and (r['subscription_match'][0] in r['subscription_match'][1])
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 | True | False | Failed |
| unrelated tracing metadata | True | True | Passed |
| repeat validation is pure | True | True | Passed |
| two independent violations in variant | False | False | Passed |
SHA-256 / 11ab052e02677112a7d4c04e63964216cc0f0018716c1a2a1edd50c6565d96a9
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.559599+00:00.
Case digest / f8b8a4f554e76f4f9e946e6e57941138c174f0f723787edb4b9478cca81d5778