FAILURE MAP
← Case archive

FA-871 / Queues / Open access

Assign a partition during consumer rebalance: Rebalance exceeds the new owner partition capacity · case 01

The partition rebalance operation is admitted even though rebalance exceeds the new owner partition capacity.

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

ROOT CAUSE

The admission path omits the capacity limit invariant while validating the other operation preconditions.

THE FAILURE

The admission path omits the capacity limit invariant while validating the other operation preconditions.

Unsuccessful approach: Adding the capacity limit check repairs the reported defect, but replacing the adjacent revocation complete 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['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 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 capacityTrueFalseFailed
unrelated tracing metadataTrueTruePassed
repeat validation is pureTrueTruePassed
two independent violations in variantFalseFalsePassed

SHA-256 / 11ab052e02677112a7d4c04e63964216cc0f0018716c1a2a1edd50c6565d96a9

2 / The unsuccessful fix

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

HELD IN THE MEMBER ARCHIVE

The verified repair and its recorded checks are member-only.

This mechanism has 9 recorded checks per implementation. The open-access tier publishes the failure and the unsuccessful fix; the repaired source that passes every check, and the observations that prove it, are available to members.

Every case sharing this mechanism uses the same contract and the same repair, so this one record is held back for all of them.

Member access is invitation-based. Sign in with your invited account to inspect the repair.

Sign in to the archive ↗

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

Case digest / 9b2ed9e6b4de365d70bcce25d003a7086b30aa67bd7299340849583b5e6dc36c