Repository navigation
fix(core) astubbs#423: an InvalidPidMappingException fails the batch, and close never leaks the producer - #429
Conversation
… of marking it succeeded What changed: ParallelEoSStreamProcessor#processAndProduceResults now rethrows after closeOnException in its InvalidPidMappingException arm, wrapped as a PCInternalRuntimeException exactly as the generic arm below it wraps, so the failure travels the route every other produce failure travels. What it changes for a user: a batch that hits that arm is now FAILED and returned to the mailbox, so its offsets can never become commit payload. The instance still closes with the same failure cause and reaches the same CLOSED state, and logs the same "PC closed due to error" line; the only new output is runUserFunction's ordinary "Exception caught in user function running stage, registering WC as failed". Records whose output was never produced are redelivered after restart rather than silently recorded as done. The diagnosis: the arm caught the exception, closed, and fell through to `return results` with a partial - usually empty - result list. runUserFunctionInternal then called onUserFunctionSuccess and addToMailBoxOnUserFunctionSuccess for EVERY WorkContainer in the batch. That is an exactly-once violation of the same shape as the partial-result-set defect #261 fixed: output missing, offset advanced. Why nobody saw it, stated because it is the whole reason this survived review: on master the wrong verdict never reaches a commit, by accident. closeOnException sets failureReason and calls closeDontDrainFirst, which transitions to CLOSING and then blocks in waitForClose(shutdownTimeout + grace) on the control future. Called from a worker thread, that worker is what innerDoClose's awaitTermination waits on, so the whole close runs - mailbox drain, commitOffsetsThatAreReady, the consumer and producer closes - before closeOnException returns to the catch arm. The succeeded containers therefore reach a mailbox nobody drains. The verdict is wrong regardless, and must not depend on the blocking shape of a close that #410 is actively redesigning. Field history: confluentinc#830 reports PC stuck retrying on an InvalidPidMappingException in PERIODIC_TRANSACTIONAL_PRODUCER, closed with no fix; confluentinc#839 wrote this arm. The reporter's own trace is the WRAPPED shape (ExecutionException from FutureRecordMetadata.get), which instanceof does not match, so it takes the generic arm instead - correctly for the batch verdict, incorrectly for liveness. That spin is #411 and is deliberately untouched here. RED then GREEN. The test (ParallelEoSStreamProcessorTest#invalidPidMappingFailsTheBatchInsteadOfMarkingItSucceeded) stubs closeOnException to a no-op so the batch's fate is observable on a live instance. Three assertions were predicted RED before the fix and all three were observed, one run each with the earlier assertion removed so the next could be reached: - produceMessages "Wanted *at least* 2 times ... But was 1 time" - the succeeded record was never re-dispatched; - assertCommits "[the batch whose output was never produced is failed, so nothing is committed] Expecting ArrayList: [1] to contain only: [] but the following element(s) were unexpected: [1]"; - getNumberOfIncompleteOffsets "expected: 1 but was : 0" - the record was marked done. After the rethrow all three hold, and the whole of ParallelEoSStreamProcessorTest and TransactionalClaimCoverageTest are green. Rejected alternatives: - Unwrapping ExecutionException causes into this arm so the wrapped shape also closes. That changes the behaviour #411 tracks - converting a retry spin into a shutdown - and collides with #410, which is replacing that spin with recovery and already pins the current behaviour with a test. - Signalling non-blocking: transitionToClosing then throw, the shape failFatallyOnUnmailboxableRecord uses. Better, but it rewrites closeOnException's contract, which #410 is already redesigning. The minimal fix keeps the blocking close and adds the throw. The register: no new claim. There is no documented sentence for "an invalid producer id fails the batch", and inventing one would mean editing published javadoc inside a data-loss fix. The hole is a second arm of C4 (OFFSET_AND_RECORDS_ATOMIC - offset without records), so the test carries @ProvesClaim(OFFSET_AND_RECORDS_ATOMIC) and C4's note records the arm and the control observed for it. The existing closePCWhenInvalidPidMappingException is unchanged in behaviour and name; its reflection block is extracted into mockProducerManagerInto so the new test does not clone it. Upstream-Issue: confluentinc#830 Upstream-PR: confluentinc#839 Applied-Upstream: no
…throws on close
What changed: ProducerManager#close(Duration) now contains the whole transaction cleanup, releases
the commit lock only if it actually took it, and calls closeProducer(timeout) from a finally. And
AbstractParallelEoSStreamProcessor#innerDoClose's producer step gets the try/catch(WARN) its three
neighbouring shutdown steps already have.
What it changes for a user: a fenced or poisoned producer no longer leaks its KafkaProducer - IO
thread, sockets, buffers - on close(), so a host that restarts PC instances stops leaking one per
fenced shutdown. And close() returns normally with the real failure cause intact, instead of
throwing "Error from poll control thread" over the top of it.
The diagnosis: abortTransaction() throws precisely in the states this library now deliberately
creates - a fenced producer (ProducerFencedException / InvalidProducerEpochException) or a poisoned
one ("we are in an error state"), which #261 makes a terminally failed send produce on
purpose. The abort sat in `try { abortTransaction(); } finally { releaseCommitLock(); }` with
closeProducer(timeout) AFTER the block, so the throw escaped past the one step that must not be
skipped, while doClose's finally still marked the instance CLOSED.
The defect-class sweep, and what it found. Naming the class rather than the symptom - a cleanup step
whose throw skips a later, more important teardown step - turned up two more instances, both fixed
here:
- The commit-lock timeout path released a lock it never took. close() catches the acquire timeout
and carries on to abort anyway, deliberately; the finally then called releaseCommitLock(), which
throws IllegalStateException("Not held be me") when this thread does not hold the write lock.
Same escape, same leak, different door. A commitLockHeld flag fixes it - and under-release was
checked: acquireCommitLock has no path that takes the write lock and then throws, because its
tryLock is the last statement that can fail, so the flag is never false while the lock is held.
- innerDoClose's producer step was the one shutdown step without a guard. Its three predecessors
(commitOffsetsThatAreReady, brokerPollSubsystem.closeAndWait, maybeCloseConsumer) each have
their own try/catch(WARN) from confluentinc#818 / #166. The leak is the smaller half of
the cost: unguarded, the throw reaches the control task's catch, which OVERWRITES failureReason
with "Error from poll control thread", re-runs doClose over already-closed subsystems, and
fails the control future - so the user's close() throws and getFailureCause() stops naming what
actually happened.
Checked and dismissed, with where. ProducerManager#commitOffsets releases through
AbstractOffsetCommitter's own `finally { postCommit(); }`, and preAcquireOffsetsToCommit() sits
OUTSIDE that try, so a failed acquire never reaches the release. The two metrics steps in doClose's
finally are already guarded individually - read the comment above deregisterMeters() there.
RetryQueue's read/write lock pairs all take the lock in the statement immediately before the try, so
the take-then-fail-then-release shape cannot arise. BrokerPollSystem#closeAndWait is a wait, not a
multi-step teardown, and its one caller is guarded with a comment saying why the consumer close must
still run.
RED then GREEN, three tests, each observed failing before its fix:
- ProducerManagerTest#closeStillClosesTheProducerWhenTheAbortThrows -
"ProducerFenced fenced" escaped close;
- ProducerManagerTest#closeStillClosesTheProducerWhenTheCommitLockCannotBeAcquired -
"IllegalState Not held be me";
- ParallelEoSStreamProcessorTest#closeCompletesWhenTheProducerCloseThrows -
"Execution org.apache.kafka.common.KafkaException: simulated producer close failure" out of the
user's close().
All three pass after the fixes, with the Producer's close verified in each and the wrapper's state
asserted CLOSE.
Rejected alternative: catching only around abortTransaction() and leaving releaseCommitLock()
unconditional - which is #410's current shape - fixes the fenced abort and leaves the
lock-timeout leak in place, because that path reaches the release without the lock. Containing the
whole cleanup and putting closeProducer in a finally covers both doors with one shape.
Retires docs/inflight/bug-eos-swallowed-produce-failures.md. Per docs/inflight/AGENTS.md's four
outcomes, what outlives the work was migrated first: the diagnosis lives in this and the previous
commit, and the defect CLASS both holes share - a catch arm that reports and carries on, on a path
whose continuation means success - is written up in
docs/solutions/logic-errors/a-catch-that-closes-and-continues-reports-the-batch-succeeded-2026-09-03.md,
including why hole 1 was masked and the sweep's dismissals. Nothing in the note remained both true
and unowned elsewhere. Its two citations by filename were repointed:
docs/inflight/pr-blockers-and-collisions.md loses the "Two pre-existing main-code holes need their
own PR" decision, which is now taken and landed, and
docs/inflight/core-241-tx-commit-failure-taxonomy.md now cites the write-up and #423 for the
two produce-path shapes its taxonomy must account for.
Prior art superseded: the unpushed local branch fix/producer-close-commit-lock carries two commits
written on 2026-08-10 against the pre-rename package, "release the commit lock only when close
actually took it" and "close the Producer even when the transaction cleanup throws". They were
re-derived here on master rather than cherry-picked, because of the package rename and because their
tests cite line-number facts that no longer hold. That branch can now be deleted.
Still out of scope, and named so nobody re-derives them: the wrapped-ExecutionException retry spin
(#411, #410's to fix with recovery); making closeOnException non-blocking (#410
redesigns that path); ProducerManager#commitOffsets's undifferentiated retry loop (#241); and
restoring the interrupt flag in close's InterruptedException arm, which is pre-existing.
…luding a leaked test thread
The two fixes in this PR are unchanged. This is what running ce-simplify and ce-code-review over the
branch before asking for review actually found, applied.
The one that mattered - a test leaked a control thread into a shared Surefire fork.
invalidPidMappingFailsTheBatchInsteadOfMarkingItSucceeded stubs closeOnException to a no-op so the
batch verdict is observable, and closeOnException is the only route to shutdown on that path. The
spy is a DIFFERENT object from parallelConsumer, and the base class's @AfterEach closes
parallelConsumer - which never ran - so nothing stopped the spy. With produceMessages stubbed to
throw and a 50ms retry delay, its non-daemon control thread retried and logged an ERROR forever,
inside a fork that reuseForks hands to every later test class in the module: CPU load and log noise
that can mask or manufacture unrelated flakiness, which is exactly what this project's no-retry
policy exists to keep visible. The spy is now closed in a finally.
That fix changed the test after its RED observation, so the control arm was re-run rather than
assumed: with the rethrow removed and everything else identical, the test fails again on
"produceMessages ... Wanted *at least* 2 times ... But was 1 time". The proof still proves.
The rest:
- ProducerManager#close's new catch logs through ThrowableUtils.logWithoutEscaping, like the
innerDoClose guard it pairs with. e came from a producer the USER configured, so rendering it
runs their getCause/getMessage inside the logging binding; an escape there would surface out of
close() as a logging stack trace instead of the diagnosis, and read one layer up as the Producer
having failed to close when the finally in fact closed it.
- closeStillClosesTheProducerWhenTheCommitLockCannotBeAcquired's commitLockAcquisitionTimeout drops
to 50ms, and its comment is corrected. It was labelled a failure-side wait, and it is not one:
the produce lock is held for the whole of close(), so tryLock is waited out IN FULL on every
passing run. It cannot flake short either - the read lock is provably held before close() is
called, so the write lock is unavailable however slow the machine.
- The same test now asserts the produce-lock holder thread actually terminated. A join that times
out returns normally, so a hanging finishProducing would have leaked the thread and left the test
passing on everything else.
- ProducerManager#close's comment now states the one place the fix CHANGES behaviour rather than
only containing it: a ConcurrentModificationException from acquireCommitLock means another thread
holds the write lock - the single-writer invariant is already broken by something else - and the
Producer is now closed rather than left running. That is the intended trade; a close that returns
having left the Producer alive is the defect this method exists to stop.
Two things the round found and this commit records rather than changes:
- brokerPollSubsystem.drain(), the FIRST statement of innerDoClose, is structurally the same defect
class and the worst instance if it ever fired, because a throw there skips every remaining
shutdown step rather than the last. It is left alone on evidence: drain() sets a field and calls
ConsumerManager#wakeup, and wakeup is the one consumer method ThreadConfinedConsumer deliberately
does not guard ("left to the delegate (thread-safe per Kafka API)"), with no documented throw.
Guarding a step with no known failure and no test is a behaviour change this PR should not smuggle
in, so it goes in the solutions write-up's dismissal list instead - where, if a throw is ever seen
there, it says the shape was already understood.
- docs/inflight/issue-response-423.md, the draft reply #423 is owed. docs/inflight/AGENTS.md
requires it before the PR merges, because the agents who did the work hold the best context at
merge time and by release time it has to be re-mined from commit logs. NOT POSTED, and deleted
when it is posted rather than when this PR merges.
Not run: the Codex cross-model adversarial pass. Recommended at merge prep for a data-loss fix, and
left to the owner because it spends a metered plan.
Dependency Review✅ No vulnerabilities or license issues or OpenSSF Scorecard issues found.Scanned FilesNone |
✅ Duplicate Code ReportTwo 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
No new clones introduced by this PR. ✅ jscpd (language-agnostic)
No new clones introduced by this PR. Powered by astubbs/duplicate-code-cross-check |
[superseded - a quarantined test changed outcome] 🧪🔒 Quarantine Lane Report
🔴 expected while the owner PR is open · 🟡🎲 flapper, pass proves nothing · 🚨 a deterministic quarantined test passing means its fix landed: delete its Updated for Superseded by a newer quarantine lane report. |
|
🟢 Throughput — OKThis branch measured the same speed as master, on the one test this measures.
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 youThe one that gets misread. 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. 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. Per-method times, not class times. A class time is Reference is the median of 10 recent 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: 28c6be6, c813942, 558fcfb, 5ae7b34, 10ed71c, 75b48bb, e028296, 8df9cc7, 88f5774, ebd8bcb Since the previous push: ratio 1.032 -> 0.995, share 1.594 -> 1.642, rate 68519 -> 64293 (-6.2%). One push of difference sits inside this test's measured spread - read it as movement, not as a result. Updated for |
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## master #429 +/- ##
============================================
+ Coverage 80.18% 82.11% +1.93%
- Complexity 1280 1471 +191
============================================
Files 84 95 +11
Lines 4511 5116 +605
Branches 483 499 +16
============================================
+ Hits 3617 4201 +584
- Misses 693 722 +29
+ Partials 201 193 -8
Flags with carried forward coverage won't be shown. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
@claude review this |
|
Claude finished @astubbs's task in 11m 55s —— View job Review
SummaryThis is a well-scoped, well-documented fix. I traced both defects against the current code (not just the PR description) and they check out. Hole 1 — swallowed Hole 2 — Tests. Ran the full
The commit-lock test's Process/docs. The inflight note retirement matches Minor notes (non-blocking)
No correctness issues found. The fix logic, its regression tests, and the exactly-once reasoning all hold up under inspection. |
Review round dispositionNo correctness findings, so no code changed on this round - HEAD is still 1. "Could not run
Line 67 is the class declaration, which is where SpotBugs anchors that detector. On master the same The other 19 findings are in the On duplication, since the note raised it: both engines report 0 new clones (PMD CPD 27/27, 2. Reconciliation section for #410 - nothing to add. Agreed, no change. Remaining red, and why neither is this branch's
|
…merge-correct What changed for a user: nothing - a one-line marker in an inflight note. Merging master for #430's CVE exclusion brought in a sentence in docs/inflight/deps-cve-backlog.md naming the three PRs the advisory went red on, this one among them. check-branch-self-reference.sh reads whole files for the PR-number arm, so an inherited line naming #429 fails on this branch even though #430 wrote it - the gate only ever looks at the branch about to change state, and after this merge that is this one. The sentence is already written in post-merge terms: it is past tense about the morning of 2026-09-03, and it stays true once this PR lands. That is precisely the case the gate's header calls permanently fine - "PRs are permanent, and citing one is fine forever"; what rots is a deleted branch or a present-tense claim, and this is neither. So the attestation is the correct response rather than a rewrite, and it is added inline so the paragraph is not split. #427 and #428 inherit the same line and will hit the same gate on their own merge of master. The marker is not branch-specific, so whichever lands first carries it for the others.
Taking master turned `repo: hygiene` red here, and the gate is right rather than confused. `docs/inflight/deps-cve-backlog.md` arrived with #430 recording that CVE-2026-82596 "went red on every open PR the morning it was published (#427, #428, #429)", and `bin/check-branch-self-reference.sh` fires on any note naming the PR whose state is about to change - which, on this branch, is #428. The sentence needs no rewrite: it is past tense about the morning of 2026-09-03 and reads correctly once this merges. What was missing is the attestation the gate exists to force, so it is added on the line rather than above it, which would have split the paragraph. This is the intended use of the marker, not a way round the gate: the gate scopes to whichever branch is merging precisely so the PR being named is the one asked to confirm. #427 and #429 are named in the same sentence and will hit the identical failure when they take master; whoever merges first gives master the marker and the other two inherit it.
…merge-correct Merging master brought in #430's backlog entry, which names the three PRs the advisory went red on - this one among them. `bin/check-branch-self-reference.sh` reads its `#NNN` arm over whole files, not just lines the branch added, so a line master wrote about #427 fails on #427's own branch. That is the gate working as designed: it fires only on the state about to change, and mine is. Marked rather than rewritten, which is the judgement the marker exists to record. The sentence is a dated record in the past tense - "It went red on every open PR the morning it was published" - and it names PRs, which the gate's own header calls permanent and fine to cite forever. What rots is a present-tense claim about a branch, and there is none here. It reads identically after any of the three merge. The attesting comment sits inside the checked block on purpose: it names #427 itself, so outside the block it tripped the same arm on its own first run. #428 and #429 will hit this identically when they merge master. This lands the marker once, on the first of the three, so the other two inherit it rather than each attesting the same line.
🧪🔒 Quarantine Lane Report
🔴 expected while the owner PR is open · 🟡🎲 flapper, pass proves nothing · 🚨 a deterministic quarantined test passing means its fix landed: delete its No quarantined test changed outcome since the previous push. Updated for |
…ng branches 34dde7f attested the LatencyUtils entry with the block form (`checked-begin`/`checked-end`) plus a comment explaining the judgement. #428 and #429 had already reached the same conclusion about the same inherited line and used the inline form - the marker appended to the end of line 199 - so three branches were about to land three different renderings of one attestation and conflict with each other over it. The gate accepts both, so this is not a correctness fix; it is choosing the form that merges. The file is now byte-identical to the copy on both sibling branches, verified by diffing against `origin/fix/168-commit-error-line-keeps-identifiers`, so whichever of the three lands first the other two see no conflict here at all. The reasoning the block comment carried is not lost - it is in 34dde7f's body, which is where prose reasoning belongs and which the release-note generator reads. In short: the sentence is a dated past-tense record naming PRs, the gate's own header calls citing a PR permanent and fine, and what rots is a present-tense claim about a branch, of which there is none here. Marked rather than rewritten, deliberately.
…esidentIT setup-guard flake What changed for a user: nothing - a ledger entry. The integration lane went red on this branch with one failure in 201 tests: RegistrationRaceStaleResidentIT.freshArrivalCollidingWithStaleShardResident\ MustStillGetProcessed, on "control thread must reach the mid-loop pause point (offset 25)". The ledger already carries that test from a single 2026-09-01 sighting on #257, with the same assertion and the same reading - the saturation/pause-point setup guard timing out, so the confluentinc#909 assertion the IT exists for is never evaluated. Why it is master-state and not this branch's: the identical failure hit feat/225-pc-built-producer at e35db1e seventeen minutes earlier, and four other branches passed the same lane inside the same 25 minutes (#428, fix/422-commit-interval-unset-detection, optimize/chaos-ci-perf, feat/225-producer-config). Two unrelated branches failing and four passing in one half-hour rules out both "this diff caused it" and "the runner was slow". Mechanism clears this branch independently, which is the part that actually settles it: the IT is PERIODIC_CONSUMER_SYNC driving pc.poll(...), so it never produces, never takes a produce lock, and never reaches the InvalidPidMappingException arm, ProducerManager#close or innerDoClose that this PR changes - and it fails in stage 2 of its own setup, before any close. Recorded rather than diagnosed or fixed, per this ledger's rule and the no-retry rule: the evidence expires with the logs, and the failing runs' 30.8s and 30.9s against 14-19s on the passing ones is the shape of a 30s budget waited out in full. Not quarantined - quarantine is master-state with evidence, and that call is not this PR's to make.
…anch, measured not argued The previous entry called the pair of sightings master-state; two consecutive failures on one branch made that overstated, so the same head was re-run once as a measurement and the outcome is recorded whichever way it went. It passed at the length every green run has, and the record now says so.
|
Integration lane: two consecutive failures on this branch, then a deliberate re-run of the same head passed. Both reds were Recorded in |
…am-memory-leak One conflict, in the untracked-flake ledger, and both sides were right about different rows. #429 raised `RegistrationRaceStaleResidentIT` from one sighting to four - three of them on 2026-09-03, twice in a row on its own head, which then passed on a deliberate re-run. This branch had added a row master does not have, for the `AmbientProbeExtensionTest` log-capture collision it diagnosed. Taking either side alone would have thrown away real evidence: master's better count, or a diagnosed flake nothing else records. Resolved by keeping master's updated row and appending this branch's new one, which is what the file is for - a ledger loses its value the moment a merge starts dropping sightings from it. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RVXBFaG4YXFP7Gjrj6pbVM
…encing-brainstorm Brings in #429 (#423), which closed the two exactly-once holes this branch had also fixed, in the shape it chose to land first. Resolved by hand, taking master's shape wherever the two overlapped: - ProducerManager.close: master's containment - the whole transaction cleanup inside a try whose finally closes the producer, the commit lock released only if it was taken, the catch logging without letting a user producer's exception escape - on top of this branch's opening, which marks recovery terminal to release parked workers and returns when recovery has already discarded the producer. This branch's narrower catch around abortTransaction alone, which master's body names as the rejected shape, is gone. - ParallelEoSStreamProcessor: master's InvalidPidMappingException arm rethrows after the close so the batch is failed; this branch had already folded that arm into produceFailureFor, which records a recoverable condition where PC built the producer and otherwise closes and rethrows the same way, so master's arm is not re-added. Master's test for the hole (invalidPidMappingFailsTheBatchInsteadOfMarkingItSucceeded) runs against this branch's shape and passes, as does its close-completes-when-the-producer-close-throws test; both use this branch's two-argument mock helper, and master's one-argument helper is dropped. - The flake ledger merges the two sides' sightings of the same tests into one count each. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01VrpH51xNDodaajE4P2nhFg
…built-producer Carries up the extraction of ProducerRecovery from ProducerManager, the master merge of #429, and the README's account of why recovery is transactional-only. Resolved by hand: the manager takes the extracted shape, and this branch's two wording differences in the replacement half - the ACL prefix in the terminal failure, the factory in the type description - move with the code into ProducerRecovery; this branch's factory-contract test reaches completeReplacement through recovery() like the rest. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01VrpH51xNDodaajE4P2nhFg
… observation that goes with it Since #429 landed, the recorded history has this guard red on about one head in four across the branches that carry it and green on the ones that do not. Too few runs to attribute, and this branch does not touch the consumer-sync registration path the test drives - so recorded as the next thing to check rather than as a finding. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01VrpH51xNDodaajE4P2nhFg
…e discriminator worth recording RegistrationRaceStaleResidentIT failed on #444 and then PASSED on a deliberate re-run of that job alone, at the identical head. That is the evidence that separates a flake from a regression, and it is the same shape the entry already records for #429 - so the entry now carries it for this sighting too rather than leaving a bare failure count. Recorded now because a re-run's result is only legible while both jobs are still in the run's history, and this ledger exists because CI logs expire and take the evidence with them. The count is left at 8 deliberately: the re-run is not a ninth sighting, it is the control for the eighth. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01WErnxQd9Ew57F9SsqzdPU5
Closes #423. Found by #262's review and parked in
docs/inflight/bug-eos-swallowed-produce-failures.md, which this PR retires.Description
Two exactly-once holes on the transactional produce path, one commit each, plus a third commit
carrying what a local simplify and review round changed before review was asked for.
1.
ParallelEoSStreamProcessor#processAndProduceResultsnow rethrows aftercloseOnException.It caught
InvalidPidMappingException, closed, and fell through toreturn results- so everyWorkContainerin the batch was marked succeeded for output that was never produced. On master theverdict never reached a commit only because
closeOnExceptionblocks the worker until the instance isCLOSED; the test stubs that close to a no-op and shows the swallowed batch produced once, marked
succeeded, and its offset committed - RED - and after the rethrow re-dispatched and uncommitted.
The wrapped
ExecutionExceptionshape (the field report's, confluentinc#830 /#411) is deliberately not touched: it already fails the batch, and its spin
is #410's to fix with recovery.
2.
ProducerManager#closecloses the producer whatever the transaction cleanup throws. A fencedor poisoned producer's
abortTransaction()always throws; it escaped pastcloseProducerand leakedthe
KafkaProducerwith the instance marked CLOSED. The class sweep found two siblings, both fixed:the commit-lock-timeout path released a lock it never took (
IllegalStateException("Not held be me"),same leak), and
innerDoClose's producer step was the one shutdown step without the try/catchconfluentinc#818 gave its neighbours - unguarded, its throw also overwrites
failureReasonwith "Error from poll control thread" and makes the user'sclose()throw. ThreeRED-then-GREEN tests.
Dismissed by the same sweep, with where each was checked:
ProducerManager#commitOffsets(releasedthrough
AbstractOffsetCommitter's ownfinally, whosepreAcquireOffsetsToCommit()sits outsidethe try), the two metrics steps in
doClose's finally (already guarded individually),RetryQueue'slock pairs (every
lock()is the statement before itstry),BrokerPollSystem#closeAndWait(await, and its one caller is guarded), and
brokerPollSubsystem.drain()- structurally the same shapeand the worst instance if it ever fired, left alone because
drain()only sets a field and callsConsumerManager#wakeup, the one consumer methodThreadConfinedConsumerdeliberately does not guard.3. The local review round.
ce-simplifyandce-code-reviewbefore asking for review, applied inthe third commit. The one that mattered: the hole-1 test leaked a non-daemon control thread that
retried forever inside a
reuseForksSurefire JVM, because it stubscloseOnException(the onlyroute to shutdown on that path) and the base class's
@AfterEachcloses a different object. Thatchanged the test after its RED observation, so the control arm was re-run rather than assumed - with
the rethrow removed and nothing else changed, the test fails again on "produceMessages ... Wanted at
least 2 times ... But was 1 time". The commit body has the rest.
Register. No new claim: no documented sentence exists for the arm. The hole-1 test carries
@ProvesClaim(OFFSET_AND_RECORDS_ATOMIC)and C4's note records the second arm and its observed control.Note retirement. Per
docs/inflight/AGENTS.md's four outcomes, the defect class both holes shareis migrated to
docs/solutions/logic-errors/a-catch-that-closes-and-continues-reports-the-batch-succeeded-2026-09-03.mdbefore the note is removed, and both files that cited it by filename are repointed. The draft reply
that directory's rules require,
docs/inflight/issue-response-423.md, is in the branch and is NOT posted.Known testing gaps, named rather than left implicit. No test exercises the REAL blocking
closeOnExceptiontogether with the new rethrow (the unit test must stub it to observe the verdict atall), and none covers a batch where earlier records were already sent before a later send raises the
exception. Both sit on the close path #410 is redesigning.
Reconciliation owed by #410
That draft carries its own versions of both fixes inside its recovery machinery. On its next merge of
master: keep its
produceFailureFor(it already contains the rethrow this PR adds - resolve the conflictin
processAndProduceResults's catch block in its favour), take this PR's fullerclose(Duration)shape(its own catch leaves
releaseCommitLock()unconditional), and drop or fold its duplicate tests(
instancePathFailsTheBatchOnASynchronousInvalidPidMappingExceptionRatherThanMarkingItSucceeded,ProducerManagerDetectionTest#closeStillClosesTheProducerWhenAbortThrows). Its body's "Two are fixedhere" sentence becomes "two were fixed on master by this PR".
Checklist
diagnosis moved to the commit bodies and
docs/solutions/logic-errors/; two citations of itrepointed; the draft issue reply that directory requires is added, unposted
docs/features/- N/A - a defect fix, no featurebodies), and the one whose test changed afterwards was re-checked against a control arm
docs/inflight/working note (pr-/branch-) started at the PR's first commit - N/A - the notethis PR resolves already existed and is retired here; what outlives it is in
docs/solutions/logic-errors/and in the requiredissue-response-423.mddraftce-simplifyandce-code-reviewlocally - simplify found two (log throughThrowableUtils.logWithoutEscaping; a test timeout mislabelled as failure-side and waited out infull), review found one P1 (the leaked test control thread above) plus a missed sweep candidate
and the missing issue-response draft. All applied in the third commit. The Codex cross-model
adversarial pass was NOT run - recommended at merge prep for a data-loss fix, left to the owner
because it spends a metered plan.