Repository navigation
fix: safely completing doClose() - #818
Merged
Roman Kolesnev (rkolesnev) merged 1 commit intoAug 14, 2024
Merged
Conversation
Roman Kolesnev (rkolesnev)
suggested changes
Aug 12, 2024
Closed
2 tasks done
Fixes confluentinc#809 Co-authored-by: Roman Kolesnev <romankolesnev@gmail.com>
Roman Kolesnev (rkolesnev)
approved these changes
Aug 14, 2024
Contributor
|
/sem-approve |
Wu Shilin (wushilin)
added a commit
to wushilin/parallel-consumer
that referenced
this pull request
Aug 15, 2024
fix: safely completing doClose() (confluentinc#818)
Roman Kolesnev (skl83)
pushed a commit
to skl83/parallel-consumer
that referenced
this pull request
Sep 2, 2024
This was referenced Sep 3, 2026
Antony Stubbs (astubbs)
added a commit
to astubbs/parallel-consumer
that referenced
this pull request
Sep 3, 2026
…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>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Ensure that
doClose()completes and successfully transitions the StreamProcessor to CLOSED state. Unhandled exception duringdoClose()results in a degraded StreamProcessor with terminated control loop but stuck in RUNNING state.This is critical to support accurate external monitoring of
isClosedOrFailed()by user code.see #809
Checklist