Skip to content

fix: safely completing doClose() - #818

Merged
Roman Kolesnev (rkolesnev) merged 1 commit into
confluentinc:masterfrom
gtassone:gt-fixclose
Aug 14, 2024
Merged

Roman Kolesnev (rkolesnev) merged 1 commit into
confluentinc:masterfrom
gtassone:gt-fixclose

Conversation

@gtassone

@gtassone gtassone commented Aug 8, 2024

Copy link
Copy Markdown
Contributor

Ensure that doClose() completes and successfully transitions the StreamProcessor to CLOSED state. Unhandled exception during doClose() 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

  • Documentation (if applicable)
  • Changelog

@gtassone
gtassone requested a review from a team as a code owner August 8, 2024 17:21
@cla-assistant

cla-assistant Bot commented Aug 8, 2024 •

Copy link
Copy Markdown

CLA assistant check
All committers have signed the CLA.

Comment thread .sdkmanrc Outdated
Fixes confluentinc#809

Co-authored-by: Roman Kolesnev <romankolesnev@gmail.com>
@rkolesnev

Copy link
Copy Markdown
Contributor

/sem-approve

@rkolesnev
Roman Kolesnev (rkolesnev) merged commit 5ce0b2f into confluentinc:master Aug 14, 2024
Wu Shilin (wushilin) added a commit to wushilin/parallel-consumer that referenced this pull request Aug 15, 2024
Roman Kolesnev (skl83) pushed a commit to skl83/parallel-consumer that referenced this pull request Sep 2, 2024
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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants