Repository navigation
fix(core): stop a terminally failed send from publishing half a result set under transactions - #261
Conversation
…ions
A terminally failed send in a transactional pollAndProduceMany result set left a
PARTIAL result set visible to a read_committed consumer - two of five records
produced from one source offset visible while three never were - and the
instance neither failed nor shut down, so it was silent.
produceMessages installed a producer Callback that throws from onCompletion,
above a comment reading "only needed if not using tx", while installing it
unconditionally. KafkaProducer#doSend invokes that callback from INSIDE its own
catch(ApiException) handler and only afterwards calls
transactionManager.maybeTransitionToErrorState. The throw escapes first, so the
transaction is never moved into an abortable state, the records already accepted
stay in it, and the next commit publishes them. Only an abort would have removed
them.
The fix honours what the comment already said: throw only when not using
transactions. Nothing is lost by not throwing - the failure still reaches the
work container either way, because processAndProduceResults waits on each
returned Future and an exceptionally-completed send fails the record for retry.
In transactional mode the throw only ever cost us the abort.
TransactionalPartialResultSetIT pins it, and was verified in both directions:
red on the partial-set assertion with the unconditional throw restored
("poison-key-0 has 2 of 5", expected any of [0, 5]), green with the fix. One
term changed, everything else identical.
The test needs one non-obvious mechanism, kept because without it the test
cannot fail. PC only attempts a commit when something has SUCCEEDED since the
last one - wm.isDirty() gates it - and the poison record fails, so with no
further traffic nothing is dirty, no commit is attempted, and the half-populated
transaction is never committed. An earlier draft passed green for exactly that
reason, with and without the fix. Feeding healthy records after the poison one
forces commit attempts against the poisoned transaction, which is what makes the
assertion reachable at all.
What the test deliberately does not assert: that PC recovers or shuts down
cleanly. It does not - the transaction sits in ABORTABLE_ERROR, subsequent sends
throw, and the broker eventually expires it. That is separate, pre-existing
behaviour. The guarantee under test is only that nothing partial becomes
visible.
Defect class swept: "a client Callback that throws, pre-empting the client own
post-callback error handling". Two callback sites exist across core, vertx and
reactor. This is the instance; ConsumerOffsetCommitter OffsetCommitCallback logs
without throwing and is ruled out. No other instances.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013nJJX2zU3hqAh89rNWSvkq
Dependency Review✅ No vulnerabilities or license issues or OpenSSF Scorecard issues found.Scanned FilesNone |
|
Claude finished @astubbs's task in 7m 15s —— View job Claude finished reviewing this PRTodo list:
SummaryThis is a well-scoped, well-evidenced fix. I independently verified the central claim rather than taking it on faith:
Code (
|
✅ 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 |
✅ SpotBugs ReportNo bugs found (new bugs only — baseline from base branch excluded). |
#261 is merged into this branch, so the arm that refuted C7 is enabled and passes, and the register follows the evidence: C2 and C7 move from REFUTED to PROVED. Eleven claims proved, none refuted. The history is kept in both register notes rather than tidied away, because the refutation is what produced the fix. The negative control is also unusually strong here - it is not synthetic. Restoring the unconditional throw, one term with everything else identical, reproduces "poison-key-0 has 2 of 5" exactly, and only that assertion fails while the sibling arms stay green. STRATEGY.md drops its caveat and states the headline unqualified, with two things kept honest: the claim is about Kafka's own topics, so the README's existing warning about duplicate replay into EXTERNAL systems still stands and is untouched; and the register, not the strategy section, remains the gate. The README's machine-checked note loses the "do not rely on all-or-none" warning, which described a live defect and would now be false. It gains the opposite statement plus something the earlier version omitted: that the two timeout guarantees are ATTRIBUTED to existing tests rather than re-proved, with no negative control observed, so neither is recorded as proved. A user reading that section should be able to tell those apart. Regenerated README.adoc from the template rather than editing it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013nJJX2zU3hqAh89rNWSvkq
#261's review surfaced a behavioural consequence worth not losing. Once a terminally failed send correctly moves the transaction to abortable-error, every subsequent record fails with "Cannot execute transactional method because we are in an error state" and the instance stays up but stops progressing, dying only at close. This is not a regression from that fix - before it, the transaction was never poisoned and PC silently published a partial result set instead. Visible wedging beats silent corruption. But it is still a liveness failure, and it is precisely the alive-but-not-progressing shape the chaos suite exists to hunt, so it should not sit unrecorded just because the fix improved on it. The test that found it cannot answer it: defaultMessageRetryDelay is 120s, longer than the window, so the failed records never reach a retry and the run cannot distinguish "recovers once the delay elapses" from "wedged until restart". Recorded with how to settle it - whether a new transaction is ever begun after the poisoned one - rather than guessed at. Cross-linked to the batching stall write-up, which has the same user-visible shape from a different cause, because telling those two apart is the whole difficulty of this area. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013nJJX2zU3hqAh89rNWSvkq
…rately
Four transactional ITs were written by four agents that could not see each
other's files, so each independently grew the same teardown mechanism and one
kept a private copy of a helper that had already been extracted. That is the
predictable cost of the isolation, and this is the pass that pays it back.
- The toClose/register/@AfterEach idiom moves to BrokerIntegrationTest, the
house-designated home for shared broker test helpers. Three of the four had
already converged on the same generic register(T) signature; the fourth had
two type-specific overloads, now subsumed.
- TransactionalCrashReplayIT's private WorkerFailureCapture is replaced by
WorkerFunctionFailureCapture, whose own javadoc had predicted this collapse
"the next time that test is touched". Seven imports verified unused before
removal - the remaining AbstractParallelEoSStreamProcessor mentions are inside
{@code} tags, which need none.
- TransactionalClaimCoverageTest imported the whole compiled test-classes tree
twice per run, once per test method. Memoized, and the duplicated scan
pipeline extracted.
- ProducerManagerSubject#commitLockHeldByCurrentThread had no caller anywhere.
Removed - a Truth subject method nothing calls rots silently, because nothing
fails if it is wrong.
- Three crash-replay arms inherited PAYLOAD_RECORDS=200, whose javadoc justifies
that volume only for the batching regression. They move to 40, the volume the
sibling recombination arm already proves sufficient for the same
crash/fence/replay mechanics; 200 stays, reserved and re-documented.
The teardown is a separately-named @AfterEach rather than folded into close(),
because TransactionMarkersTest declares its own close() that OVERRIDES the base
one - so anything added there would silently not run for that class. Found while
checking all 29 subclasses for collisions. Recorded separately: that override
also means kcu.close() never runs for that class today.
A false statement is also corrected. TransactionalBatchVisibilityIT's class
javadoc still said its failed-send arm "ships @disabled" and that the claim was
refuted, both untrue since #261 landed and the arm was re-enabled. Two
independent reviewers caught it. The evidence record is kept; only the tense and
the @disabled claim were wrong.
Two findings were deliberately NOT applied. The 20s commit-attempt observation
window was reviewed as "likely-margin but not proven-margin", and a ~10s saving
does not justify a test that stays green while detecting less. The warm-up-record
helper would have changed producer lifetime from teardown-close to
immediate-close - unobservable, but not worth adding to a change already this
wide.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013nJJX2zU3hqAh89rNWSvkq
… fix Six reviewers produced more than this branch should absorb. Applying everything would have meant redesigning the gate at PR time on a 5000-line change, so the decisions are written down instead of rushed. bug-eos-swallowed-produce-failures.md - two PRE-EXISTING main-code holes. The serious one: ParallelEoSStreamProcessor catches InvalidPidMappingException, closes, and does not rethrow, so control reaches "return results" and every WorkContainer in the batch is marked SUCCEEDED - offsets committed for records whose output was never produced. Same shape as the defect #261 just fixed, and the single exception to the rationale that justified it ("the failure still reaches the work container either way"). The second: a throwing abortTransaction skips closeProducer, leaking a producer per fenced shutdown - and fencing is a state these very tests now create deliberately. Both want their own PR with a reviewer looking only at them. next-transactional-register-hardening.md - the register's own weaknesses, ranked by how much false assurance each buys. Worst: -Dexcluded.groups= transactions is a documented, supported invocation that runs ZERO claim proofs while the register still reports every claim covered and every sentence intact. The gate is untagged, every proof is tagged, and the gate reads annotations rather than results, so it cannot tell a proof that passed from one never selected. Next: the drift check compares by substring, so a guarantee can be weakened IN PLACE - prepend "Except where a send fails terminally" and append "though this is best effort" and the recorded sentence still matches verbatim. Qualification is the form a walked-back guarantee actually takes, and it is exactly the form that check cannot see. Also recorded there: two ITs now prove the same defect and only one is gated, the new shared teardown makes "PC failed to close" unobservable for ~29 subclasses, and the gate that guards every other gate has no self-test of its own. bug-wedged-after-poisoned-transaction.md gains its answer, settled from the code rather than the experiment it proposed: the only abortTransaction call site is inside close(Duration), and a poisoned transaction stays BEGIN, so lazyMaybeBeginTransaction never opens a replacement. There is no recovery path short of close. What remains open is the design decision - abort-and-reopen versus fail fast - which belongs with the deferred chaos work. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013nJJX2zU3hqAh89rNWSvkq
…hanges) Directory move only, so every path is 100% similar and git's exact-rename detection cannot fail on it. The content edits follow in the next commit. Generated by bin/rename-packages.sh.
Text edits only. No file moves in this commit, so it cannot dilute the rename detection in its parent. Generated by bin/rename-packages.sh.
…roduce-callback-abort # Conflicts: # .idea/runConfigurations/All_examples.xml # .idea/runConfigurations/_Tag__transactions__.xml # README.adoc # bin/rename-packages.sh # parallel-consumer-core/src/main/java/bz/stub/parallelconsumer/internal/Documentation.java # parallel-consumer-core/src/test/java/bz/stub/parallelconsumer/TestConventionRules.java # parallel-consumer-core/src/test/java/bz/stub/parallelconsumer/internal/ProducerManagerTest.java # parallel-consumer-core/src/test/java/bz/stub/parallelconsumer/offsets/OffsetEncodingBackPressureTest.java # parallel-consumer-core/src/test/resources/logback-test.xml # parallel-consumer-examples/parallel-consumer-example-core/src/test/resources/logback-temp-test.xml # parallel-consumer-examples/parallel-consumer-example-reactor/src/test/resources/logback-test.xml # parallel-consumer-examples/parallel-consumer-example-streams/src/test/resources/logback-test.xml # parallel-consumer-examples/parallel-consumer-example-vertx/src/test/resources/logback-test.xml # parallel-consumer-mutiny/src/test/resources/logback-test.xml # parallel-consumer-reactor/src/test/resources/logback-test.xml # parallel-consumer-vertx/src/test/resources/logback-test.xml # src/docs/README_TEMPLATE.adoc
🧪🔒 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 |
Four of the ten review findings held up against the code; the reasoning for the three that did not is on their threads. **cleanUpContext could destroy the failure it was reporting alongside.** It runs in runUserFunction's finally, and an exception thrown from a finally *replaces* the one the catch above is propagating - plain try/finally does not attach it as suppressed the way try-with-resources does. So a throw from either the orElseThrow or from finishProducing's ensureProduceStarted check would have reported "no producer manager" in place of the user function's real failure. It is now caught and logged instead. The log has to carry everything, because the lock is already unrecoverable by that point (takeProducingLock has claimed it out of the context) and the alternative destination, WorkContainer#future, is read by nothing in main. **setProducingLock silently overwrote a lock it was already holding.** The one-lock-one-release invariant was enforced only by caller convention: the two call sites in processAndProduceResults are mutually exclusive on isAllowEagerProcessingDuringTransactionCommit, so today a second set cannot happen. A future call site that broke that exclusivity would have orphaned the first lock - never released - and the next commit's write-lock acquisition would then block forever on a read hold nobody can free. It now fails loudly and attributably instead of hanging. **The surefire-collects rule lived in two places.** TransactionalClaimCoverageTest had re-derived TestConventionRules' four-suffix list verbatim, so surefire's includes changing in one and not the other would let a class pass one gate while failing the other. Extracted as TestConventionRules#surefireCollects and called from both. **ProduceLockReleaseTest now uses the Subject this PR ships.** Its two hold-count assertions went through raw Truth while ProducerManagerSubject existed for exactly that; both are now hasNoProduceLockHolders, keeping their explanatory messages via ManagedTruth#assertWithMessage. ProducerManagerTest's remaining raw reads are deliberate: they snapshot lock state inside a Mockito answer, which is a capture rather than an assertion. Also corrects the cleanUpContext javadoc, which cited ExternalEngine as a path whose produce lock is released here. ExternalEngine's constructor rejects PERIODIC_TRANSACTIONAL_PRODUCER outright, so it can never hold one. Not applied, with the evidence on each thread: extending the produce-callback throw removal into non-transactional mode (a behaviour change to a path #261 deliberately scoped out - raised for a decision rather than taken here); replacing the bespoke log appender with ListAppender (which would retain every event from a hot logger for the test's duration, the thing the append-time filter exists to avoid); and closing registered test clients concurrently (which reintroduces the hang vector the per-client bound exists to prevent). Verified on Temurin 17: parallel-consumer-core unit suite 346 tests, 0 failures, 10 skipped. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…failure one that was missing These were the last three items the stack left open. They were blocked on #260, which rewrote both files they touch; that merged and is now merged into this branch, so the blocker is gone and they are finished here rather than deferred again. Both dark tests failed 100% deterministically, and in both cases the LIBRARY was right and the test was wrong. `offsetsAreNeverCommittedForMessagesStillInFlightLong` asserted, through the transactional producer spy, that NOTHING commits while work is in flight. PC does commit - the base offset - which is a statement about where to resume, not progress. Rewritten onto the mode-agnostic helpers so it means the same thing in all three commit modes, asserting the frontier cumulatively: {3}, then {3,4}, then {3,4,6}. 3 rather than 2 after 0,1,2 complete, because a committed offset is exclusive. `processInKeyOrder` asserted flattened offsets {0,2,3,6,8}. Two separate errors: the same exclusive-offset confusion, and partition 1's base offset is 4 (record creation uses a global counter), which the flattened helper cannot tell from progress because it only trims a genesis of 0. Now asserted per-partition with assertCommitLists, ending (p0=[3,4], p1=[4,7,9]). The invariant it exists for is step E: partition 1 will not advance past 4 while key-2's record 4 is in flight even though 5, 6 and 8 are done, and then jumps to 7 rather than 9, because 8 is complete but not contiguous. `userSucceedsButProduceToBrokerFails` is new - a path the audit found reachable and untested. It asserts the consequence rather than the throw: the offset does not advance, and the record is retried. No producer-callback-thread variant: that is #261's, and it is unreachable here because the harness gives every test one auto-completing MockProducer that calls back on the calling thread. Each was mutation-checked rather than trusted for being green - breaking one expected value (of(3,4,6)->of(3,4,5); p1 4,7->4,9; removing the injected produce failure) reds all three modes. The `@Disabled` import goes with the annotations, its last users. `docs/test-hardening/dark-core-tests-measured-commit-sequences-2026-08-12.md` is deleted in the same change. It was salvaged one commit ago to stop these measurements dying with the plan that held them; they now live in the tests that assert them, with the reasoning in comments. A prose copy of numbers an assertion already pins is the drift this PR removed from `docs/todo-index.md`. Core unit suite: 345 tests, 0 failures under -Pci, up 7 with 2 fewer skips. All four repo gates pass. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…roduce-callback-abort
…al-result-set guard Review of this PR found two ways TransactionalPartialResultSetIT could go green while the defect it pins was present, plus one altitude issue in the fix itself. - Pin max.request.size rather than inherit it. The whole reproduction turns on the producer rejecting the oversized record client side, so the limit it is judged against is now part of the test; OVERSIZED_VALUE_BYTES is derived from it so the two cannot drift apart. Adds createNewProducer(CommitMode, Properties) to KafkaClientUtils, mirroring the existing consumer-side createNewConsumer(boolean, Properties). Overrides apply last; no existing caller passes any, so their behaviour is unchanged. - Assert POISON_KEY exactly. Its oversized result can never be delivered, so NONE visible is the only legal outcome. isAnyOf(0, RESULTS_PER_INPUT) also accepted a full result set that cannot physically exist - weaker than the guarantee being pinned, and weaker than this class's own javadoc already claimed. The other keys keep the all-or-none rule. - Build the send Callback once, in the constructor. Transactionality is settled before ProducerManager exists - ProducerWrapper resolves it into a final field in its own constructor - so the mode decision does not belong in a branch evaluated per completion. - Correct an overclaim in the new comment. The throw was only ever observable on the synchronous pre-accumulator path: when a send fails asynchronously Kafka's own ProducerBatch catches and logs whatever a callback throws, so it was already inert there in both modes. - Use StringUtils.repeat instead of a hand-rolled char[] fill; commons-lang3 is already on this test's classpath. The fix's behaviour is unchanged. Verified with a full build and TransactionalPartialResultSetIT green. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…roduce-callback-abort
… while running The liveness half of what #261 fixed. A result record that can never be sent now correctly moves the transaction into an abortable state, but nothing aborts it while PC runs: abortTransaction() is reached only from close, and the commit that would surface the error is gated on something else having succeeded. No partial set is published - the data guarantee holds - so this is liveness and observability, not corruption. Linked to the retry and dead-letter work because they share a root: PC has no terminal-failure concept, so nothing can decide a record is undeliverable and route it away. Names #149 / confluentinc#310 (no DLQ exists), confluentinc#196 (no max-attempt count - retry is unbounded and purely time-based), confluentinc#291 (terminal vs retryable exception types, closed unmerged), and the adjacent bug-max-failure-history-is-inert note. Records the user workaround available today - wrap the user function so an unsendable record never reaches the producer - which is the same answer the README already gives for poison messages generally. Also notes that "poison" IS discussed here and upstream, with where, so the next session does not conclude from a failed search that the concept is unrecorded. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…duce-send failure Raised while reading back the #261 fix. The surviving throw in ProducerManager's send callback is reachable in ordinary operation - client-side rejection of a result record the user's function returned, RecordTooLargeException being the case the new IT exercises - so it is an expected state, not an internal fault. Calling it InternalRuntimeException makes every reader re-derive that whole chain (which path is reachable, why the async one is not, what actually throws) before concluding it is ordinary failure handling. That derivation was done from source and kafka-clients bytecode during review and was recorded nowhere. Preferred fix is non-breaking: throw a specific subclass so existing catch (InternalRuntimeException) keeps working while the type states the situation. Renaming the base type would be user-visible and is called out as belonging in the breaking-changes section instead. Also corrects this note's own account of confluentinc#291. It is a PR, not an issue, it implements confluentinc#242, and both concern exceptions the USER FUNCTION throws - not send-failure classification. Reading them as cover for send retries was wrong; that work sits under #141 / #149, epic #239. Records that both are already accounted for, so neither gets re-mirrored: confluentinc#291 was closed unmerged on 2023-06-15, in the PR half of the administrative sweep, and is carried in upstream-map.yaml under sweep-2023-admin-closure - it has no dedicated fork issue only because that cohort mirrored the 28 swept issues, with the 35 swept PRs held in the map. confluentinc#242 is not sweep-affected at all, closed as completed by astubbs in 2022. Notes the one real gap: #239's "Prior work" cites confluentinc#366 from that same cohort but not confluentinc#291, which is equally prior art for the retry half of that epic. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…escape Round five, and the same conflation as round four in a new place. The entry claimed asynchronous failures "never reach" the callback. They do: ProducerBatch invokes the same onCompletion, so the log.error fires for effectively every failed send. What is unique to doSend's synchronous catch (ApiException) path is that a throw from the callback can escape to the caller; on the asynchronous path ProducerBatch catches and logs it. Only the second question governs the exception type, since only there is the thrown type observable to a caller - so the scope statement is unchanged in substance. But a refactor built on the wrong reachability model would drop logging for every asynchronous send failure, so the entry now states both questions separately and warns against that specific mistake. The javadoc #261 left in ProducerManager was already right about this - it says the throw was inert on the async path, not that the callback is not invoked. The error was introduced here while paraphrasing it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
… swept PR heads only upstream had (#301) Two records, both from reading back the #261 fix. Documentation only - the sole code change is a TODO marker. `InternalRuntimeException` misnames the produce-send failure ----------------------------------------------------------------------------- `ProducerManager`'s send callback throws `InternalRuntimeException` when a send fails in non-transactional mode, but that is an expected operational state, not an internal fault. The cost is rediscovery: a name saying "internal" sends every reader, human or agent, to re-derive where the callback runs, where its throw can actually escape, and what reaches it, before they can conclude this is ordinary failure handling. That derivation was done from source and kafka-clients bytecode during review and was recorded nowhere. The entry keeps two questions apart, because every loose version of this in review collapsed them. The callback is invoked on asynchronous failures too - `ProducerBatch` calls the same `onCompletion`, so the log fires for effectively every failed send. But a throw from it escapes to the caller only on `KafkaProducer#doSend`'s synchronous `catch (ApiException)` path; asynchronously, `ProducerBatch` catches and logs it. Only the second question governs the exception type, since only there is the thrown type observable to a caller. That set is narrower than "any pre-accumulator failure" - a serializer throwing `SerializationException` propagates without invoking the callback at all - and wider than "oversized records": `RecordTooLargeException` is the case `TransactionalPartialResultSetIT` exercises, while metadata and authorization failures reach the same handler. A refactor that instead models "the callback only runs synchronously" would drop logging for every asynchronous send failure. The recorded fix is non-breaking - a specific subclass, so existing `catch (InternalRuntimeException)` keeps working - with the caveat that makes it actually work: `processAndProduceResults` re-wraps into a generic `InternalRuntimeException`, so the subtype has to be preserved through that outer wrapper or callers still cannot catch it and the refactor buys nothing but a better log line. The lambda parameter should become `sendFailure`; the entry warns against naming it after user code, which is the confusion it exists to stop. A `TODO(refactor)` marker now sits at the site, since a backlog entry with no code-local marker is invisible to grep. Its summary stays on one physical line because `bin/todo-index.sh` indexes only that line - the first attempt wrapped, and generated the fragment "expected state, not an" with diagnosis and remedy dropped. This also corrects the in-flight note #261 added. The retryable half of terminal-vs-retryable exceptions *was* built: `PCRetriableException` is public and `AbstractParallelEoSStreamProcessor` honours it. What is missing is the terminal half - nothing can say "stop, this will never succeed". `confluentinc#291` is a PR, not an issue, implementing `confluentinc#242`, and both concern exceptions the user function throws rather than send-failure classification. Six swept PR heads existed only upstream - now preserved ----------------------------------------------------------------------------- A mirror records what a closed PR said; it does not keep the code. The 35 PRs closed in the 2023-06-15 administrative sweep were reachable through `refs/pull/<n>/head` in the upstream repository, which is not a copy we control. Two things turned out **not** to be loss events, and both were assumed to be in earlier drafts. Deleting the branch behind a PR does not lose the commit, since the base repository holds the pull ref. Nor does a contributor's fork vanishing: every one of the 35 heads was raised from a `confluentinc:` branch, so no third party holds any of them - including `confluentinc#443`, whose author's fork is already gone with no effect. The single upstream exposure is loss of `confluentinc/parallel-consumer` itself. The containment check is restricted to `origin/*`, and that restriction is the whole point: a full clone of this fork also carries `upstream`, and a head contained only in `upstream/*` is exactly the case being hunted. A bare `git branch -r --contains` searches every remote and would have reported every head preserved, so nothing would have been tagged. Reconciled against a live `git ls-remote --heads origin`, because stale tracking refs would otherwise count as safe. That split the cohort: 29 heads are contained by some branch still on this fork - contained by, not raised from; `confluentinc#271`'s own branch is long gone - and six are contained by nothing here. Those six are pinned as annotated `archive/upstream-pr-<n>` tags, whose messages carry the upstream title, author, head branch and closure date so provenance survives without the thread. Tags rather than branches: they are not swept by branch-cleanup tooling, and an annotated tag is fetched by every clone, which a pull ref is not. The tag names, target SHAs and check date live only in `sweep-2023-admin-closure.preserved_heads` in `upstream-map.yaml`. They are deliberately not duplicated into prose - a corrected SHA updated in one copy while the other still read as authoritative is the drift the record exists to prevent. What this does not do ----------------------------------------------------------------------------- The recurring check is **not** automated. `--audit` covers tracking and mirroring, not reachability, and would report clean with every archive tag deleted. It is a manual step until a containment check is wired into `upstream-sweep.sh`. There is also a second, fork-side risk this cannot remove: those 29 heads are safe only because an `origin/*` branch contains them, so deleting one of those branches re-orphans a head that reads as preserved today. That, not any upstream branch, is what a re-run has to watch. #300 tracks this alongside the rest of the upstream-closure decisions nobody has made yet. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…mode-battle-test #262's declared parent had moved a long way: #257 has since merged master up to #403, so taking it also takes nearly all of master, and 262's own CONFLICTING state against master goes with it. 846 files, nine conflicts. Package rename: both sides were already on `bz.stub.*`, so the merge was ordinary and `bin/rename-packages.sh` was not needed. Verified by `git ls-tree -d` on all three refs before starting, not assumed. Conflicts, and which side won: - ProducerManager - TAKE 257. Master has already absorbed 262's produce-callback fix and improved it: the callback is hoisted to a `sendCallback` field with the reasoning in javadoc rather than a comment, and it carries the fact 262 established last (Kafka's own ProducerBatch catches and logs whatever a callback throws, so the throw was only ever observable on the synchronous doSend path). 262's local copy is deleted as the duplicate it now is. `InternalRuntimeException` is `PCInternalRuntimeException` on master; the three places 262 still named the old class are renamed. - AbstractParallelEoSStreamProcessor#cleanUpContext - BOTH. Code takes 262's catch-and-log, which master does not have: this runs in runUserFunction's `finally`, and a throw from a finally REPLACES the exception the catch above is propagating, destroying the user function's real failure. Javadoc takes 262's correction (an ExternalEngine can never reach here holding a produce lock - its constructor rejects PERIODIC_TRANSACTIONAL_PRODUCER, so 257's claim that this is that path's only release is wrong) plus 257's new paragraph on the lock being taken rather than read. - KafkaClientUtils#createNewProducer - BOTH, and neither could just win. 262 added a typed (mode, transactionTimeout, stableTransactionalId) overload; master added (mode, Properties overrides). The auto-merged body already used all three parameters, so both overloads now delegate to one private four-arg builder. - TransactionalPartialResultSetIT - add/add, TAKE 257 plus one thing. 257's is the reviewed descendant of the same test: it pins MAX_REQUEST_SIZE_CONFIG rather than inheriting it, holds POISON_KEY to NONE rather than all-or-none (its oversized result can never be sent, so isAnyOf(0, n) would also accept a full set that cannot physically exist), and uses commons-lang3 `repeat` instead of a local Java-8 helper. Only 262's `@Timeout(600)` guard is carried over. - ProducerManagerTest - TAKE 262's extracted `acquireProduceLockInto` / `assertProduceLockStillOwnedByContext` helpers over 257's inline copies of the same code and comments; the fact 257's comment had and the helper javadoc did not (#257 made cleanUpContext the ONLY release point) is folded into the javadoc. - ProduceLockReleaseTest, WorkContainer - imports only. - docs/quarantined-tests.md is now EMPTY, by two independent routes meeting. #351 diagnosed and fixed OffsetEncodingBackPressureTest on master (it asserted an offset back pressure exists to stop advancing), and this branch performs its own rule-3 re-enable of ProducerManagerTest.producedRecordsCantBeInTransactionWithoutItsOffsetDirect. No @Quarantined annotation remains anywhere in the tree and the registry check confirms 0 entries. - docs/inflight/test-untracked-ci-flakes.md - master's rows (both entries 262 knew about are fixed and gone; simpleBatchTest and processInKeyOrder are new) plus 262's now-current account of the BlockedThreadAsserter collision. Fixed a blank line that was splitting the table in two. Worth flagging for the handoff: master's new `processInKeyOrder` section is about one of the three unexplained failures this branch pushed with at 48d210f, and says it is a solved flake still firing rather than a fresh one. - docs/inflight/pr-blockers-and-collisions.md - master's collision list, plus 262's transactional-stack section updated for reality: #261 merged on 2026-08-14, so the chain is now two PRs, not three. 261 is named rather than deleted because a reader who knows the stack as three links needs telling which one is already master. Verified locally: full-reactor `test-compile` green (main and both test source roots), check-quarantine-registry 0 entries, check-issue-refs clean. The unit suite has not run yet.
…ound Merging #257 brought master's newer repo-hygiene gates onto this branch, and they immediately found real debt in this branch's own notes - not in master's. `bin/check-all.sh` goes from 4 failing gates to 0. **check-inflight-tags** - nine notes here predate the tag scheme and carried no `inflight-type`, so the session index could not sort them and they were invisible at session start. Tagged by CONSEQUENCE per docs/inflight/AGENTS.md, not by filename prefix: the swallowed produce failures are `data-loss` (a batch marked succeeded for records never produced), the commit-interval override is `config-lie`, the poisoned transaction is `stall`, the unclosed Kafka clients are `crash`, and the two "tell someone" notes are `stranded-work`. Register hardening is `task` + `test-debt` rather than `misdirection`: the false-green at the top of its ranking is the strongest item, but `misdirection` is a bug impact and the note is a list of decisions. **check-branch-self-reference** - 24 sentences described #262 or its branch in the present tense, which the merge falsifies. Rewritten in post-merge terms and attested, not marked-and-left: the transactional-stack section now states what each PR IS rather than whether its check is red, the flakes ledger records the #265 collision and the rule-3 lift as history, and two "on this branch" phrases that named nothing are gone. `branch-transactional-battle-test.md` takes `post-merge: exempt-file` because being #262's issue map is the file's whole subject. One claim was wrong rather than merely stale, and is corrected: `pr-strategy-doc-merge-triggers.md` would have read as though the exactly-once headline had to move. It does not. Two claims were refuted against a master without #257 and both read PROVED with it merged in, which is what STRATEGY.md already says. **check-file-refs** - 22 path citations in the dated battle-test plan still said `io/confluent/`, broken by this branch's own rename commits. Repaired to `bz/stub/`, which docs/citations.md requires: the address of the thing the record pointed at may be fixed, its claims may not. Three of them name files the plan proposes to create for the deferred Phase B, so those paragraphs carry `file-refs: N/A` instead - a proposal is not a citation. **Deleted `handoff-262-master-merge.md`.** It was written for an agent picking this branch up cold, and this merge answers three of its four open items: master's #351 fixed the back-pressure flake, its de1620d un-quarantined PCMetricsTest, and it absorbed the produce-callback fix the one unresolved review thread was about. #261 merged, so the fourth is half gone. Its one piece of evidence that was not recorded elsewhere - the three `(CommitMode)[3]` failures seen locally on 2026-08-13 - is folded into the flakes ledger beside master's own independent `[3]` sighting, where someone looking for that failure will actually find it. The two corroborate: master saw `[3]` fail on a branch carrying no main Java, and this branch's own main-code changes are ruled out by grep. The assertion messages from the local run were not kept, so it is NOT established that the two are the same failure, and the note says so. The control that would settle it - the same suite on plain `origin/master` - is still unrun.
…ose never leaks the producer (#429) Two exactly-once holes on the transactional produce path, found by the review of #262 and recorded there rather than fixed. Both closed here, one commit's worth of reasoning each. Hole 1: ParallelEoSStreamProcessor#processAndProduceResults caught InvalidPidMappingException, called closeOnException, and fell through to return a partial - usually empty - result list, so runUserFunctionInternal marked EVERY WorkContainer in the batch succeeded. Output never produced, offset advanced: the same shape as the partial-result-set defect #261 fixed. It now rethrows after the close, wrapped as a PCInternalRuntimeException exactly as the generic arm does, so the batch is FAILED and returned to the mailbox and its offsets can never become commit payload. The instance still closes with the same failure cause. Why nobody saw it: on master the wrong verdict never reached a commit by accident, because closeOnException blocks the worker until the instance is CLOSED and the succeeded containers land in a mailbox nobody drains. The verdict was wrong regardless, and #410 is redesigning exactly the blocking close that masked it, which is why this had to land first. Field history: confluentinc#830 reported the spin, confluentinc#839 wrote this arm; the reporter's wrapped ExecutionException shape takes the generic arm and is #411, deliberately untouched here. Hole 2: ProducerManager#close(Duration) aborted an open transaction inside try/finally that released the lock but did not contain the abort's throw, with closeProducer(timeout) after the block. A fenced or poisoned producer's abort always throws, so close() leaked the KafkaProducer - IO thread, sockets, buffers - while doClose still marked the instance CLOSED. The whole cleanup is now contained, the commit lock is released only if it was taken, and closeProducer runs from a finally. The defect-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, different door), and innerDoClose's producer step was the one shutdown step without the try/catch(WARN) its three neighbours carry from confluentinc#818, so a throwing producer close overwrote failureReason with "Error from poll control thread" and ran doClose twice. Checked and dismissed: commitOffsets releases through AbstractOffsetCommitter's own finally; RetryQueue's lock pairs take the lock in the statement before the try; brokerPollSubsystem.drain() sets a field and calls the one consumer method deliberately left unguarded, and is recorded in the write-up rather than guarded without evidence. RED then GREEN, four tests, each observed failing before its fix: invalidPidMappingFailsTheBatchInsteadOfMarkingItSucceeded (three predicted assertions all observed: the record never re-dispatched, the offset committed, the record marked done; control arm re-run after the review round changed the test), closeStillClosesTheProducerWhenTheAbortThrows, closeStillClosesTheProducerWhenTheCommitLockCannotBeAcquired, and closeCompletesWhenTheProducerCloseThrows. The register gains no new claim - there is no documented sentence for it - so the hole-1 test carries @ProvesClaim(OFFSET_AND_RECORDS_ATOMIC) as a second arm of C4, recorded in C4's note. Rejected: unwrapping ExecutionException into the pid arm (changes the behaviour #411 tracks and collides with #410's recovery); a non-blocking signal-then-throw close (rewrites closeOnException's contract, which #410 owns); catching only around abortTransaction with the lock release unconditional (#410's current shape, which leaves the lock-timeout leak in place). The local review round found the hole-1 test leaking a non-daemon control thread into a shared Surefire fork - the spy that stubs closeOnException is a different object from the one the base @AfterEach closes - now closed in a finally. Retires docs/inflight/bug-eos-swallowed-produce-failures.md; the defect class both holes share is written up in docs/solutions/logic-errors/a-catch-that-closes-and-continues-reports-the-batch-succeeded-2026-09-03.md and the note's two citations are repointed. #410 must reconcile when it next merges master; the list is on that PR. The local branch fix/producer-close-commit-lock carried an older hole-2 fix and is superseded. Also records a RegistrationRaceStaleResidentIT setup-guard flake seen twice on this branch and once elsewhere, with a deliberate re-run of the same head passing at normal length, and attests the inherited LatencyUtils backlog line naming this PR. Upstream-Issue: confluentinc#830 Upstream-PR: confluentinc#839 Applied-Upstream: no Co-authored-by: Claude Fable 5.1 (1M context) <noreply@anthropic.com>
Description
A terminally failed
sendin a transactionalpollAndProduceManyresult set published a partial result set to aread_committedconsumer. Two of the five records produced from one source offset became visible while three never did, and the instance neither failed nor shut down - so it was silent. That is exactly what the all-or-none guarantee of transactional mode denies.Cause: a callback that threw before the client could react
ProducerManager#produceMessagesinstalled a producerCallbackthat threw fromonCompletion, directly beneath a comment reading "only needed if not using tx" - while installing it unconditionally, transactional mode included.KafkaProducer#doSendinvokes that callback from inside its owncatch (ApiException)handler, and only afterwards callstransactionManager.maybeTransitionToErrorState(e). The throw escaped first, so a terminally failed send never moved the transaction into an abortable state. The records already accepted stayed in it, the next commit succeeded, and the half-populated set was published. Only an abort would have removed them.The fix
Throw only when not using transactions, honouring what the comment already said. Because transactionality is settled before
ProducerManagerexists -ProducerWrapperresolves it into afinalfield in its own constructor - the callback is built once in the constructor rather than rebuilt per call, so the mode is not re-decided on every completion.Nothing load-bearing is lost by not throwing: the failure still reaches the work container either way, because
processAndProduceResultswaits on each returnedFutureand an exceptionally-completed send fails the record for retry. In transactional mode the throw only ever cost us the abort.The replacement comment is deliberately precise about how narrow the throw's effect always was. It was only ever observable on the synchronous pre-accumulator path at all: when a send fails asynchronously, Kafka's own
ProducerBatchcatches and logs whatever a callback throws, so it was already inert there in both modes.Evidence
TransactionalPartialResultSetITwas verified in both directions rather than merely observed to pass:poison-key-0 has 2 of 5The narrowness is the point: only that assertion fails, so the behaviour is attributable to the single condition rather than to a broad perturbation.
What makes the test mean anything
Four design decisions are load-bearing, not decoration:
WorkManager#isDirtygates it), and the poison record fails - so with no further traffic nothing is dirty, no commit is attempted, and the half-populated transaction is never committed. Nothing would be visible either way and the test could not tell the fix from the regression.max.request.sizeis pinned, not inherited. The reproduction turns on the producer rejecting the oversized record client side, so the limit it is judged against belongs in the test; the oversized payload is derived from that constant so the two cannot drift apart. This adds acreateNewProducer(CommitMode, Properties)overload toKafkaClientUtils, mirroring the existing consumer-sidecreateNewConsumer(boolean, Properties). Overrides apply last; no other caller passes any, so their behaviour is unchanged.The poison key is held to more than all-or-none: its oversized result can never be delivered, so none visible is the only legal outcome. Asserting
isAnyOf(0, RESULTS_PER_INPUT)there would also accept a full result set that cannot physically exist - weaker than the guarantee being pinned. Every other key, including keys the verifier saw but the test never sent, keeps the all-or-none rule.Defect-class sweep
The class is "a client
Callbackthat throws, pre-empting the client's own post-callback error handling". Two Kafka client callback sites exist across core, vertx, reactor and mutiny.ProducerManageris the instance fixed here;ConsumerOffsetCommitter'sOffsetCommitCallbacklogs without throwing and is ruled out. No other instances.What this does not fix
The guarantee restored here is all-or-none, not liveness. A terminally failed send now leaves the transaction in an abortable state, but nothing aborts it while PC is running:
abortTransaction()is reached only fromclose, and the commit that would surface the error is itself gated on something else having succeeded. So if the failing record is the only uncommitted work, no commit is attempted, no error surfaces, and the instance keeps running with a dead transaction open until close.No partial set is published - the data guarantee holds - and under ordinary traffic the next commit attempt fails loudly. Users can avoid the case today by wrapping their function so a record that can never be sent never reaches the producer.
This is recorded in
docs/inflight/bug-poisoned-transaction-not-aborted-while-running.md, and linked there to the dead-letter and retry work it shares a root with: no DLQ exists (#149, confluentinc#310), retry is unbounded and purely time-based with no max-attempt count (confluentinc#196), and terminal-vs-retryable exception types were proposed and never merged (confluentinc#291). Making PC act on a poisoned transaction is deliberately out of scope here and should be settled with that work.Provenance
Found by the transactional battle-test suite being built on
test/transactional-mode-battle-test, which refutes two documented claims on the strength of this defect -PRODUCE_MANY_ALL_OR_NONEandALL_OR_NONE_PER_SOURCE_OFFSET. Split out into its own PR, offmaster, so the fix reaches other branches without waiting for the larger suite to land.Checklist
KafkaProducer#doSendand why throwing there is unsafe under transactions, plus the in-flight note for the liveness gap left opendocs/features/- N/A - no new or changed user-facing option or feature; this restores documented behaviour rather than adding anyTransactionalPartialResultSetIT, verified red without the fix and green with it