FAILURE MAP
← Case archive

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.

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

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 fixtureActualExpectedOutcome
valid operationTrueTruePassed
The new consumer starts while the old owner remains activeTrueFalseFailed
A stale rebalance overwrites a newer assignmentFalseFalsePassed
A new owner starts without the last committed checkpointFalseFalsePassed
A partition is assigned to a consumer not subscribed to its topicFalseFalsePassed
Rebalance exceeds the new owner partition capacityFalseFalsePassed
unrelated tracing metadataTrueTruePassed
repeat validation is pureTrueTruePassed
two independent violations in variantFalseFalsePassed

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 fixtureActualExpectedOutcome
valid operationTrueTruePassed
The new consumer starts while the old owner remains activeFalseFalsePassed
A stale rebalance overwrites a newer assignmentTrueFalseFailed
A new owner starts without the last committed checkpointFalseFalsePassed
A partition is assigned to a consumer not subscribed to its topicFalseFalsePassed
Rebalance exceeds the new owner partition capacityFalseFalsePassed
unrelated tracing metadataTrueTruePassed
repeat validation is pureTrueTruePassed
two independent violations in variantFalseFalsePassed

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 fixtureActualExpectedOutcome
valid operationTrueTruePassed
The new consumer starts while the old owner remains activeFalseFalsePassed
A stale rebalance overwrites a newer assignmentFalseFalsePassed
A new owner starts without the last committed checkpointFalseFalsePassed
A partition is assigned to a consumer not subscribed to its topicFalseFalsePassed
Rebalance exceeds the new owner partition capacityFalseFalsePassed
unrelated tracing metadataTrueTruePassed
repeat validation is pureTrueTruePassed
two independent violations in variantFalseFalsePassed

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