Skip to content

Fix deadlock between pc-control and pc-broker-poll threads where partitions are revoked - #548

Merged
John Byrne (johnbyrnejb) merged 5 commits into
confluentinc:masterfrom
nachomdo:bugs/fix-rebalance-eos-deadlock
Apr 3, 2023
Merged

John Byrne (johnbyrnejb) merged 5 commits into
confluentinc:masterfrom
nachomdo:bugs/fix-rebalance-eos-deadlock

Conversation

@nachomdo

@nachomdo Nacho Muñoz Gómez (nachomdo) commented Feb 15, 2023 •

Copy link
Copy Markdown
Contributor

Fixes #541
image

Checklist

  • Documentation (if applicable)
  • Changelog

@what-the-diff

what-the-diff Bot commented Feb 15, 2023 •

Copy link
Copy Markdown
  • The copyright year was changed from 2022 to 2023.
  • A new method called maybeAcquireCommitLock() is added in the onPartitionsRevoked(). This will acquire a lock before committing offsets, which prevents multiple threads from accessing this function at the same time and causing race conditions.
  • In PartitionState class, we have removed some unnecessary code that checks if there are any partitions assigned to it or not because now we can be sure that only one thread accesses each partition state object at a given point of time due to locks acquired by maybeAcquireCommitLock(). So no need for synchronization here anymore!
  • We also remove an unused variable "partition" in ProcessingShard class as well as another check condition (if(workContainer != null)) inside WorkContainer's getWork() method since these were used when parallel consumer had more than 1 worker per topic-partition pair but now with just 1 worker per tp pair, they're redundant and cause confusion/bugs so removing them makes things simpler :)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Please add a relevant test.
Please add changelog entry and do a local build or run process-sources to re-generate readme with updated changelog.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

LGTM

@johnbyrnejb
John Byrne (johnbyrnejb) merged commit 2738fb3 into confluentinc:master Apr 3, 2023
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Jul 28, 2026
Cache the fork<->upstream relationship once, machine-readably, so it stops being
re-derived by hand every session. This fork (bz.stub.parallelconsumer) tracks the
effectively-archived confluentinc/parallel-consumer, whose issues/PRs are a backlog
worth mining and back-linking.

Source of truth:
- src/docs/development/upstream-map.yaml -- one entry per unit of work mapping fork
  branch/PR <-> upstream issue/PR, work group, lifecycle status
  (none|in-progress|ready|pr-open|merged|released), optional reconciliation, todo,
  and a public-facing backlink message. Header documents the schema; carries a
  last_swept date. Design follows Debian DEP-3 / Yocto Upstream-Status / OpenShift
  UPSTREAM.
- src/docs/development/upstream-pr-analysis.adoc slimmed to editorial judgement
  (rankings/verdicts/merge order) with anchors the manifest links to; the manifest
  wins for facts. docs/inflight.md points at the manifest for the durable mapping.

Tooling (scripts/):
- upstream-map.py -- validate | table | refs | show | meta | tracked | posted-refs | todo
- upstream-backlink.sh -- post a "fixed in the fork" / "maintained in a fork" comment
  to an upstream issue/PR, driven by the manifest. Dry-run by default; anti-spam:
  idempotent (skips already-forwarded), per-run cap, delay, status guard. Comment
  body comes from the entry's backlink field (single source of truth) or a template.
- upstream-sweep.sh -- read-only check for NEW upstream activity since last_swept and
  drift on tracked refs; --publish updates a single fork tracking issue.

Conventions: .gitmessage adds DEP-3-style upstream commit trailers (unforced);
AGENTS.md documents the whole system.

Seeded from the analysis doc, inflight notes, git and memory, and reconciled against
a live gh sweep -- which caught drift (upstream confluentinc#541/confluentinc#548 now closed, confluentinc#866 is Kafka
v4 not v7) and new items (confluentinc#892 merged, confluentinc#917/confluentinc#918/confluentinc#919/confluentinc#920/confluentinc#902). confluentinc#859 reconciled:
upstream confluentinc#892 fixed the per-commit meter churn; fork PR #57 fixes the tracking-List
(List->Set) plus assignment-path OffsetMapCodecManager caching.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Jul 28, 2026
Cache the fork<->upstream relationship once, machine-readably, so it stops being
re-derived by hand every session. This fork (bz.stub.parallelconsumer) tracks the
effectively-archived confluentinc/parallel-consumer, whose issues/PRs are a backlog
worth mining and back-linking.

Source of truth:
- src/docs/development/upstream-map.yaml -- one entry per unit of work mapping fork
  branch/PR <-> upstream issue/PR, work group, lifecycle status
  (none|in-progress|ready|pr-open|merged|released), optional reconciliation, todo,
  and a public-facing backlink message. Header documents the schema; carries a
  last_swept date. Design follows Debian DEP-3 / Yocto Upstream-Status / OpenShift
  UPSTREAM.
- src/docs/development/upstream-pr-analysis.adoc slimmed to editorial judgement
  (rankings/verdicts/merge order) with anchors the manifest links to; the manifest
  wins for facts. docs/inflight.md points at the manifest for the durable mapping.

Tooling (scripts/):
- upstream-map.py -- validate | table | refs | show | meta | tracked | posted-refs | todo
- upstream-backlink.sh -- post a "fixed in the fork" / "maintained in a fork" comment
  to an upstream issue/PR, driven by the manifest. Dry-run by default; anti-spam:
  idempotent (skips already-forwarded), per-run cap, delay, status guard. Comment
  body comes from the entry's backlink field (single source of truth) or a template.
- upstream-sweep.sh -- read-only check for NEW upstream activity since last_swept and
  drift on tracked refs; --publish updates a single fork tracking issue.

Conventions: .gitmessage adds DEP-3-style upstream commit trailers (unforced);
AGENTS.md documents the whole system.

Seeded from the analysis doc, inflight notes, git and memory, and reconciled against
a live gh sweep -- which caught drift (upstream confluentinc#541/confluentinc#548 now closed, confluentinc#866 is Kafka
v4 not v7) and new items (confluentinc#892 merged, confluentinc#917/confluentinc#918/confluentinc#919/confluentinc#920/confluentinc#902). confluentinc#859 reconciled:
upstream confluentinc#892 fixed the per-commit meter churn; fork PR #57 fixes the tracking-List
(List->Set) plus assignment-path OffsetMapCodecManager caching.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Jul 28, 2026
Cache the fork<->upstream relationship once, machine-readably, so it stops being
re-derived by hand every session. This fork (bz.stub.parallelconsumer) tracks the
effectively-archived confluentinc/parallel-consumer, whose issues/PRs are a backlog
worth mining and back-linking.

Source of truth:
- src/docs/development/upstream-map.yaml -- one entry per unit of work mapping fork
  branch/PR <-> upstream issue/PR, work group, lifecycle status
  (none|in-progress|ready|pr-open|merged|released), optional reconciliation, todo,
  and a public-facing backlink message. Header documents the schema; carries a
  last_swept date. Design follows Debian DEP-3 / Yocto Upstream-Status / OpenShift
  UPSTREAM.
- src/docs/development/upstream-pr-analysis.adoc slimmed to editorial judgement
  (rankings/verdicts/merge order) with anchors the manifest links to; the manifest
  wins for facts. docs/inflight.md points at the manifest for the durable mapping.

Tooling (scripts/):
- upstream-map.py -- validate | table | refs | show | meta | tracked | posted-refs | todo
- upstream-backlink.sh -- post a "fixed in the fork" / "maintained in a fork" comment
  to an upstream issue/PR, driven by the manifest. Dry-run by default; anti-spam:
  idempotent (skips already-forwarded), per-run cap, delay, status guard. Comment
  body comes from the entry's backlink field (single source of truth) or a template.
- upstream-sweep.sh -- read-only check for NEW upstream activity since last_swept and
  drift on tracked refs; --publish updates a single fork tracking issue.

Conventions: .gitmessage adds DEP-3-style upstream commit trailers (unforced);
AGENTS.md documents the whole system.

Seeded from the analysis doc, inflight notes, git and memory, and reconciled against
a live gh sweep -- which caught drift (upstream confluentinc#541/confluentinc#548 now closed, confluentinc#866 is Kafka
v4 not v7) and new items (confluentinc#892 merged, confluentinc#917/confluentinc#918/confluentinc#919/confluentinc#920/confluentinc#902). confluentinc#859 reconciled:
upstream confluentinc#892 fixed the per-commit meter churn; fork PR #57 fixes the tracking-List
(List->Set) plus assignment-path OffsetMapCodecManager caching.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
… at us

The sweep searched `updated:>=<last_swept>`, and the backlink run
commented on all 78 mirrored upstream issues - which bumped `updated` on
every one of them. The report became the entire upstream tracker, with
no way to see the two items a human had actually touched. A watcher that
fires on everything is the same as no watcher.

It now reports an item only when someone other than us moved it: opened
inside the window, or commented in inside the window by an account that
is not the authenticated user. On the current window that is 78 items
down to 2 - upstream confluentinc#885, where vivahu replied to our 0.5.3.3 answer,
and upstream PR confluentinc#905, where flashmouse replied about us carrying their
PR. Both are exactly the outreach responses the sweep exists to surface.

Also corrects a tracked ref the sweep caught: upstream PR confluentinc#548 was
recorded `open`, the file header claimed `CLOSED`, and GitHub says
merged 2023-04-03. The entry now says merged; the header's STATE NOTE is
gone, since it duplicated per-entry status and had already drifted from
it.

last_swept bumped to 2026-08-06.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01QqHpNSXC39ANv9kG1ZvUzn
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
…in this fork

The bug-857 entry lists upstream PR confluentinc#548 without saying what it is or
where it stands, so it reads like work still to carry. It is not: it
merged upstream in 2023, its merge commit is an ancestor of this fork's
master, and RebalanceEoSDeadlockTest ships with it. Verified both ways -
git merge-base --is-ancestor, and the test file's presence on master.

That matters for scoping the 857 family: the upstream deadlock fix is
accounted for, and the three fork PRs address different defects behind
the same symptom.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01QqHpNSXC39ANv9kG1ZvUzn
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
…e the mirrors

Two gaps in the Upstream tracking instructions, both of which this
session hit for real.

Nobody owned the `upstream:` half. Every lifecycle transition listed was
our own work; upstream-sweep.sh only reports drift and never writes, and
nothing told anyone to act on what it reports. That is how upstream confluentinc#548
sat recorded `open` while merged since 2023, with a header note in the
same file asserting a third answer. Now: fix the entry and bump
last_checked in the same pass, and verify against GitHub rather than the
entry.

The section also still told agents to record upstream *issues* in the
manifest, contradicting the Backlinking section directly below it and
the slimming this PR performed. Issues live in the 78 fork mirrors,
which is where diagnosis, labels and closing belong. The file-map table
carried the same stale claim and is corrected too.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01QqHpNSXC39ANv9kG1ZvUzn
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
…tream confluentinc#857

"upstream" describes a relationship, not a repository, and it is not
stable: this fork is itself upstream to anyone who forks it. That is the
same objection that ruled out "fork" as a qualifier, and it applies
equally here. `#119` / `confluentinc#857` names both ends, is
symmetric, and matches the form GitHub autolinks.

181 tokens renamed across the 26 files this PR already touches. The rest
of the tree is left to the tree-wide sweep. `upstream #NN` still passes
the gate, so nothing existing breaks - the gate's own tests keep
exercising that form deliberately.

The gate now also accepts the spaced variants, so the issue-vs-PR
distinction survives qualification: `confluentinc PR confluentinc#548` reads better
than dropping the word.

Also fixes the CI this work broke. Qualifying references touched three
upstream-derived files (OffsetMapCodecManager, PartitionState,
WorkContainer), and AGENTS.md requires any such file modified after the
fork point to carry a Modifications Copyright line - comment-only edits
still count, which is right, since the check asks whether the file
changed rather than how much. check-copyright-headers.sh runs in Maven's
validate phase, so this did not merely fail its own job: Mutation Tests,
Chaos Pain Suite and Performance all died before running anything. One
missing line, five red checks.

Two of my own added lines were also flagged by my own gate - a "#329" in
its explanatory comment and a quoted "PR #100" in the follow-up note -
and the sweep had produced one double qualification ("upstream is at
upstream confluentinc#922"). All fixed.

Lesson worth keeping: run the gate over `git diff --cached origin/master`
before pushing. It is one node call, and it would have caught both.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01QqHpNSXC39ANv9kG1ZvUzn
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
…r repo

Mirrors all 78 open issues from confluentinc/parallel-consumer into this
fork (#44, #117-#195, label upstream-mirror), each
carrying a code-backed diagnosis, and each backlinked from its upstream
original while that is still possible - archival kills writes, not reads,
so the backlinks were the half with a deadline. Seven are closed against
a released version.

Everything else here follows from that.

Mirroring made bare issue numbers ambiguous. The fork numbers from 1 and
confluentinc reaches confluentinc#922, so the ranges overlap completely: of the 51
numbers cited across the files this touches, 48 exist in BOTH repos
meaning different things. #29 is our rebalance fix and
confluentinc#29 is an async-sending request; #114 is a docs PR and
confluentinc#114 is a GPG key issue. So a reference now names its repo
below #1000, and a CI gate enforces it on added lines.

The gate went through three designs, and the discarded two look
plausible enough to be worth recording. Comparing against "the fork is
at #N" raced - CI read 196 while #197 already existed. Checking
whether a number resolves here fails worse: `#200` resolves, to a fork
issue about ManagedTruth, while the author meant confluentinc#200,
shared-nothing architecture. A wrong reference that resolves is worse
than a broken one, because nothing looks amiss. The rule is textual, so
it makes no API calls and cannot race.

The qualifier names the owner rather than the role - confluentinc#857,
not "upstream confluentinc#857". "Upstream" describes a relationship and is not
stable: this repo is upstream to anyone who forks it. "Fork" is out for
the same reason.

Also swept every reference in the files touched here, fixed the source
comments behind the generated TODO index rather than the index, and
stopped the quarantine fixtures borrowing real PR numbers - #80 and #123
are live fork PRs, so the fixtures read as genuine references.

The map shrinks to match: upstream-map.yaml tracks upstream PRs only,
because issues now live in the mirror, and the manifest-driven backlink
tooling is retired - it commented one issue per map entry, and the map no
longer holds issues.

Two corrections the work surfaced: the sweep was reporting our own
backlink comments as upstream activity, hiding the two real replies among
all 78; and confluentinc#548 was recorded open when it merged in 2023 and
is already in this fork.

Remaining tree-wide references are deliberately out of scope, tracked in
docs/inflight/next-qualify-remaining-refs.md with the Java set already
classified.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01QqHpNSXC39ANv9kG1ZvUzn
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
…r repo

Mirrors all 78 open issues from confluentinc/parallel-consumer into this
fork (#44, #117-#195, label upstream-mirror), each
carrying a code-backed diagnosis, and each backlinked from its upstream
original while that is still possible - archival kills writes, not reads,
so the backlinks were the half with a deadline. Seven are closed against
a released version.

Everything else here follows from that.

Mirroring made bare issue numbers ambiguous. The fork numbers from 1 and
confluentinc reaches confluentinc#922, so the ranges overlap completely: of the 51
numbers cited across the files this touches, 48 exist in BOTH repos
meaning different things. #29 is our rebalance fix and
confluentinc#29 is an async-sending request; #114 is a docs PR and
confluentinc#114 is a GPG key issue. So a reference now names its repo
below #1000, and a CI gate enforces it on added lines.

The gate went through three designs, and the discarded two look
plausible enough to be worth recording. Comparing against "the fork is
at #N" raced - CI read 196 while #197 already existed. Checking
whether a number resolves here fails worse: `#200` resolves, to a fork
issue about ManagedTruth, while the author meant confluentinc#200,
shared-nothing architecture. A wrong reference that resolves is worse
than a broken one, because nothing looks amiss. The rule is textual, so
it makes no API calls and cannot race.

The qualifier names the owner rather than the role - confluentinc#857,
not "upstream confluentinc#857". "Upstream" describes a relationship and is not
stable: this repo is upstream to anyone who forks it. "Fork" is out for
the same reason.

Also swept every reference in the files touched here, fixed the source
comments behind the generated TODO index rather than the index, and
stopped the quarantine fixtures borrowing real PR numbers - #80 and #123
are live fork PRs, so the fixtures read as genuine references.

The map shrinks to match: upstream-map.yaml tracks upstream PRs only,
because issues now live in the mirror, and the manifest-driven backlink
tooling is retired - it commented one issue per map entry, and the map no
longer holds issues.

Two corrections the work surfaced: the sweep was reporting our own
backlink comments as upstream activity, hiding the two real replies among
all 78; and confluentinc#548 was recorded open when it merged in 2023 and
is already in this fork.

Remaining tree-wide references are deliberately out of scope, tracked in
docs/inflight/next-qualify-remaining-refs.md with the Java set already
classified.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01QqHpNSXC39ANv9kG1ZvUzn
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
…r repo

Mirrors all 78 open issues from confluentinc/parallel-consumer into this
fork (#44, #117-#195, label upstream-mirror), each
carrying a code-backed diagnosis, and each backlinked from its upstream
original while that is still possible - archival kills writes, not reads,
so the backlinks were the half with a deadline. Seven are closed against
a released version.

Everything else here follows from that.

Mirroring made bare issue numbers ambiguous. The fork numbers from 1 and
confluentinc reaches confluentinc#922, so the ranges overlap completely: of the 51
numbers cited across the files this touches, 48 exist in BOTH repos
meaning different things. #29 is our rebalance fix and
confluentinc#29 is an async-sending request; #114 is a docs PR and
confluentinc#114 is a GPG key issue. So a reference now names its repo
below #1000, and a CI gate enforces it on added lines.

The gate went through three designs, and the discarded two look
plausible enough to be worth recording. Comparing against "the fork is
at #N" raced - CI read 196 while #197 already existed. Checking
whether a number resolves here fails worse: `#200` resolves, to a fork
issue about ManagedTruth, while the author meant confluentinc#200,
shared-nothing architecture. A wrong reference that resolves is worse
than a broken one, because nothing looks amiss. The rule is textual, so
it makes no API calls and cannot race.

The qualifier names the owner rather than the role - confluentinc#857,
not "upstream confluentinc#857". "Upstream" describes a relationship and is not
stable: this repo is upstream to anyone who forks it. "Fork" is out for
the same reason.

Also swept every reference in the files touched here, fixed the source
comments behind the generated TODO index rather than the index, and
stopped the quarantine fixtures borrowing real PR numbers - #80 and #123
are live fork PRs, so the fixtures read as genuine references.

The map shrinks to match: upstream-map.yaml tracks upstream PRs only,
because issues now live in the mirror, and the manifest-driven backlink
tooling is retired - it commented one issue per map entry, and the map no
longer holds issues.

Two corrections the work surfaced: the sweep was reporting our own
backlink comments as upstream activity, hiding the two real replies among
all 78; and confluentinc#548 was recorded open when it merged in 2023 and
is already in this fork.

Remaining tree-wide references are deliberately out of scope, tracked in
docs/inflight/next-qualify-remaining-refs.md with the Java set already
classified.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01QqHpNSXC39ANv9kG1ZvUzn
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
…r repo

Mirrors all 78 open issues from confluentinc/parallel-consumer into this
fork (#44, #117-#195, label upstream-mirror), each
carrying a code-backed diagnosis, and each backlinked from its upstream
original while that is still possible - archival kills writes, not reads,
so the backlinks were the half with a deadline. Seven are closed against
a released version.

Everything else here follows from that.

Mirroring made bare issue numbers ambiguous. The fork numbers from 1 and
confluentinc reaches confluentinc#922, so the ranges overlap completely: of the 51
numbers cited across the files this touches, 48 exist in BOTH repos
meaning different things. #29 is our rebalance fix and
confluentinc#29 is an async-sending request; #114 is a docs PR and
confluentinc#114 is a GPG key issue. So a reference now names its repo
below #1000, and a CI gate enforces it on added lines.

The gate went through three designs, and the discarded two look
plausible enough to be worth recording. Comparing against "the fork is
at #N" raced - CI read 196 while #197 already existed. Checking
whether a number resolves here fails worse: `#200` resolves, to a fork
issue about ManagedTruth, while the author meant confluentinc#200,
shared-nothing architecture. A wrong reference that resolves is worse
than a broken one, because nothing looks amiss. The rule is textual, so
it makes no API calls and cannot race.

The qualifier names the owner rather than the role - confluentinc#857,
not "upstream confluentinc#857". "Upstream" describes a relationship and is not
stable: this repo is upstream to anyone who forks it. "Fork" is out for
the same reason.

Also swept every reference in the files touched here, fixed the source
comments behind the generated TODO index rather than the index, and
stopped the quarantine fixtures borrowing real PR numbers - #80 and #123
are live fork PRs, so the fixtures read as genuine references.

The map shrinks to match: upstream-map.yaml tracks upstream PRs only,
because issues now live in the mirror, and the manifest-driven backlink
tooling is retired - it commented one issue per map entry, and the map no
longer holds issues.

Two corrections the work surfaced: the sweep was reporting our own
backlink comments as upstream activity, hiding the two real replies among
all 78; and confluentinc#548 was recorded open when it merged in 2023 and
is already in this fork.

Remaining tree-wide references are deliberately out of scope, tracked in
docs/inflight/next-qualify-remaining-refs.md with the Java set already
classified.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01QqHpNSXC39ANv9kG1ZvUzn
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
…mption

The `upstream #NNN:` mirror-title prefix was never a decision. The bulk-import
plan specified `confluentinc#NNN:`; the run shipped `upstream #NNN:`, and the
deviation was later written up in AGENTS.md as though it were intent - complete
with a rule that the prefix "never changes". This PR had then added a further
paragraph exempting it from the owner-not-role convention, which entrenched an
accident as a principle.

Nothing justified the exemption:

- Length, the only plausible defence, is not real. `confluentinc#903:` is three
  characters longer than `upstream confluentinc#903:`.
- The auto-link argument in Mirror format applies to issue *bodies*. GitHub does
  not render links in titles at all, so neither form links there and the role
  word bought nothing.
- Nothing executable parsed the prefix. The couplings were documentation
  `--search` strings, one assertion, and the README line - all of which move
  with the titles.
- GitHub search handles the owner token. Verified before committing to it:
  `--search "confluentinc#622"` finds the mirror by body text, so retitling
  costs no lookup.

Left alone, the exemption would have kept the deprecated form on 78 issue
titles - its most visible surface anywhere in the project - while the tree-wide
sweep removed it everywhere else.

So all 78 mirrors are retitled `confluentinc#NNN: <description>`, and the
documented lookups move with them: AGENTS.md in two places, formatFailure()'s
mirror hint, and the README (edited in src/docs/README_TEMPLATE.adoc and
regenerated - README.adoc is generated). The Issue references exemption
paragraph is deleted rather than reworded; there is now one rule and no
carve-out. Mirror format records what the prefix used to be and why it changed,
so the next reader does not rediscover the question, and the ledger's phase-2
carve-out is replaced by the general habit that produced it: a reference that is
*shown* rather than *made* - a search string, a template, an example of a wrong
form - is not a reference, and rewriting it can break what it documents.

No space in the prefix, matching the house prose form (`#119`,
`confluentinc#857`); the spaced variant exists only to carry the issue-vs-PR
distinction, as in `confluentinc PR confluentinc#548`.

Verified: all 78 retitled with 0 failures and 0 skips, no mirror left on the old
form, and the newly documented lookup returns the expected mirror. Gate tests
26/26 with the assertion updated to require the owner form and reject a stale
role-form hint, bin/check-issue-refs.sh clean, copyright clean.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01YQxeTopekrjwJHKaK3LVeF
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
…mption

The `upstream #NNN:` mirror-title prefix was never a decision. The bulk-import
plan specified `confluentinc#NNN:`; the run shipped `upstream #NNN:`, and the
deviation was later written up in AGENTS.md as though it were intent - complete
with a rule that the prefix "never changes". This PR had then added a further
paragraph exempting it from the owner-not-role convention, which entrenched an
accident as a principle.

Nothing justified the exemption:

- Length, the only plausible defence, is not real. `confluentinc#903:` is three
  characters longer than `upstream confluentinc#903:`.
- The auto-link argument in Mirror format applies to issue *bodies*. GitHub does
  not render links in titles at all, so neither form links there and the role
  word bought nothing.
- Nothing executable parsed the prefix. The couplings were documentation
  `--search` strings, one assertion, and the README line - all of which move
  with the titles.
- GitHub search handles the owner token. Verified before committing to it:
  `--search "confluentinc#622"` finds the mirror by body text, so retitling
  costs no lookup.

Left alone, the exemption would have kept the deprecated form on 78 issue
titles - its most visible surface anywhere in the project - while the tree-wide
sweep removed it everywhere else.

So all 78 mirrors are retitled `confluentinc#NNN: <description>`, and the
documented lookups move with them: AGENTS.md in two places, formatFailure()'s
mirror hint, and the README (edited in src/docs/README_TEMPLATE.adoc and
regenerated - README.adoc is generated). The Issue references exemption
paragraph is deleted rather than reworded; there is now one rule and no
carve-out. Mirror format records what the prefix used to be and why it changed,
so the next reader does not rediscover the question, and the ledger's phase-2
carve-out is replaced by the general habit that produced it: a reference that is
*shown* rather than *made* - a search string, a template, an example of a wrong
form - is not a reference, and rewriting it can break what it documents.

No space in the prefix, matching the house prose form (`#119`,
`confluentinc#857`); the spaced variant exists only to carry the issue-vs-PR
distinction, as in `confluentinc PR confluentinc#548`.

Verified: all 78 retitled with 0 failures and 0 skips, no mirror left on the old
form, and the newly documented lookup returns the expected mirror. Gate tests
26/26 with the assertion updated to require the owner form and reject a stale
role-form hint, bin/check-issue-refs.sh clean, copyright clean.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01YQxeTopekrjwJHKaK3LVeF
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 6, 2026
…mption

The `upstream #NNN:` mirror-title prefix was never a decision. The bulk-import
plan specified `confluentinc#NNN:`; the run shipped `upstream #NNN:`, and the
deviation was later written up in AGENTS.md as though it were intent - complete
with a rule that the prefix "never changes". This PR had then added a further
paragraph exempting it from the owner-not-role convention, which entrenched an
accident as a principle.

Nothing justified the exemption:

- Length, the only plausible defence, is not real. `confluentinc#903:` is three
  characters longer than `upstream confluentinc#903:`.
- The auto-link argument in Mirror format applies to issue *bodies*. GitHub does
  not render links in titles at all, so neither form links there and the role
  word bought nothing.
- Nothing executable parsed the prefix. The couplings were documentation
  `--search` strings, one assertion, and the README line - all of which move
  with the titles.
- GitHub search handles the owner token. Verified before committing to it:
  `--search "confluentinc#622"` finds the mirror by body text, so retitling
  costs no lookup.

Left alone, the exemption would have kept the deprecated form on 78 issue
titles - its most visible surface anywhere in the project - while the tree-wide
sweep removed it everywhere else.

So all 78 mirrors are retitled `confluentinc#NNN: <description>`, and the
documented lookups move with them: AGENTS.md in two places, formatFailure()'s
mirror hint, and the README (edited in src/docs/README_TEMPLATE.adoc and
regenerated - README.adoc is generated). The Issue references exemption
paragraph is deleted rather than reworded; there is now one rule and no
carve-out. Mirror format records what the prefix used to be and why it changed,
so the next reader does not rediscover the question, and the ledger's phase-2
carve-out is replaced by the general habit that produced it: a reference that is
*shown* rather than *made* - a search string, a template, an example of a wrong
form - is not a reference, and rewriting it can break what it documents.

No space in the prefix, matching the house prose form (`#119`,
`confluentinc#857`); the spaced variant exists only to carry the issue-vs-PR
distinction, as in `confluentinc PR confluentinc#548`.

Verified: all 78 retitled with 0 failures and 0 skips, no mirror left on the old
form, and the newly documented lookup returns the expected mirror. Gate tests
26/26 with the assertion updated to require the owner form and reject a stale
role-form hint, bin/check-issue-refs.sh clean, copyright clean.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01YQxeTopekrjwJHKaK3LVeF
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 18, 2026
…disagreed

An earlier draft of this plan said BrokerPollSystem.isResponsibleForCommits() and
AbstractParallelEoSStreamProcessor.isResponsibleForCommits() contradict each
other, and proposed reconciling them. Archaeology disproves it.

Both were created in the same commit - 60e3981, 2020-11-23, "feature: Choose
between Consumer commit or Producer transactional commits" (confluentinc#25) -
which also introduced ConsumerManager, ConsumerOffsetCommitter,
AbstractOffsetCommitter and ProducerManager. Both carry the identical javadoc,
still present verbatim today: "To keep things simple, make sure the correct thread
which can make a commit, is the one to close the consumer. This way, if partitions
are revoked, the commit can be made inline."

They are an XOR over commit mode, not a contradiction: committer.isPresent() is
true iff consumer-commit mode, committer instanceof ProducerManager is true iff
transactional. Exactly one closes the consumer. Neither has been edited since 2020.

So the collision is between #29's new invariant - whole-consumer poll-thread
confinement - and the 2020 invariant it replaced without noticing, that the
committing thread owns the close and which thread that is varies by mode. Control
closing the consumer in transactional mode is original design, not an accident.

Records what history constrains, since this is not a two-line fix:

- Poll-thread ownership was never chosen. af1fa5d (2020-06) extracted the
  blocking poll into its own thread and the consumer went with it; no commit states
  the decision.
- Moving the consumer to control was attempted three weeks after the mode split -
  9dc92e5 (2020-12-03, move-cons-to-pc) - and abandoned. The record does not say
  why.
- Four incidents at this seam, each patched rather than restructured:
  confluentinc#25, the AK 2.7 groupMetadata collision (the metaCache), and
  confluentinc#548, whose fix is the Thread.sleep(100) spin that is now the
  transactional revoke wait behind #44 - then
  confluentinc#857.
- "Control blocks on something only the poll thread produces" outlived the 2020
  lock/Condition scheme, its queue rewrite, and the 2022 actor draft, which swapped
  queues for Future.get() and kept the edge. Mechanism changes have never removed
  it; only moving ownership would.
- Confinement needs ownership transfer, which nothing expresses today -
  ThreadConfinedConsumer.claimOwnership() is claim-only.

Canonical tracker for the redesign is confluentinc#200, fork mirror #142,
still open.
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 18, 2026
…why mechanism fixes never help

The reasoning behind PC's commit machinery existed only in the author's head.
Three separate 2026 investigations each re-derived part of it before acting, and
one of them got it backwards. Reconstructed from the commit record.

The findings that were not previously written down anywhere:

- Poll-thread consumer ownership was never decided. af1fa5d (2020-06) extracted
  the blocking poll into its own thread and the consumer travelled with it. No
  commit states the choice; everything since is machinery to live with it.
- The two isResponsibleForCommits() methods do NOT contradict each other. Both were
  created in 60e3981 (2020-11-23, confluentinc#25), carry identical javadoc, and
  are an XOR over commit mode - exactly one thread closes the consumer. Neither has
  been edited since. The trap is the shared name: one name on two classes reads as
  one question with one answer.
- Moving the consumer to the control thread was attempted three weeks after the mode
  split - 9dc92e5 on move-cons-to-pc - and abandoned. The record does not say why.
  That branch was missing from the catalogue; added to docs/refactoring.md.
- The "owner commits directly" carve-out exists only on the consumer-commit side. In
  transactional mode the revoke callback runs on the poll thread while the committer
  is the producer on control, so there is no direct path. Both confluentinc#548 and
  confluentinc#857 live in that gap.
- Four incidents at this seam, each patched rather than restructured: confluentinc#25,
  the AK 2.7 groupMetadata collision (the metaCache), confluentinc#548 (whose
  Thread.sleep(100) spin is now the transactional revoke wait behind
  #44), and confluentinc#857.
- The fatal edge - control blocks on something only the poll thread can produce -
  outlived the 2020 lock/Condition scheme, its queue rewrite, and the 2022 actor
  draft, which swapped queues for Future.get() and kept the edge. Mechanism changes
  have never removed the deadlock class; only moving ownership would. That is
  confluentinc#200 / #142, "Consider a shared nothing
  architecture, to reduce thread complexity", still open.

Adds a small standalone refactoring item for the naming, which is cheap and
independent of the thread-model work: name the mode-keyed XOR rather than leaving
it for the next archaeologist.
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 18, 2026
…able rule

The invariant was already written down - in the confluentinc#857 solutions doc, under
Prevention - and writing it down is what we did last time too. confluentinc#548 (2023)
and confluentinc#857 are the same defect at the same seam three years apart: a
rebalance callback waiting on something the control thread held. A document fires when
someone chooses to read it; a rule fires because the build runs it.

The rule walks the transitive call graph from every onPartitionsRevoked /
onPartitionsAssigned / onPartitionsLost in our packages, stopping at the JDK and Kafka
client boundary, and fails on Thread.sleep, Lock.lock, lockInterruptibly and
CountDownLatch.await. The reasoning it encodes: a rebalance callback runs on the poll
thread INSIDE consumer.poll(), the whole group waits while it runs, and overrunning
max.poll.interval.ms evicts the member - so anything it cannot get immediately it must
decline rather than wait for. tryLock is the shape that is allowed.

VERIFIED NON-VACUOUS, which is the point of adding it rather than the risk of it. With
the exemption removed the rule fails and names the real defect:

  onPartitionsRevoked reaches blocking call java.lang.Thread.sleep(long)

That is the unbounded `while (isTransactionCommittingInProgress()) Thread.sleep(100)`
wait - which arrived as confluentinc#548's own fix and is now the defect behind
#44, the only issue upstream ever labelled a verified bug. A
rule that passed by finding nothing would have been worse than no rule.

The single exemption is therefore documented as open debt with an owner and a tracking
file, not as an accepted design, and says to delete the entry when #44 lands.
Adding to that list should feel like taking on a defect, because it is one.

Full unit suite green with the rule in place.
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 18, 2026
…probes as candidate work

STRATEGY.md, Reliability track

The correctness bugs this project cares about recur at the same seams, and the record
shows why: confluentinc#548 and confluentinc#857 are one defect three years apart -
both a rebalance callback waiting on the control thread - and both were fixed by hand
with the invariant written into a document afterwards. A document fires when somebody
chooses to read it.

So the track now states the commitment: where an invariant can be expressed
mechanically it becomes a check that runs - an ArchUnit rule, a gate, a probe that can
tell an internal counter from the truth - and where it cannot, saying so beats
pretending a paragraph will hold. STRATEGY.md is a claims document nothing tests, so
this is a claim the work must keep falsifiable; the ArchUnit rule added in a6132f5 is
the first instance of it.

The supporting evidence is this project's own history rather than a principle: a fix
sat unproven for four months because the test written to prove it could not observe it,
and two of the four changes shipped alongside it made things worse.

docs/inflight/next-truth-probes-for-internal-state.md

Ranks the second prevention idea as candidate work. Nothing in the repo could
distinguish "the counter says 0" from "there are 0 records in flight" until one had to
be built to settle a question - and that absence is exactly what let a fix be written
for drift that does not exist, which then caused drift of its own.

Names the two probes that now exist as the model, generalises their shape (arrange the
state deliberately, then assert the internal view against an independently computed
truth - not against itself, and not against "did it stall"), and ranks four candidates
by how much behaviour they gate: shard/queue depth, incomplete-offset tracking, paused
state, and epoch fencing.

Filed in docs/inflight/ rather than docs/refactoring.md because it is a testing
capability rather than a code refactor, and it is unowned.
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 18, 2026
…ed offsets, not the method that once carried them

RebalanceEoSDeadlockTest is upstream confluentinc#548's regression guard for
the confluentinc#541 transactional-mode rebalance deadlock. It proved "the
revoke path committed" by overriding commitOffsetsThatAreReady() and counting
a latch when the override ran on pc-broker-poll. This branch's revoke path
commits via the private tryCommitOffsetsOnRevoke() instead, so the commit
still happened but the latch never counted: measured, the test passed 5/5 on
the defective build and failed 5/5 on the fixed one - a textbook inverted
instrument (see docs/solutions/workflow-issues/
prove-the-problem-exists-before-writing-the-fix.md). It is why the
Integration Tests lane was red on this PR.

What the defect actually was, from the record (upstream 2738fb3): in
PERIODIC_TRANSACTIONAL_PRODUCER, pc-control takes the producer transaction
write lock (A, via maybeAcquireCommitLock -> preAcquireOffsetsToCommit) and
then the commit mutex (B, in commitOffsetsThatAreReady); the pre-fix revoke
callback on pc-broker-poll entered the commit path B-first, then A inside
retrieveOffsetsAndCommit. AB-BA: the poll thread wedged in the rebalance
callback (or died on ProducerManager's cross-thread write-lock
ConcurrentModificationException) and the revoked partitions' completed work
was never committed. The confluentinc#548 fence: the revoke callback waits
out any in-flight transactional commit (the isTransactionCommittingInProgress
spin - its predicate is "lock A held"), and isRebalanceInProgress stops
control starting a new commit cycle mid-rebalance.

The repaired detector observes behaviour instead of an internal method name:

- Deterministic overlap: each pc-control commit cycle dwells 4s between
  acquiring A and touching B; the revoke callback waits (bounded, asserted
  via overlapForced so the test cannot silently go vacuous) for a fresh dwell
  before proceeding - the exact pre-fix fatal window, every run.
- The outcome: the group's committed offsets for the revoked partitions,
  read via AdminClient at callback entry and again at callback return, must
  advance strictly during the revoke window - by the waited-out control
  commit or by the revoke path's own commit. Post-revoke processing is slowed
  500x so redelivery cannot re-advance the offsets before the read.
- Bounded revoke (30s), no recorded failure cause, and post-rebalance
  progress complete the guard.

Mode stays PERIODIC_TRANSACTIONAL_PRODUCER deliberately: confluentinc#541 is
a transactional-mode defect - lock A only exists there. That also bounds the
scope: this test does NOT exercise the confluentinc#857 AB-BA cycle, whose
second edge lives in ConsumerOffsetCommitter and only exists in the
consumer-commit modes; Rebalance857CommitSyncDeadlockProbeIT covers that.

Both arms measured, byte-identical test on each. Defect arm = this head with
the confluentinc#548 fence disabled behind a temporary hardcoded boolean
(spin and gate short-circuited; flag deleted before this commit), per the
red-black-switch technique: fixed arm 10/10 pass, defect arm 10/10 fail,
each failure naming the frozen committed offsets and each run logging the
predicted "write lock already held by another thread" abandonment.

Also fixes an accidental aliasing the upstream test always had:
setupTopic("output-topic") overwrites the inherited topic field, so PC
consumed from and produced back into the "output" topic in a feedback loop.
The output topic is now genuinely separate.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01QiZYbruWmTFGkfepmZn2er
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Sep 2, 2026
…lined, not that it waited

RebalanceEoSDeadlockTest asserted that the group's committed offsets for the revoked partitions had
advanced by the time onPartitionsRevoked returned. On master that held BECAUSE of the defect: the
confluentinc#548 spin waited the in-flight pc-control commit out, that commit landed, and the
assertion was satisfied by the wait #44 (confluentinc#803) is about. Measured before this
change: 5/5 pass on master, 5/5 fail with the decline. It was not a referee the fix upset; it
encoded the wait.

The design question the test kept open - decline, or a bounded wait once producer fencing is
recoverable - was settled by #225's plan on #410. Its KTD11 makes the revoke commit a
third detection site, recorded and declined on the poll thread, and states that once recovery
bounds the write-locked region neither a bounded wait nor deadlining the holder is needed for
liveness. Both stay viable and both are deliberately not taken. Decline stands, and the note
records why, alongside the cost: with control holding the producer write lock when the revoke
lands, the revoked partitions' completed-but-uncommitted work is redelivered to the next owner.

What the test asserts now, with the forced overlap and its vacuity guard unchanged:

- the callback returns inside a quarter of the dwell it was timed into (1000ms of 4000ms). The
  deleted spin returned only when the dwell ended, so it fails this by construction; a deadlock
  still fails the 30s bound. Negative control: with the spin restored in main, 5/5 fail on this
  assertion's own message at 4028-4029ms; reverted.
- the window was resolved, not skipped: the revocation declined the held lock (counted), or - had
  the dwell already ended - committed inline and the offsets moved. Either/or for the same reason
  the probe uses it: after the fix the callback is fast precisely because it declined, and "it was
  fast" alone cannot tell a fix from a window that never opened.
- survival and liveness, as before.

The decline counter is the probe's instrument, extracted: DeclineCountingProducerManager in
integrationTests/utils, with the module that hands it to PC. Revoke857TransactionalWaitProbeIT's
DwellingProducerManager now extends it and keeps only the dwell.

Run on this tree: RebalanceEoSDeadlockTest 5/5, callbacks 43ms, one decline each, offsets unmoved
inside the callback; the probe's defect arm 5/5, ten declines; ArchitectureTest and the test
convention rules green; bin/check-all.sh clean.

Also: a docs/inflight note inherited from master cites this PR as a dated observation and reads
correctly after merge, so it carries the post-merge marker the self-reference gate asks for.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PVL2FEJ645T76PbybEBUZ6
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Sep 7, 2026
…the control thread (#466)

In PERIODIC_TRANSACTIONAL_PRODUCER mode the revocation-time commit ran inline on the broker-poll
thread and never drained the controller's work mailbox, so it could publish a transaction whose
committed offsets omitted records that transaction contained: output committed, input offset not,
and the partition's next owner reprocessed the input and produced the output again.
#436 diagnosed it, refuted claims C9 and C4 of the transactional claim
register, and quarantined the proof with a passing control arm. This is the fix.

What changed for a user: exactly-once holds across a rebalance in transactional mode.
RebalanceEoSDeadlockTest reads the output topic with a read_committed consumer after the revoked
partitions return and fails on a repeated result; against the old code it found three to five
duplicated results in about a hundred and ten, five runs of five, and reads zero five of five with
the fix. C9 and C4 read PROVED again with their one documented exception, the proof left quarantine
and carries both claims, and the README caution is gone.

The fix. In transactional mode onPartitionsRevoked no longer commits on the poll thread. It posts
a request, wakes the control loop with a mailbox message (never an interrupt, which a periodic
commit's timed lock wait would read as a shutdown), and waits, bounded by
commitLockAcquisitionTimeout. The control loop takes the request at the top of a pass and runs its
ordinary sequence - write lock, flush, drain, fence, commit - then completes the request; a pass
that throws fails it; a request the control thread has already taken is waited through to
completion so truncation strictly follows the commit; the close serves whatever is pending at its
own commit and the callback declines at once while the instance is closing. Waiting on the control
thread is safe in this mode and only this mode: the deadlock behind confluentinc#857 is the control
thread blocking on the poll thread, and a transactional commit needs nothing from it. The
consumer-commit modes keep the inline tryLock commit.

The fence. With the drain alone the broker-level check still read two or three duplicates per
run: once the served commit released the write lock, a worker parked on the produce lock resumed
with a record of a partition about to be truncated. PartitionState.fenceForRevocation is set by the
served pass inside the write lock, for the assignment epoch the request was posted with, so a late
pass cannot fence a re-assignment; the produce wrapper checks it right after the produce lock in
both transactional modes and refuses with PCRetriableException.

Rejected before writing code, by experiment: the position recorded at #436's merge prep was to
decline the commit unconditionally. Two arms in ProducerManagerTest showed that a revoke which
commits nothing leaves the output in the open transaction for the next commit to publish without
its offset - the same defect through a different door. Declining is the deadline fallback, logged
at WARN with that cost named, never the fix.

Also: processWorkCompleteMailBox declares @ThreadConfined to the control thread with a runtime
assertion; the confluentinc#548 sleep spin and its ArchUnit exemption are gone; the revoke-drain
inflight note is retired into docs/solutions/logic-errors/the-revoke-path-commit-did-not-drain-the-mailbox-2026-09-07.md;
STRATEGY.md's account of the register is brought current. Collides with
#408 on tryCommitOffsetsOnRevoke; whichever lands second resolves it.

Co-authored-by: Claude Fable 5.1 (1M context) <noreply@anthropic.com>
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Sep 8, 2026
 wins the transactional revoke path

Master landed #466 (the revoke-path commit hands itself to the control thread), #468
(WorkContainer equality is identity) and #471, and this PR went CONFLICTING. #466 is a
different answer to the question this branch answered: in transactional mode onPartitionsRevoked no
longer commits on the poll thread at all - it posts a request, wakes the control loop through the
mailbox, and waits, bounded by commitLockAcquisitionTimeout; the confluentinc#548 spin is gone; and
it refutes this branch's decline by experiment (a revoke that commits nothing leaves its output in
the open transaction for the next commit to publish without the offset - the same duplicate through
a different door). Master's own re-premising of docs/inflight/bug-857-transactional-revoke-wait.md
says what is left for this PR: not the absence of a bound, but whether the bound is the right value
- five minutes, against a max.poll.interval.ms it can exceed.

So the resolution is master's on that path, per file:

- AbstractParallelEoSStreamProcessor: master's onPartitionsRevoked, commitOnRevokeViaTheControlThread
  and the no-argument consumer-commit tryCommitOffsetsOnRevoke replace this branch's parameterised
  decline, its performCommit extraction and its post-catch wake (moot: the served commit runs on the
  control thread, which recovers itself); commitOffsetsThatAreReady is master's again; the mailbox
  loop keeps master's wake-up message skip in front of #410's first-failure try; one of two
  identical assertOnControlThread helpers (#410's and master's) is kept - master's, which
  names the new design.
- RebalanceEoSDeadlockTest: master's whole file. Its unamended assertion that committed offsets
  advance inside the callback holds again by construction under #466, and it now reads the
  output topic at read_committed for duplicates; this branch's decline amendment is superseded.
- ArchitectureTest: master's whole file (#465). This branch's interface-hop widening does not
  merge onto it; its note now records that the blind spot is still open on master and the widening
  is to be re-applied on top of #465 as its own change.
- ProducerManagerTest: master's revoke-request tests and this branch's five revocation tests, both
  kept; PartitionState: master's onSuccess(long) javadoc with #410's ledger cross-reference
  folded in; TransactionalClaim: master's scope note plus #410's C15.
- config/infer-known-findings.txt, bug-857-family.md, core-recoverable-producer-fencing.md: master;
  bug-wedged-after-poisoned-transaction.md: master's deletion (the grooming sweep); the three vetted
  notes keep master's markers and this branch's concurrency label; test-untracked-ci-flakes keeps
  master's rows and this branch's three later sightings.
- ProducerRecoveryTest's revoke-path fence test is retargeted to the served-commit shape: the fence
  now fires on the control thread inside the served pass, the callback returns promptly on the failed
  pass, nothing is stranded, and the replacement is built. Its wake assertion is gone with the wake.

What is red, on purpose: Revoke857TransactionalWaitProbeIT, 5/5, with the callback at 19.2s of a 20s
in-flight dwell against its 10s poll-interval budget - and never the 79s starvation the spin
produced. That is the measurement of #466's bound, and it is this PR's remaining acceptance
test, not a broken instrument. What is now dead main code, held for the owner's call: the three
ProducerManager revocation lock helpers and the DeclineCountingProducerManager instrument, which
count a decline the transactional path no longer makes.

Verified: ProducerRecoveryTest, ProducerManagerTest, ProducerManagerDetectionTest, ArchitectureTest
and the convention rules, PartitionStateAbortedTransactionReplayTest, TransactionalClaimCoverageTest;
RebalanceEoSDeadlockTest 5/5, ProducerFencingRecoveryIT 2/2.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PVL2FEJ645T76PbybEBUZ6
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.

Transactional PConsumer stuck while rebalancing

3 participants