Skip to content

fix(core): the revoke sweep names the registration it means, not whatever holds the offset - #492

Merged
astubbs merged 3 commits into
masterfrom
fix/conditional-by-key-shard-removal
Sep 9, 2026
Merged

astubbs merged 3 commits into
masterfrom
fix/conditional-by-key-shard-removal

Conversation

@astubbs

@astubbs astubbs commented Sep 9, 2026 •

Copy link
Copy Markdown
Owner

Owning note: docs/inflight/bug-shard-displacement-orphans-the-retry-queue-entry.md - the note #483's defect-class sweep reported this removal on. Its state and its PROPOSED closed marker are untouched; only a dated outcome section is added.

Description

The last unconditional by-key shard removal in main code becomes conditional.
ShardManager.removeWorkFromShardFor - the revoke and lost sweep - removed with
removeWorkAtOffset(consumerRecord.offset()), so whatever occupied the offset when the removal
landed is what left the shard. It now removes only the container the revoked record was registered
as, or a stale one, and never a live container that a later registration put at the offset.

This is #468's defect class at its second site.
#468 fixed the poller's stale sweep by making WorkContainer equality
identity, so Map.remove(key, value) became a true compare-and-remove;
#483's sweep reported this site as "the one remaining unconditional by-key
shard removal" and deliberately left it, saying it belonged with #468's
mechanism rather than with a reachability proof. This is that.

The caller holds no container, so the fix is shaped differently

The stale sweep inspects a container and wants that one gone. This sweep is handed the
ConsumerRecords one generation was still carrying as incomplete, so what it can name is the
registration. That works because PartitionState.maybeRegisterNewPollBatchAsWork puts the same
record instance into incompleteOffsets and into the WorkContainer it builds, in one loop body -
so "the container this generation registered" is a reference comparison, and a later generation's
re-delivery of the offset is a different object because it came from a different fetch.

Two legs, and the middle option was wrong

  • Registration identity alone declines for a container from another registration that is
    nonetheless STALE. Declining there would leave a stale container holding an offset of a revoked
    partition - the state this sweep exists to prevent.
  • Staleness alone declines for a live container built from the record the sweep was handed,
    which is what a revoke sweep driven before the partition's state has been swapped sees - and what
    ShardManagerLincheckTest.revokeSweep models on every invocation. That harness's removal
    operation would have become a silent no-op while staying green.

So the guard declines for exactly one thing, a live container from another registration, which
makes this a no-op everywhere except in the defect case. Each leg has its own ablation arm.

Reachability, stated plainly: not today

Both rebalance callbacks run on the broker-poll thread, so the revoke sweep for a generation
completes before the assignment that could register a replacement at one of its offsets even begins

The harm if it is ever reached is the lost record, not a misdirection: the live container is
gone from the shard while its own PartitionState carries the offset as incomplete, so nothing
selects it and the commit high-water mark cannot pass it until the partition is re-polled.

Evidence

ShardRevokeSweepReplacementEvictionTest, three arms, driven through the production entry point
(PartitionState.onPartitionsRemoved) so the sweep's argument list is built by production code.
Red-then-green, plus an ablation matrix predicted before running - one arm red per ablation, and no
other:

Ablation Arm that goes red
whole guard removed (back to removeWorkAtOffset(offset)) the displacement arm
staleness leg removed aStaleContainerFromAnotherRegistrationIsStillSweptOut
registration-identity leg removed aRevokeSweepStillRemovesALiveContainerBuiltFromTheRecordItWasHanded
decline retires the occupant instead of leaving it the new accounting assertion in the displacement arm

The third ablation reported no failure on its first run and that was a silent false negative:
dropping the leg left an unused local, the module never recompiled, and the grep for failures
matched nothing. Re-run with the build outcome asserted, it is red as predicted. Recorded rather
than tidied away.

Green: bin/ci-unit-test.sh (whole reactor), bin/lincheck-test.sh (ShardManagerLincheckTest included - the harness whose revokeSweep the staleness-only guard would have silently neutered), bin/check-all.sh, and ten repeat runs of the shard suite for order dependence.

Same-defect-class sweep

Every by-key or by-offset removal and lookup in ShardManager, ProcessingShard and
PartitionState was walked. Fixed: the one above. Dismissed, with reasons:

The one red check is inherited, not this PR's

deps: whole-tree CVE scan fails on CVE-2026-59296 in io.micrometer:micrometer-core:1.13.15 - a
whole-tree dependency advisory published after master's last green run. This diff touches no pom and
no dependency, and #480 fails the same check on the same component.
#489 has since split that scan out of scan: repo into its own
non-required check for exactly this situation, so it no longer gates anything. Not fixed here: a
dependency bump belongs on its own branch cut from master.

Everything else is green: Unit, Integration (+heavy), Chaos 1-4, Lincheck, Performance, PIT,
static: analysis, repo: hygiene, shell: macos, CodeQL, codecov patch, claude-review,
review: human LGTM.

Checklist

  • Docs updated - the class's owner (docs/solutions/logic-errors/a-by-key-removal-cannot-say-which-container-it-meant-2026-09-07.md) gains a dated section for the second site; docs/inflight/bug-shard-displacement-orphans-the-retry-queue-entry.md records the outcome beside test(core): the shard-displacement retry-queue orphan is unreachable, and the guard is outside the class #483's report
  • User-facing feature documentation data added under docs/features/ - N/A - an internal correctness fix with no user-facing option, behaviour or default to describe
  • Tests added/updated - ShardRevokeSweepReplacementEvictionTest, three arms with an ablation matrix
  • docs/inflight/ working note (pr-/branch-) started at the PR's first commit - N/A - one commit, no cross-branch state a future branch inherits and nothing here gh cannot show; the durable knowledge went to the two documents above at the same commit
  • Title & body reflect the final content of this PR
  • Ran ce-simplify and ce-code-review locally - ce-simplify N/A (one new method plus its tests; nothing to consolidate). A scoped local code review WAS run and its findings are all fixed in 0381d0daa - the cleared suspicion naming the wrong object, the half-true thread claim, an unstated exception to the registration-identity premise, an inert @SuppressWarnings, a stale javadoc, and two test gaps including a decline branch with no accounting assertion (proven by ablation). The owner also requested @claude review this; that review is clean, and my reply records that it ran against the previous head.

🤖 Generated with Claude Code

https://claude.ai/code/session_01Xoi3HYae8pjsEatuNFKieD

…ever holds the offset

WHAT CHANGED FOR A USER: a rebalance can no longer remove a record that a later assignment had
already re-registered. `ShardManager.removeWorkFromShardFor` - the revoke and lost path - removed by
KEY with `removeWorkAtOffset(consumerRecord.offset())`, so whatever occupied the offset when the
removal landed is what left the shard. It now removes only the container the revoked record was
registered as, or a stale one, and never a live container from a later registration.

THE DEFECT CLASS, AND WHERE THIS SITE SITS IN IT. #468 closed the same
class at the poller's stale sweep: a removal keyed on the offset cannot say WHICH of two containers
at that offset it meant, and the fix was to make `WorkContainer` equality identity so that
`Map.remove(key, value)` is a true compare-and-remove. #483's defect-class
sweep reported this site alongside it - "the one remaining unconditional by-key shard removal" - and
deliberately left it, on the grounds that it belonged with #468's mechanism rather than with
a reachability proof. This is that.

THE CALLER HOLDS NO CONTAINER, WHICH IS WHY THE FIX IS SHAPED DIFFERENTLY. The stale sweep inspects
a container and then wants that one gone. This sweep is handed the `ConsumerRecord`s ONE generation
was still carrying as incomplete (`PartitionState.onPartitionsRemoved` ->
`removeAnyShardEntriesReferencedFrom`), so what it can name is the REGISTRATION, not an inspected
container. That works because `PartitionState.maybeRegisterNewPollBatchAsWork` puts the same record
instance into `incompleteOffsets` and into the `WorkContainer` it builds, in one loop body - so
"the container this generation registered for this record" is `occupant.getCr() == revokedRecord`,
and a later generation's re-delivery of the same offset is a different object because it came from
a different fetch.

TWO LEGS, AND THE MIDDLE OPTION WAS WRONG. Registration identity alone declines for a container from
another registration that is nonetheless STALE, and declining there would leave a stale container
holding an offset of a revoked partition - the state this sweep exists to prevent. Staleness alone
declines for a live container built from the record the sweep was handed, which is what a revoke
sweep driven before the partition's state has been swapped sees, and is what
`ShardManagerLincheckTest.revokeSweep` models on every invocation: that harness's removal operation
would have become a silent no-op while staying green. So the guard declines for exactly one thing -
a LIVE container from another registration - which makes the change a no-op everywhere except in the
defect case, and each leg has its own ablation arm.

THE HARM, IF IT IS REACHED, IS THE LOST RECORD - #468's, not a misdirection. The live
container is gone from the shard while its own `PartitionState` carries the offset as incomplete, so
nothing selects it and nothing completes it, and the commit high-water mark cannot pass it until the
partition is re-polled. (The scope note this came from said the cost was "misdirection bounded to
one control-loop tick by #481's purge"; that is the DISPLACEMENT orphan's cost, not this
one's. This removal is a shard removal, and the purge collects retry-queue entries.)

REACHABILITY, STATED PLAINLY: NOT TODAY. Both rebalance callbacks run on the broker-poll thread, so
the revoke sweep for a generation completes before the assignment that could register a replacement
at one of its offsets even begins - the same single-poll-thread argument #483 used for
`addWorkContainer`'s displacement branch. That is an argument about callers and nothing checks it.
The conditional form costs one reference comparison and does not rest on it, which is exactly why
#468 wrote `getWorkIfAvailable`'s last-resort sweep conditionally against a race it had shown
was unreachable there. The arms therefore drive the sweep directly rather than through a rebalance,
and each says so.

THE NPE GUARDS ON THIS PATH ARE UNTOUCHED, AND A THIRD SUSPICION IS CLEARED AT THE SITE. The
single-read `getShard` idiom #345 put in stays exactly where it was - the new work is inside
`ProcessingShard`, below it. The staleness question the second leg asks reaches
`PartitionStateManager.getPartitionState`, which answers null for a partition never assigned; the
discriminator is that `resetOffsetMapAndRemoveWork` installs `RemovedPartitionState` for the
partition BEFORE calling in, and every record here came out of that state's own tracking. Dated and
recorded on the method, with what would reopen it.

EVIDENCE. `ShardRevokeSweepReplacementEvictionTest`, three arms, RED-then-GREEN through the
production entry point (`PartitionState.onPartitionsRemoved`) so the argument list is built by
production code and nothing is hand-picked. Ablation matrix, predicted before running, one arm red
per ablation and no other:

- whole guard removed, back to `removeWorkAtOffset(offset)` -> the displacement arm alone
- staleness leg removed -> `aStaleContainerFromAnotherRegistrationIsStillSweptOut` alone
- registration-identity leg removed -> `aRevokeSweepStillRemovesALiveContainerBuiltFromTheRecordItWasHanded` alone

The third ablation reported NO failure on its first run and that was a silent false negative, not a
result: dropping the leg left an unused local under `-Xlint`, the module never recompiled, and the
grep for failures matched nothing. Re-run with the build outcome asserted, it is red as predicted.
Recorded rather than tidied away, because a green ablation is exactly the shape that gets believed.

Reproduce: `bin/ci-unit-test.sh`; `bin/lincheck-test.sh`; `bin/check-all.sh`;
`./mvnw -pl parallel-consumer-core -am test -Dtest='ShardRevokeSweepReplacementEvictionTest,ShardStaleSweepReplacementEvictionTest,ShardManagerRevokeSweepNpeTest,ShardDisplacementOrphanReachabilityTest'`.

RECORDS. The class's owner - `docs/solutions/logic-errors/a-by-key-removal-cannot-say-which-container-it-meant-2026-09-07.md`
- gains a dated section for the second site, because its own sweep's only hit was a false positive
and the real one was found a day later.
`docs/inflight/bug-shard-displacement-orphans-the-retry-queue-entry.md` records the outcome beside
#483's report; its state and its PROPOSED close are untouched, both being the owner's.

Co-Authored-By: Claude Opus <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Xoi3HYae8pjsEatuNFKieD
@github-actions

github-actions Bot commented Sep 9, 2026

Copy link
Copy Markdown

Dependency Review

✅ No vulnerabilities or license issues or OpenSSF Scorecard issues found.

Scanned Files

None

@github-actions

github-actions Bot commented Sep 9, 2026 •

Copy link
Copy Markdown

✅ Duplicate Code Report

Two engines run in parallel for cross-validation. Each has its own thresholds tuned to its baseline - the real safety net is the per-engine "max increase vs base" check.

✅ PMD CPD

PR Base Change
Clones 27 27 ➖ 0
Duplicated lines 949 949 ➖ 0
Duplication 0.37% 0.37% ➖ 0
Rule Limit Status
Max duplication 0.5% ✅ Pass (0.37%)
Max increase vs base +0.1% ✅ Pass (+0.00%)

No new clones introduced by this PR.

✅ jscpd (language-agnostic)

PR Base Change
Clones 106 106 ➖ 0
Duplicated lines 1503 1503 ➖ 0
Duplication 0.83% 0.84% 🙂 -0.01%
Rule Limit Status
Max duplication 2% ✅ Pass (0.83%)
Max increase vs base +0.1% ✅ Pass (-0.01%)

No new clones introduced by this PR.

Powered by astubbs/duplicate-code-cross-check

@github-actions

github-actions Bot commented Sep 9, 2026 •

Copy link
Copy Markdown

🟢 Throughput — OK

This branch measured about 5% faster than master, on the one test this measures. That is INSIDE this test's own run-to-run spread of about 17%, so read it as a reading and not as a result - re-running the same commit moves it by about as much.

What Value Meaning
Compared with master 1.046 above 1.00 is faster, below is slower
Subject test took 44.1s the test under measurement
Control tests took 25.0s the other tests in this same run
Subject ÷ controls, here 1.765 not a speed - a shape that cancels machine speed
Subject ÷ controls, master 1.847 median of recent master runs
Reported rate 88967 rec/s this machine only; not comparable across runners

Allowable range 🟢 ≥ 0.70 · 🟡 0.50–0.70 (about a 30% loss) · 🔴 < 0.50 (about a 50% loss)

What the numbers mean, and what they cannot tell you

The one that gets misread. Subject ÷ controls is a shape, not a speed. 1.7 means the subject took 1.7 times as long as the control tests in the same run - it says nothing about master on its own, and a reviewer has already read it as "1.7x faster than master". Only Compared with master answers that question.

Why a shape and not a rate. A rate depends on which runner you drew. A shape does not: every test here processes a fixed number of records, so a runner twice as slow doubles the subject and the controls together and leaves their ratio alone. That is the whole trick, and it is why the reported rate is shown last and labelled as this machine only.

Reading the comparison. master ÷ this run. Above 1.00 the subject is proportionally quicker here than on master; below 1.00 it is slower. 0.50 means it takes twice as long relative to its controls - that is the failing bound, not a small one.

By conservation, not by correction. Every test in this lane processes a fixed number of records, so within one run the ratio of one test's time to another's is invariant under machine speed — a runner twice as slow doubles both terms and leaves the ratio alone. There is no machine-index correction to be wrong, because nothing needed correcting. share = subjectSeconds / controlSeconds, both from this same run.

Per-method times, not class times. A class time is work + setup, and container startup and @BeforeAll do not scale with work — they are the non-conserved term, and leaving them in breaks the invariant.

Reference is the median of 9 recent perf baseline (master) run(s), read from their artifacts. There is no committed baseline to go stale, and a share is dimensionless, so an old entry stays comparable to a new one without re-baselining. Shares observed: 1.780 – 2.032.

What this still cannot do. It removes machine-to-machine variance. It does not remove this test's own run-to-run variance, measured at about 30% on a single unchanged commit while its controls stayed within 5%. That is a property of the test, not of the comparison, and no arithmetic here can touch it — which is why the reference is a median and the bounds are deliberately coarse. 🟡 means look at this; only 🔴 is outside the measured spread.

Runs used: 51d9bb2, 0ca787c, b654cb2, a055248, 65e11e3, e05b399, ee3d8a9, 890aa53, fd3f913

Since the previous push: ratio 0.9 -> 1.046, share 2.047 -> 1.765, rate 81173 -> 88967 (+9.6%). One push of difference sits inside this test's measured spread - read it as movement, not as a result.

Updated for 0381d0d · run 34298550025 · 2026-09-09 01:22 UTC

@codecov

codecov Bot commented Sep 9, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 83.04%. Comparing base (b654cb2) to head (0381d0d).
⚠️ Report is 3 commits behind head on master.
✅ All tests successful. No failed tests found.

Additional details and impacted files
@@             Coverage Diff              @@
##             master     #492      +/-   ##
============================================
+ Coverage     82.50%   83.04%   +0.54%     
- Complexity     1563     1588      +25     
============================================
  Files            96       96              
  Lines          5383     5431      +48     
  Branches        537      547      +10     
============================================
+ Hits           4441     4510      +69     
+ Misses          747      726      -21     
  Partials        195      195              
Flag Coverage Δ
chaos 53.14% <54.54%> (?)
integration 61.63% <54.54%> (+0.10%) ⬆️
performance 47.65% <0.00%> (?)
unit 79.30% <100.00%> (+0.12%) ⬆️

Flags with carried forward coverage won't be shown. Click here to find out more.

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@github-actions

github-actions Bot commented Sep 9, 2026 •

Copy link
Copy Markdown

⚠️ SpotBugs Report

354 bug(s) found (rule-level exclusions only - see docs/inflight/static-spotbugs-rule-registry.md). See the annotations on the Files Changed tab for details.

Updated for 0381d0d · run 34298550025 · 2026-09-09 01:24 UTC

@astubbs astubbs left a comment

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@astubbs

astubbs commented Sep 9, 2026

Copy link
Copy Markdown
Owner Author

@claude review this

@claude

claude Bot commented Sep 9, 2026 •

Copy link
Copy Markdown

Claude finished @astubbs's task in 2m 30s —— View job


Review

  • Read PR body, comments, and prior review
  • Read changed files (ProcessingShard.java, ShardManager.java, test, docs)
  • Review diff for correctness, style, edge cases
  • Check ablation claims against the test
  • Post final review

Summary: looks correct. I traced the guard logic by hand against all three test arms and the production call graph; nothing new to add on top of @astubbs's LGTM.

Correctness

ProcessingShard.removeWorkForRevokedRecord (parallel-consumer-core/src/main/java/bz/stub/parallelconsumer/state/ProcessingShard.java:394-410) declines only when !isTheRegistrationBeingRevoked && !isWorkContainerStale(occupant) — i.e. it evicts the occupant unless it's both (a) from a different registration and (b) still live. I hand-traced all three ShardRevokeSweepReplacementEvictionTest arms against this:

  • theRevokeSweepMustNotEvictAFreshContainerThatDisplacedTheStaleOne: occupant is the fresh, live, different-registration container → declines → matches assertion.
  • aRevokeSweepStillRemovesALiveContainerBuiltFromTheRecordItWasHanded: occupant.getCr() == revokedRecord is true (same instance, since it's planted directly from revokedRecord) → short-circuits past the staleness check → always evicted → matches assertion.
  • aStaleContainerFromAnotherRegistrationIsStillSweptOut: different registration but stale → evicted → matches assertion.

The evictIfStillResident compare-and-remove (workMap.remove(offset, inspected)) is sound given WorkContainer equality is reference identity (no equals/hashCode override, confirmed by inspection) — so it's a true CAS, not a coordinate match.

I also confirmed PartitionState.maybeRegisterNewPollBatchAsWork (PartitionState.java:716-717) puts the same ConsumerRecord instance into both incompleteOffsets and the WorkContainer it builds in one loop body, which is the premise the registration-identity comparison depends on.

The Lincheck-harness claim checks out

I verified ShardManagerLincheckTest.revokeSweep/addWork build sweepArguments from the same records[] instances used by addWork (ShardManagerLincheckTest.java:160-163,196-198,221-223), with a fixed EPOCH throughout — so a staleness-only guard genuinely would have made every revokeSweep operation in that harness a silent no-op, as the PR body claims. That's a real regression this guard's registration-identity leg prevents.

No remaining production callers of removeWorkAtOffset

Grepped the tree — every remaining call site is in a test; the javadoc's claim that it's now "the unconditional primitive four test classes use" is accurate.

Docs

docs/solutions/logic-errors/a-by-key-removal-cannot-say-which-container-it-meant-2026-09-07.md and docs/inflight/bug-shard-displacement-orphans-the-retry-queue-entry.md are both updated consistently — the inflight note's PROPOSED closed marker is untouched as claimed, and the new dated sections correctly scope this as a distinct fix from the (still-open, still-unreachable) displacement-branch orphan.

Nits (non-blocking)

None worth raising — I didn't find a correctness, concurrency, or style issue in the diff. The @SuppressWarnings("ReferenceEquality") on the one intentional identity comparison is appropriately scoped and commented.

astubbs and others added 2 commits September 9, 2026 13:13
…branch is pinned by its accounting

Review pass on the commit before this one. No behaviour changes except the removal of a suppression
that suppressed nothing; the rest corrects claims that were wrong or incomplete, and closes two gaps
where nothing would have gone red.

THE ARGUMENT WAS ABOUT THE WRONG OBJECT. The cleared suspicion on removeWorkForRevokedRecord
justified the staleness call's null-safety with a fact about `revokedRecord` - that it came out of
the removed state's own tracking. But isWorkContainerStale resolves the state from the OCCUPANT
(pm.getPartitionState(workContainer)), and in the only branch that asks, the occupant is by
definition a DIFFERENT registration. What actually closes it is that every ShardKey variant is
partition-scoped - KeyOrderedKey holds a TopicPartition, TopicPartitionKey is one - so an occupant of
this shard necessarily belongs to the revoked partition. Both legs are now stated, and the reopener
list gains the one that follows: A TOPIC-SCOPED SHARD KEY, which would admit an occupant from a
partition that may never have been assigned, and which reopens addWorkContainer's cleared suspicion
at the same time. That entry already named it; this one did not.

A STRONGER UNREACHABILITY ARGUMENT WAS SITTING IN THE SAME PARAGRAPH, UNUSED. resetOffsetMapAndRemoveWork
installs RemovedPartitionState BEFORE it calls in, so checkIfWorkIsStale is true for EVERY occupant and
the staleness leg carries all of them: on the shipped path this method is behaviourally identical to the
by-key removal it replaced, and the decline branch is defence for a caller that sweeps before the swap -
which is exactly what ShardManagerLincheckTest.revokeSweep does. That is structural and says nothing
about threads, so it is strictly stronger than the thread-ordering argument, and the previous commit used
the same fact only as the NPE discriminator without noticing what else it decided. Now stated on the
method.

THE THREAD CLAIM WAS HALF TRUE, IN THREE PLACES. "Both rebalance callbacks run on the broker-poll
thread" is contradicted by the javadoc on the method that calls into this sweep: the close route
(maybeCloseConsumer -> onLeavePrepare -> onPartitionsRevoked) runs it on the CONTROL thread, and
resetOffsetMapAndRemoveWork calls that live in every commit mode. The conclusion survives - on that
route the sweep and addWorkContainer are on the same thread - but the stated reason was wrong for one
of the two production routes. Corrected in the test javadoc and the solutions write-up; the commit
body that also carries it is squashed away, and the squash message is written against this wording.

THE REGISTRATION-IDENTITY CLAIM HAD A NAMED EXCEPTION IT DID NOT NAME. addNewIncompleteRecord puts
unconditionally while addWorkContainer DROPS an arrival whose resident is not stale, so an
in-generation replay of an already-registered offset leaves the partition naming record B while the
shard holds the container over record A - and the comparison is then false about a container that IS
the right one. Harmless, because the staleness leg covers it, and it is the same reopener
addWorkContainer's cleared suspicion already lists. Said at the site rather than left for the next
reader to find.

TWO TEST GAPS, ONE OF WHICH A MUTATION WALKED STRAIGHT THROUGH.

- The decline branch is a new exit that must retire nothing and release nothing, and NOTHING went red
  if it did. Verified by ablation rather than asserted: making the decline call
  population.onRetired() + excludeFromSelection(occupant) left the offset assertion GREEN and nothing
  else noticed - a permanent RecordPopulation deficit, which is this class's recurring defect and has
  no clamp and nothing to reconcile it. The arm now samples population, the shard available-counter
  sum and the container's own claim before the sweep and asserts all three unchanged, and that
  mutation is red.
- aStaleContainerFromAnotherRegistrationIsStillSweptOut asserted only that the occupant came from
  another registration, never that it was STALE - the property the arm exists to pin - leaning on the
  fixture's epoch instead. It now asserts staleness directly, the way its sibling arm asserts
  liveness, and its closing assertion gains the message it was the only one in the file to lack.

REMOVED: @SuppressWarnings("ReferenceEquality"). Measured by deleting it and asserting BUILD SUCCESS -
neither WorkContainer nor Kafka's ConsumerRecord overrides equals, so Error Prone never fired. It
suppressed nothing and disagreed with isResident's bare == two methods away; a suppression that
suppresses nothing teaches the next reader something untrue. The reason it is absent is now written
where it was.

removeWorkAtOffset's javadoc still opened "From the onPartitionsRemoved callback", the caller the
previous commit moved. It now says what it is - a test-only modelling primitive with no production
caller - and names which of the two conditional removals a main-code site should reach for instead.

The docs/inflight/ section is cut to a pointer. It restated the mechanism, which the solutions
write-up owns, and that directory's rule is to track what is OPEN: the item is closed, so what
belongs there is only that a reader arriving at #483's four-item list can see which one moved.

ABLATION MATRIX, re-run against the strengthened arms - one arm red per ablation, no other:
whole guard removed -> the displacement arm; staleness leg removed -> the stale-other-registration
arm; identity leg removed -> the live-same-registration arm; decline retires the occupant -> the new
accounting assertion.

Reproduce: `bin/ci-unit-test.sh`; `bin/lincheck-test.sh`; `bin/check-all.sh`.

Co-Authored-By: Claude Opus <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Xoi3HYae8pjsEatuNFKieD
@astubbs

astubbs commented Sep 9, 2026

Copy link
Copy Markdown
Owner Author

Responding to the automated review, the two duplication engines and the SpotBugs report - and leading with the one thing that qualifies all of it.

The review ran against the previous head

It completed at 01:10:18Z against 381df6fb7; I pushed 0381d0daa at 01:14:02Z. So it reviewed commit 1 only, and two of its observations are now stale in my favour rather than against it:

  • "The @SuppressWarnings("ReferenceEquality") on the one intentional identity comparison is appropriately scoped and commented." That annotation is gone as of 0381d0daa. I measured it rather than reasoned about it - deleted it, rebuilt under -Pci, BUILD SUCCESS. Neither WorkContainer nor Kafka's ConsumerRecord overrides equals, so Error Prone never fired: it suppressed nothing, and it disagreed with isResident's bare == two methods away. A suppression that suppresses nothing teaches the next reader something untrue, so what is there now is the comment explaining why there is no annotation.
  • "the javadoc's claim that it's now 'the unconditional primitive four test classes use' is accurate." Still true, but that javadoc was rewritten in 0381d0daa - it now leads with what the method is (a test-only modelling primitive with no production caller) rather than opening "From the onPartitionsRemoved callback", which named the caller commit 1 had just moved away.

0381d0daa also merges master (#480, #488, #489). A re-review is the owner's call, not mine - I am not requesting one; flagging it so nobody reads the clean verdict as covering the current head.

Point by point, including what I verified rather than took

  • The guard's decline condition, traced against all three arms. Agreed, and independently ablated rather than only traced: each leg has its own arm, and each ablation reds exactly one. Whole guard removed → the displacement arm; staleness leg removed → aStaleContainerFromAnotherRegistrationIsStillSweptOut; registration-identity leg removed → aRevokeSweepStillRemovesALiveContainerBuiltFromTheRecordItWasHanded. Re-run on the merged tree; the red arm is still red there.
  • evictIfStillResident is a true compare-and-remove given identity equality. Agreed. That premise has its own standing tripwire on master - ShardStaleSweepReplacementEvictionTest.twoContainersAtOneOffsetMustNotBeInterchangeable, from fix(core)!: WorkContainer equality is identity, so the stale sweep removes only the container it inspected #468 - so this PR did not need to re-assert it.
  • maybeRegisterNewPollBatchAsWork puts the same ConsumerRecord instance into both. Agreed, and I have since written the exception into the javadoc, because the unqualified claim was mine and it is not quite true: addNewIncompleteRecord puts unconditionally while addWorkContainer drops an arrival whose resident is not stale, so an in-generation replay of an already-registered offset (a seek, an offset-reset, a truncation replay) leaves the partition naming record B while the shard holds the container over record A. Harmless - the staleness leg covers it - and it is the same reopener addWorkContainer's own cleared suspicion already lists.
  • The Lincheck-harness claim. Agreed, and I checked the same thing from the other end: incrementPartitionAssignmentEpoch starts from KAFKA_OFFSET_ABSENCE = -1, so the first assignment yields epoch 0 and the harness's EPOCH = 0L containers are live. A staleness-only guard really would have made every revokeSweep in it a silent no-op.
  • No remaining production callers of removeWorkAtOffset. Agreed - same grep, same answer.
  • Docs consistency, and the PROPOSED closed marker untouched. Agreed. 0381d0daa cuts that inflight section further, to a pointer: it restated the mechanism, which the solutions write-up owns, and docs/inflight/AGENTS.md says to track what is OPEN. What is left is only that a reader arriving at test(core): the shard-displacement retry-queue orphan is unreachable, and the guard is outside the class #483's four-item list can see which one moved.
  • "Nits: none worth raising." Noted. The gaps that did exist were found by a separately-scoped local review and are fixed in 0381d0daa: the cleared suspicion justified the null-safety with a fact about revokedRecord when isWorkContainerStale resolves the state from the occupant (the real discriminator is that every ShardKey variant is partition-scoped, and the reopener list gains "a topic-scoped shard key"); the "both rebalance callbacks run on the broker-poll thread" claim was half true, because the close route runs the sweep on the control thread; and the decline branch had no accounting assertion - proven by ablation, since making it population.onRetired() + excludeFromSelection(occupant) left every existing assertion green while creating a permanent RecordPopulation deficit. That arm now samples population, the shard available-counter sum and the container's own claim, and that mutation is red.

Duplicate Code Report

Read, both engines, refreshed for 0381d0da: PMD CPD 27→27 clones (0.37%, +0.00%) and jscpd 106→106 (0.83%, -0.01%), "No new clones introduced by this PR" from each. Nothing for me to remove. There is no separate file-similarity comment on this PR.

SpotBugs Report

Read rather than counted from the tick. The total is tree-wide and pre-existing; what matters is that bin/check-pr-analysis-surfaces.sh 492 reports none on lines this PR wrote and none inherited in the files it touches, at head 0381d0daa. I also ran spotbugs:spotbugs spotbugs:check plus Error Prone/NullAway over parallel-consumer-core locally and got the same answer on my files.

Checks

Green: Unit, Integration (+heavy), Chaos 1-4, Lincheck, Performance, PIT, static: analysis, repo: hygiene, shell: macos, CodeQL, codecov patch. The one red is deps: whole-tree CVE scan on CVE-2026-59296 in io.micrometer:micrometer-core:1.13.15 - whole-tree, this diff touches no pom, and #489 has already made that check non-required for exactly this reason. Not fixed here; a dependency bump belongs on its own branch from master.

@astubbs
astubbs marked this pull request as ready for review September 9, 2026 01:30
@astubbs
astubbs merged commit 9c67c89 into master Sep 9, 2026
36 of 37 checks passed
@astubbs
astubbs deleted the fix/conditional-by-key-shard-removal branch September 9, 2026 01:36
astubbs added a commit that referenced this pull request Sep 9, 2026
…y-key removal are on master

The intake-stall measurement landed with the retry-cadence arm and the
interim latch warning as what remains, and its box is ticked. The last
by-key shard removal the #483 sweep reported is fixed
on master, red first with an ablation arm per leg of the guard. The
known unknown that named it is struck through with the outcome, and the
retry queue's same-shaped removal is recorded as that queue's keying
model rather than a remaining instance.

Claude-Session: 460f7df9-dcc2-4b00-a9f9-62f3a2c6d5e4
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
astubbs added a commit that referenced this pull request Sep 9, 2026
… named, and five stale lines are corrected

Owner's decisions, 2026-09-09: the merge queue is closed as of today,
with later finds 0.6.0.x unless data loss on a default configuration;
the poisoned-transaction wedge is the second named exception beside
#44, and the release-note draft now carries it; the gate-latch
warning #487 argued for is v6-sized and joins tier 1 as the last
item; the upstream flat-counter reporters are not asked.

Corrections from the owner's read of the note: the release page body
is posted by hand on the day with gh release edit, because release.yml's
exact heading match misses the unreleased heading on master - so
#199 follows the tag rather than gating it, and the two lines that
said the workflow already publishes the curated section are fixed; the
#468 line no longer asks the reader to check a PR body for two
by-key removals that #468 dismissed and #492 fixed; the
vetting sweep's opening claim that the quarantine registry is non-empty
is struck as the sweep's dated reading; and the disposition list names
the three deferred bug notes it omitted.

Claude-Session: 460f7df9-dcc2-4b00-a9f9-62f3a2c6d5e4
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
astubbs added a commit that referenced this pull request Sep 9, 2026
This branch is docs-only - two markdown files differ from master - so neither
red it drew can be its own. Both are recorded against the notes that own them,
and reading each occurrence turned up a correction the note needed.

INTEGRATION TESTS, on this PR's own head. The shape in
ci-broker-container-exit-126-is-undiagnosable.md, exactly: one class fell slowly
at the container-start timeout and every other broker class fell in
milliseconds with NoClassDefFoundError on BrokerIntegrationTest. Codecov renders
that as "20 Tests Failed"; it is one failure.

What is new is that the cause was in the log all along. Testcontainers prints
the failed container's own output at GenericContainer#tryStart, one line below
the "Wait strategy failed" line the note's signature block quotes, and it reads
"sh: /tmp/testcontainers_start.sh: Text file busy" - ETXTBSY, exec refused
because the starter script was still open for writing. The container command
waits for that script to EXIST and then executes it, so a file the daemon has
created but not finished extracting is executable-shaped and not executable, and
the shell reports the refusal as exit 126. A Testcontainers start race, widened
by a busy runner; nothing in the product, the image or the Kafka configuration.

Refetching #347's 2026-08-25 job, the run this note was written from,
shows the identical two lines. So the note's premise under item 1 - "the
container's stdout is nowhere in the job log" - was false of its own founding
evidence. The instrument was fine; the triage stopped one line short. That
correction is proposed in the note's vetting marker rather than applied, because
the note's impact is misdirection and those are the owner's to close.

CHAOS PAIN SUITE 4/4, on #495 - also docs-only, one
README paragraph. ChaosChurnStormIT NO_PROGRESS at 96632/100000 for 30s against
a 30s bound, seed 3717713223451201639. It goes in
test-no-progress-window-may-not-transfer-to-w1.md as one appended row, with the
part that makes it worth having: the fleet KEPT CONSUMING, reaching 99569 by the
settle summary, so the outstanding count fell from 3368 to 431 - inside the
TAIL_SLACK of 500. That is the "drains" branch of the deciding experiment the
note states. It is the weak form and the row says so: no recovery diagnostic, so
the counter compared is the ledger's rather than the probe's, and the conductor's
churn ended 10s after the firing, so it is recovery-once-churn-stops.

The bigger finding is that the deciding experiment had already been answered
twice and this note never took delivery. test-857-churn-storm-async-stalls.md
drained six for six on seed 9086872209853284830 with the diagnostic engaged, and
its 2026-09-08 sighting drained seed 5650361238717170909 from 93487 to
101070/100000 at an outstanding count of 6513 - larger than every row in the
table. That sighting says outright that this note owns the question; the pointer
was written and nobody followed it. Proposed in the vetting marker for the same
reason as above.

RULED OUT, with a control arm rather than an argument. The six merges that
landed on master today - #480, #487, #488, #491,
#492 and #493 - are the obvious suspects for a chaos red, and
#491 does touch ProgressProbe.java. Its diff does not touch the
NO_PROGRESS path at all - it adds the UNCOMMITTED_COMPLETIONS detector, edits
javadoc, and refactors the finding sink - and the same test PASSED on two heads
that carry every one of those merges, four and six minutes either side of the
failing run. A deterministic regression is excluded; a rate change is not, and
one failure could not establish one.

Nothing quarantined. The container fault has no test to quarantine and the
exit-126 note says a re-run is the correct response there. The chaos firing has
no rate that rule 1 would accept, and docs/quarantined-tests.md is empty - which
is the state to preserve.

Co-Authored-By: Claude Opus <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Xoi3HYae8pjsEatuNFKieD
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant