Skip to content

feat(streams) astubbs#255: give a Kafka Streams topology PC's per-key concurrency - #271

Draft
astubbs wants to merge 47 commits into
feats/ks-streams-reconciledfrom
feats/ks-on-pc-spike
Draft

astubbs wants to merge 47 commits into
feats/ks-streams-reconciledfrom
feats/ks-on-pc-spike

Conversation

@astubbs

@astubbs astubbs commented Aug 10, 2026 •

Copy link
Copy Markdown
Owner

Serves #255.

depends on #395
depends on #396
depends on #398
depends on #391
depends on #454

roadmap-stage: N/A - this PR no longer carries the artifact the entry is about. roadmap.yaml's
streams-parallelism-preview still names this PR because it named it when this branch held the whole
module; the module now lives in the rungs above, and what merges here is the study's documentation.
Merging it does not change how real that entry is, which is the question the gate asks. Whether the
entry's pull_request should be re-pointed at a rung is a decision for whichever rung makes the
module publishable - its done_when is an opt-in alpha a reader can follow, and publication is still
off.

Description

This PR has been retargeted onto the top of the Wagon B stack, and what is left is the record, not
the code.
The Kafka Streams execution seam, the machinery under it, the semantics around it and the
demo on top of it are each reviewed in their own rung. What is reviewed here is the feasibility
study itself
: the plan documents, the result write-ups, the branch handover, the ranked open work,
the eight docs/solutions/ learnings this study produced, and the two tests the stack deliberately
left behind.

The base is now feats/ks-streams-reconciled, the branch on which the four sibling rungs were merged
together and the dispatch default was decided against one measurement of the whole. That branch was
merged into this one as an ordinary merge: no history was rewritten and nothing was force-pushed.
The pre-merge tip is preserved on the remote as backup/pre-stack-merge-271.

This is the plan's "retarget move" -
the god-branch decomposition plan, Wagon B, which lives on its own branch and is read with
git show 90188f0bb:docs/plans/2026-08-31-001-process-god-branch-decomposition-plan.md - and it mirrors what
#293 already received on the proxy continent.

The stack this now sits on

Rung PR What it owns
Fork/build machinery #379 patching Kafka Streams at build time, the regeneration discipline, publication off, Kafka's own suite as the behaviour-preservation oracle
Execution seam #388 PcTaskDispatcher, records through PC's WorkManager, wake-on-work, the switch
Refusal envelope #389 joins, windows, suppression, versioned and session stores and exactly-once refused by name, in three layers
Task lifecycle #394 suspend, recycle, close, revive, cooperative rebalance, and the commit frontier across a change of hands
Error surfacing #395 typed control-flow exceptions keeping their type and their timing, the commit fence, the memory bound
Stream time #396 stream time as a low-water mark, punctuation of both kinds, and the divergences it measured
Evidence suite #398 the seam-on divergence lane, the broker-backed laws with their controls, the benchmark with its noise floor
Example #391 a runnable demo with the controls that keep it honest

They merge bottom-up; this one merges last.

What is actually in this diff

Documentation, and two tests. The plan and result documents, the branch handover, the ranked
next-work notes, eight docs/solutions/ write-ups, and:

The dispatch default: measured on the reconciled tree, and it stays OFF

Every rung declined to move it on its own evidence and reserved the decision for the merged module,
because the seam-on numbers move with each of them. The measurement was taken.

What was run. Kafka's own suite twice through the evidence lane - seam off as the control arm,
seam on as the measurement, report directories deleted first. The control arm was clean.

The acceptance test ran red first, which is what makes it a test. The two attributions the
sibling rungs owned were deleted before the run, so a fix the merge had lost would match no
mechanism and be reported as unexplained rather than keep its label by inertia. It exited 1 with
three unexplained divergences.

A control arm then separated "the merge lost a fix" from "the mechanism only partly closed." The
same seam-on suite was run at #396's own tip, in a detached worktree - that
tip carries the stream-time work and not the error-surfacing work, so the one term that differs is
the reconciliation itself.

Case At the sibling's own tip Reconciled Reading
shouldReinitializeRevivedTasksInAnyState fails on all three parameters passes on two, fails on the third the error-surfacing fix is present and did what its rung claimed; the third is Kafka's private processing-threads mode
shouldRecordBufferedRecords fails passes the backpressure work closed it
shouldPunctuateActiveTask, shouldPunctuateWithTimestampPreservedInProcessorContext pass pass closed on the stream-time rung already
shouldPunctuateOnceStreamTimeAfterGap 7 expected, 0 produced 7 expected, 6 produced the clock is no longer stuck; the mark is one record behind
shouldRespectPunctuateCancellationStreamTime fails one assertion fails a different assertion a different line of the test, so a different behaviour

Nothing was lost by the merge, and the residue is a narrowed mechanism rather than the old one:
STREAM_TIME punctuation now lags stock by the work still in flight. That is the low-water mark
doing its job - it never passes a record still in flight, which is what stops a punctuation closing a
window over work still inside the chain - and it is still a divergence a user gets silently, running
in the direction the module's documentation does not state. The README, CONCEPTS.md and the
punctuation note all pin the mark overtaking stock and say nothing about it lagging.

So the default stays off, and that is the fourth time the measurement has named the next reason -
the first from the reconciled whole rather than from one rung. Closing it needs no code fix: it needs
the lagging direction recorded to the same standard as the overtaking one. The three sites that carry
the decision are re-pointed together - PcDispatchSwitch's javadoc, which owns it and now states the
whole chain once; the module README; and the pom's module description. The reservation note is
deleted, its own "delete when" having been the default being re-decided against a fresh measurement
whether or not it moved.

With the two mechanisms re-attributed, the lane exits 0 with every divergence explained and no stale
entry.

The ten findings on this PR, and which rung answers each

These are the unresolved review threads. No reply has been posted to any of them - the mapping is
here so the thread can be read against the rung that changed the code under it.

Finding Answered by
No backpressure on the PC dispatch path #395 - the memory bound applied from the same buffered.records.per.partition stock uses, with occupancy derived from PC's own incomplete set and a two-arm control
A worker's processing failure can be committed past #395 - a pending failure fences commit-data collection, and deliberately does not fence hasUncommittedWork(), which is the opposite question
STREAM_TIME punctuation silently never fires, with no guard #396 - stream time advances, both punctuation types warn at registration, and the residual lag is the default's fourth named trigger above
revive() permanently broken after any dirty close #394 - revival rebuilds the dispatcher, measured seam-on before and after with nine upstream cases going green and none regressing
abortClose() racing dispatchAvailable() throws an uncaught RejectedExecutionException #396 - named as one of three holes the stream-time work recorded against itself and closed, with the in-flight count and the stream-time hold both compensated
recordFailure() silently drops every failure's Throwable after the first #395 - every failure after the first is logged with its record and its cause
close() is missing the owner-thread guard its siblings enforce #394 - answered as a contract rather than a guard: the class javadoc now splits its surface into guard-enforced and convention-only instead of claiming enforcement three of its methods do not have
registerRecords/dispatchAvailable lack assertOwnerThread #394 - the same split; these two are named convention-only
An explicit context().commit() can become a silent no-op Nothing. Still open - pcAwareCommitNeeded() never reads commitRequested, the one field the patch makes volatile for a worker to write
Control-marker-only offset advancement is never detected as needing a commit Nothing. Still open - the consumer-position sweep is unreachable on this path because consumedOffsets is gone by design

The two open ones were re-read against the reconciled tree rather than against the state they were
written on, and are recorded in
docs/inflight/core-streams-two-commit-signals-the-pc-path-cannot-see.md, with
the reason a one-line || commitRequested is not the fix. An eleventh thread asks about the upstream
naming policy and is a question for the owner, not a finding.

How the merge decided which side won

The rule was: the rungs' refined versions win wherever both sides have a file; the residue is kept.
The second half is only safe if this branch contributed nothing to the shared files, and that was
measured.
Both trees were exported with package paths normalised and compared file by file against
this branch's own merge base: outside the streams module and the streams example, every difference
this branch carried is the package rename itself plus a copyright-header line - the same rename
master performed independently.

93 conflicted files took the rungs' version, 32 residue files were kept, 12 files master deleted
stayed deleted, and two shared ledgers were merged as unions rather than resolved to a side. The
check that makes a mis-paired rename impossible to hide is that the merged tree was differenced
directly against the base branch: the only differences are those additions and the two unions.

Both sides turned out to be already renamed, which the plan did not predict - so the
near-identical TestConventionsArchTest.java files surfaced as rename/rename conflicts on the right
files rather than silently landing one module's edit in another, and each was resolved into its own
module.

Verified

JDK 17, macOS. Full reactor test: BUILD SUCCESS, all 13 modules, zero failures - including
core's whole unit suite, the streams module's own suite and Kafka's upstream oracle. The base branch
was verified in its own right first: seam-off oracle unchanged, and all 35 broker-backed integration
arms green against a real broker.

bin/check-all.sh passes every runnable gate on the base branch. On this branch it exits 1 on
check-branch-self-reference alone, with nine findings across five notes that are byte-identical
to their copies on the base branch
, and therefore master's. Marking those would be attesting to
somebody else's note, which that gate's own header says not to do; they are visible here only because
this branch's merge base with master is old, so the gate's phrase arm scopes over master's files.
Three other gates failed on first pass, all of them gates that did not exist on this branch before
the merge
, and all of their findings were on content older than they are - the citations broken by
the rename were repaired to verified successors, the untagged notes were tagged, and the two bare
issue references were qualified.

Two flakes met and neither chased. The full unit suite was run three times on the merged tree:
one run failed ParallelEoSStreamProcessorTest.inFlightMessagesCommittedIfProcessedDuringShutdown,
the next failed a different case in the same class in its own preamble sanity check, and the third
was entirely green. Both are in a module this branch does not modify -
git diff --name-only <base>...HEAD -- parallel-consumer-core is empty - and both are recorded in
the load-tightness ledger, one of them citing a diagnosis on another branch that reddened unmodified
master harder than the branch it was seen on.

Checklist

  • Docs updated - N/A - this PR is now documentation. The code it used to carry is reviewed in the eight rungs above, and the three sites that state the dispatch default were re-pointed on the base branch in the commit that took the measurement
  • User-facing feature documentation data added under docs/features/ - N/A - the module is unpublished, so there is no artifact for a user to depend on and nothing for the ladder to point at. Same reason every rung of this stack gave
  • Tests added/updated - N/A - the two tests still here are inherited, not new: HeadOfLineBlockingBenchmarkTest, which astubbs/parallel-consumer#398 deliberately left on the forest, and ProcessorContextConfinementTest. Every test this study produced that the stack wanted is in the rung that wanted it
  • Title & body reflect the final content of this PR
  • Ran ce-simplify and ce-code-review locally - N/A - there is no code left in this diff for ce-simplify to work on. What a reviewer should spend attention on is whether the residue is the right residue: whether anything in these plans and write-ups should have been extracted into a rung and was not

astubbs added 12 commits August 10, 2026 20:24
…el Consumer

Kafka Streams parallelises across partitions; within one partition
StreamTask hands records to the processor chain strictly one at a time.
Where per-record cost is dominated by blocking IO, that serialisation has
no semantic justification - records on different keys are independent -
and the only way to buy concurrency is to buy partitions.

The investigation asks how little would have to change for PC's
WorkManager to select those records instead, and answers four routes with
their costs. The plan that follows commits to the narrowest one: a seam
inside processor.internals, reached by patching Kafka's own classes at
build time so no Apache Kafka source is committed here.

Two earlier drafts were wrong in ways worth recording rather than
deleting. The first called the whole thing impossible on offset-commit
grounds; the owner's rebuttal - we do not have to commit offsets that way,
and the work-shard manager can replace partitionGroup.nextRecord - was
correct, and the report carries the retraction. The second was re-cut
after review found two P0s in the plan itself.
…nd prove the harness is neutral (U1, U3)

The module unpacks the released kafka-streams sources, applies a tracked
patch, and compiles the result into its own target/classes - which
precede the kafka-streams jar on the classpath, so the same package and
classloader give package-private access without a fork. Only the patch is
tracked: no Apache Kafka source enters this repository.

Proved in the order that makes a green result mean something. First that
the generated classes actually win the classpath race, because if the
jar's copies loaded instead, every downstream comparison would be stock
against stock and would pass beautifully. Then the control arm: with an
EMPTY patch - byte-for-byte the released sources - a stateless topology
behaves identically. Without that baseline a later failure could not be
attributed between the technique and the change.
Stock Kafka Streams keeps "which record am I on" and "which node am I in"
in two fields on AbstractProcessorContext, because exactly one
StreamThread ever owns a task. Running N workers inside one task makes
those the first thing to corrupt, and they corrupt silently - as a record
processed under another record's headers, offset and topic.

Both are confined per thread, and RecordInfo becomes per-record rather
than a single instance fixed for the task's lifetime.
…in (U5-U7)

The central question, answered yes. addRecords feeds PC's WorkManager
instead of the partition group, and process() takes whatever PC hands out
and runs each record's chain on a worker - a switch, never a fan-out.
Measured on the stateless arm: thirty records offered, thirty accepted,
thirty dispatched, thirty succeeded, peak chain concurrency four, and no
key ever concurrent with itself.

Correctness is asserted against genuinely stock Kafka Streams, not against
our own expectations: a fixture generated by a separate module with none
of these classes on its classpath, replayed here, and compared as a
multiset across the run with sequence equality within each key. Order is
asserted where it is actually claimed - per key, which is what PC's KEY
ordering is supposed to preserve - because global order is precisely what
parallel dispatch is allowed to change.

Kafka's own tests are the behaviour-preservation evidence: 188 of them
(StreamTaskTest, RecordCollectorTest, ProcessorContextImplTest) pass
unmodified against the patched classes with the seam off, zero skipped.

And the confinement is proven load-bearing rather than assumed, by two
single-term control arms: un-confining both slots produces a
ClassCastException out of the topology graph, and un-confining only the
record context produces changelog records carrying another record's
timestamp. A green result now distinguishes "confinement works" from
"confinement was never needed here".
…fault

The module publishes rather than staying a throwaway: depending on the
artifact IS the opt-in, so a user who adds it should not then have to
switch it on. Every control arm therefore disables the seam explicitly,
and Kafka's 188 say so on their own surefire execution rather than riding
on whatever the default happens to be - that claim was one default-flip
away from silently becoming 155/188 while still being quoted as 188.

Proved by control arm rather than assertion: with the property removed
entirely so only the default decides, StreamTaskTest gives the same
seam-on result as setting it explicitly true, and returns to 188/188 when
set false.

NOTICE gains the Apache 2.0 changed-files statement the licence requires
once compiled modified classes ship. The module's README leads with the
blast radius, because an existing PC user reading that parallel-consumer
now patches Kafka Streams internals must not conclude their own plain
usage got riskier: the module is a leaf, the patched classes ship only in
its own jar, and the seam is unreachable from core, vert.x and reactor.
…wait that throttles it (U8)

The property this module exists for, measured - and it is latency, not
throughput. One partition, a 1500ms record at the head of the queue,
twentyfour 25ms records behind it on other keys, both arms in one JVM on
the same patched classes with only the seam switched:

              min       p50       p99
  stock     1541ms    1858ms    2205ms
  PC          27ms     232ms      637ms

The minimum is the result worth quoting, because it states the claim
rather than summarising it: under stock even the LUCKIEST record behind
the blocker waited for it. Under PC the quickest paid its own 25ms and
nothing else.

Predictions were written before the run and the negative control passed
the interesting way: with every record on ONE key, PC measures 0.69x -
slower, which is what it must be when KEY ordering forbids concurrency
and the pool handoff still costs. A control that merely tied would be
weaker evidence.

Two corrections during the run, both recorded because they change what
the numbers mean. The control originally varied two terms at once
(cost was keyed on the blocker's key, so the single-key arm made every
record slow); and the p99 was the wrong statistic to assert on, since at
n=24 it is the single worst sample - the last record queued through the
pool - and tracks pool size rather than the seam.

Chasing the control's own anomaly then found the throttle: StreamThread
is one thread that both polls and processes, so blocking up to poll.ms
costs stock nothing. Under the seam that assumption is false - workers
run in the background and a blocked poll stalls dispatch. Changing only
poll.ms, 100 to 1, moves the single-key penalty from ~1695ms to ~24ms and
experiment A's p50 from 8.0x to 19.1x. The benchmark keeps the DEFAULT so
it reports what a user gets today; the fix is recorded with the asymmetry
that shapes it - quiescent needs no wake, since work can only arrive
through the poll itself, while the busy case wants to react to a
completion, which never arrives through the consumer at all.
…mmit-frontier fix (U9)

With the seam on, 33 of Kafka's 101 StreamTaskTest cases fail. Read and
grouped rather than fixed, to find how much of the gap is work and how
much is design: offset and commit accounting 14, close/suspend/recycle 5,
buffering 4, error surfacing 3, EOS gating 3, stream-time punctuation 2,
ordering 1, metrics 1. Twentysix of thirtythree are work, not design.

The root cause was settled by probe rather than assumed. Most failures
read "expected true but was false", which two very different causes
produce: work dispatched but the assertion beating the worker, or work
never becoming dispatchable because these tests drive StreamTask with a
mock consumer. Instrumenting dispatchAvailable gives dispatched=2
available=2 inFlight=2 - the harness is fine, and these are synchronous
assertions against an asynchronous dispatcher.

That makes the largest pile PC's own competency, and the fix a deletion
rather than a repair: Streams keeps one Long per partition, a high-water
mark, which is sufficient under sequential processing and incapable under
concurrent dispatch - no locking lets one number express "12 done, 10 and
11 still in flight". The data structure is the defect.

Also settles how the unsupported API surface should be refused - by
annotating and throwing, never by deleting the methods, since Kafka's own
suite calls them and deletion would forfeit the 188-test evidence - and
records that the commit metadata field has one owner, with anything
Streams needs riding inside PC's encoding as an opaque blob rather than a
second writer.
…r position (U9)

The only shortcoming that could lose data, closed. On the PC path the
consumer-group commit now carries PC's frontier - the lowest incomplete
offset - plus its encoded map of completed offsets beyond it, instead of
Streams' consumedOffsets.

Written red-first, and the red is the demonstration: with the partition
group empty on the PC path, stock's committableOffsetsAndMetadata falls
through findOffset to consumer.position() and commits EVERY polled
record, the parked one included. The test asserts the committed offset
while a record sits parked inside the chain - 11 before, 0 after.

Both halves of the commit protocol are wired. Both commitNeeded gates
re-point (the method, and the private field prepareCommit checks before
ever asking for offsets - either one missed would mean nothing ever
commits), collection happens BEFORE the producer flush (workers complete
during a flush, and KafkaProducer.flush does not cover sends enqueued
after it was invoked), and the success acknowledgement lives in
updateCommittedOffsets, which Kafka invokes only from the success
branches - postCommit is reached after swallowed commit failures and with
no commit at all, so acking there would mark work clean that was never
durably committed.

Crash safety is proven end to end: a genuine abort - no drain, no
feed-back, workers interrupted - leaves the frontier committed, and the
restart redelivers the record that was in flight. A seam-off instance can
then take the same group over and run, decoding PC's payload leniently.

Two product defects surfaced by chasing refuted predictions rather than
by inspection: process() reported progress as dispatched-to-pool, so
consuming a batch of corrupted records reported no progress to the caller
that paces on it; and corrupted records were shipped to workers as no-op
runs, which under KEY ordering stalled the poison pill's successors by a
full pump cycle. Both restore stock parity.

The pile-A test-name delta is zero, and that is a discovery rather than a
failure: those unit tests compare metadata ENCODINGS, which this design
deliberately diverges from, so crash safety was only ever provable in the
integration arm - where it now is, red-then-green.
Three reviewers; the findings worth taking, plus one that turned out to
be a defect rather than a simplification.

One helper replaces the dispatcher-aware commit gate that had been
copy-pasted into three call sites. Both branches survive at every site -
the stock path stays byte-identical - but the three gates must always
agree about whether work is outstanding, and three copies is how they
would silently stop agreeing.

The dispatcher's unit test is now @isolated. It calls abortAllActive(),
which kills every live dispatcher in the JVM, while SAME_THREAD only
serialises that class's own methods - and this module inherits concurrent
test-class execution from core's test jar, so a sibling's dispatcher was
reachable. The registry's javadoc now carries that invariant instead of
leaving it unwritten.

Two comments had gone stale against their own code, and the integration
test's end-offset lookup uses the shared AdminClient rather than
constructing a KafkaConsumer for one stateless metadata call.

Not applied: prepareRecycle never closes the dispatcher, leaking the
registry entry, the worker pool and the WorkManager's partitions on
active/standby recycling. That is a defect, so it belongs to the
rebalance work rather than to a refactor commit.
… and fail loudly on revive

Status words out of the identifiers. "Spike" and "experiment" describe
the work's STATUS, not what the code IS - status is temporary and
external, so it belongs in the plan, the inflight entry and the module
description, and the moment the experiment graduates every identifier
carrying it starts lying in a way that reads as accurate to a newcomer.
So the module, the package, the patch files, the system property and the
pom properties all lose it, and the patch's comment markers now say "PC
dispatch", which is what those hunks actually are. Done now because
publication is what makes a name expensive to change, and publication is
currently gated. Credit to the Connect agent, who made the same mistake
at smaller scale and passed on the correction.

Kept deliberately: "Experiment A" and "the control arm" in the benchmark.
That is the scientific term for a controlled measurement with a stated
prediction and a negative control - method vocabulary describing method,
which is correct usage exactly as status vocabulary describing status is.

Revive now fails loudly instead of stalling silently. closeDirtyAndRevive
resurrects the same StreamTask, whose dispatcher was closed on the way
down, so a revived task registered records into a dispatcher that never
handed them out: no progress, no exception, nothing logged. Overriding
revive() in the already-patched StreamTask keeps the surface at four
classes. It is the loud-failure floor; recreating the dispatcher stays
open.

And the plan gains a shipping target it did not have. The piles are the
roadmap - every one now has an owner, a deferral with a reason, or a
by-design verdict, because an unassigned pile is how a divergence
survives to a release. Pile B goes with rebalance, since B's five
failures ARE close/suspend/recycle and rebalance is what drives those
transitions. Rebalance itself is preview-DISCLOSURE rather than
preview-blocking: the preview's contract is at-least-once inside a stated
envelope and rebalance duplicates sit inside it - an earlier draft ranked
it blocking by applying a production bar to a preview. What the envelope
must say is recorded too, including that "stateful works" is true while
"stateful works at the same cost" is not.
The commit design's central idea had been described by paraphrase -
highest sequentially succeeded offset, safe resume point, lowest
incomplete - depending on which document was speaking. One name for it:
the frontier, and frontier semantics for the commit-the-frontier,
encode-the-holes design built on it.

Defined against its contrast (the high-water mark, which assumes
sequential completion) and its structural twin (TCP's cumulative ACK plus
SACK blocks), because the contrast is what makes the Kafka Streams
module's commit fix legible as a deletion rather than a repair.
…adata collision actually did

STRATEGY.md arrived from master carrying a merge-trigger note that asked
whichever of the Streams or Connect branches landed first to settle the
persona question, so the later one inherits a decision rather than
reopening it. Streams arrived first, so it is settled here.

A second persona, not a widened primary one. "My downstream outscales my
partitions" and "my topology has one slow stage" are different
self-descriptions reached from different starting points - the second
group arrives with a topology or a sink connector, is not shopping for a
consumer library, and what they are buying is that their code does not
change. One sentence covering both would describe neither. Connect
inherits this and should extend the second persona rather than
re-litigate whether one exists.

"Our approach" gains its second clause: the same client-side engine can
run underneath another framework, not only inside your own application,
with the measured head-of-line result as the evidence. And Target problem
gains the reason Share Groups do not reach this - they operate at the
consumer layer, while Streams and Connect build their own execution model
above it and process a task's records one at a time regardless.

That framework-layer point had been sitting in our promotional notes as a
line to use; it is an argument about what the product IS, so it belongs
upstream in the strategy rather than only in material quoting it.

The metadata-collision trigger has half fired, and the entry now says so:
the module met the collision and did not merely survive it - PC owns the
field outright on its dispatch path and Streams' own metadata is
deliberately not written. Still not a STRATEGY.md caveat, since a guiding
policy should not carry footnotes for handled risks, but the shape of the
resolution is strategic and worth watching: a generalised rider inside
PC's encoding, one owner and one decoder. If that direction ever fails,
the substrate role has a real limit and that IS a claim in "Our approach".

Also corrects our own note, which sent readers to STRATEGY.md for the
Share Groups comparison. The substantive table is in the generated
README; STRATEGY.md carries one sentence of it.
@chatgpt-codex-connector

Copy link
Copy Markdown

You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard.
To continue using code reviews, you can upgrade your account or add credits to your account and enable them for code reviews in your settings.

@claude

claude Bot commented Aug 10, 2026 •

Copy link
Copy Markdown

Claude finished @astubbs's task in 4m 10s —— View job


Claude is reviewing this PR

Task list

  • Gather context (PR diff, changed files)
  • Understand the request (code review via /code-review)
  • Run code review across the diff (line-by-line scan, cross-file trace, removed-behavior audit, cleanup/efficiency scan) — in progress
  • Verify findings and post results

This is a large PR (44 files, ~10.9k additions) adding a new parallel-consumer-streams module. Review is running in the background; this comment will be updated with findings once complete.

@github-actions

github-actions Bot commented Aug 10, 2026 •

Copy link
Copy Markdown

Dependency Review

✅ No vulnerabilities or license issues or OpenSSF Scorecard issues found.

Scanned Files

None

@github-actions

github-actions Bot commented Aug 10, 2026 •

Copy link
Copy Markdown

✅ Duplicate Code Report

Two 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

PR Base Change
Clones 38 36 🫤 +2
Duplicated lines 1257 1157 :face_with_raised_eyebrow: +100
Duplication 0.15% 0.14% 🫤 +0.01%
Rule Limit Status
Max duplication 0.5% ✅ Pass (0.15%)
Max increase vs base +0.1% ✅ Pass (+0.01%)
⚠️ 2 new clones introduced
  • 28 lines (96 tokens): parallel-consumer-examples/parallel-consumer-example-streams-pc/src/main/java/bz/stub/parallelconsumer/examples/streams/pc/Latencies.java:34 <-> parallel-consumer-streams/src/test/java/bz/stub/parallelconsumer/streams/integrationTests/HeadOfLineBlockingBenchmarkTest.java:338
  • 22 lines (72 tokens): parallel-consumer-examples/parallel-consumer-example-streams-pc/src/main/java/bz/stub/parallelconsumer/examples/streams/pc/ArmRunner.java:195 <-> parallel-consumer-streams/src/test/java/bz/stub/parallelconsumer/streams/integrationTests/HeadOfLineBlockingBenchmarkTest.java:246

✅ jscpd (language-agnostic)

PR Base Change
Clones 119 115 🫤 +4
Duplicated lines 1696 1656 :face_with_raised_eyebrow: +40
Duplication 1.11% 1.15% 🙂 -0.04%
Rule Limit Status
Max duplication 2% ✅ Pass (1.11%)
Max increase vs base +0.1% ✅ Pass (-0.04%)
⚠️ 5 new clones introduced
  • 11 lines: parallel-consumer-streams/src/test/java/bz/stub/parallelconsumer/streams/integrationTests/HeadOfLineBlockingBenchmarkTest.java:6 <-> parallel-consumer-streams/src/test/java/bz/stub/parallelconsumer/streams/integrationTests/WakeOnWorkShutdownTest.java:7
  • 9 lines: parallel-consumer-streams/src/test/java/bz/stub/parallelconsumer/streams/integrationTests/HeadOfLineBlockingBenchmarkTest.java:231 <-> parallel-consumer-examples/parallel-consumer-example-streams-pc/src/main/java/bz/stub/parallelconsumer/examples/streams/pc/ArmRunner.java:180
  • 10 lines: parallel-consumer-streams/src/test/java/bz/stub/parallelconsumer/streams/integrationTests/HeadOfLineBlockingBenchmarkTest.java:349 <-> parallel-consumer-examples/parallel-consumer-example-streams-pc/src/main/java/bz/stub/parallelconsumer/examples/streams/pc/Latencies.java:54
  • 14 lines: parallel-consumer-streams/src/test/java/bz/stub/parallelconsumer/streams/integrationTests/BackpressureBoundIntegrationTest.java:10 <-> parallel-consumer-streams/src/test/java/bz/stub/parallelconsumer/streams/integrationTests/ShadowedStreamsControlTest.java:10
  • 10 lines: docs/solutions/best-practices/choose-the-statistic-that-states-the-claim.md:169 <-> parallel-consumer-streams/src/test/java/bz/stub/parallelconsumer/streams/integrationTests/HeadOfLineBlockingBenchmarkTest.java:127

Powered by astubbs/duplicate-code-cross-check

@github-actions

github-actions Bot commented Aug 10, 2026 •

Copy link
Copy Markdown

📌 Duplicate code detection tool report

The tool analyzed your source code and found the following degree of similarity between the files:

🆕 New file similarities introduced

File A File B Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerOptions.java parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/PartitionState.java 30.1

🔺 Increased similarities

File A File B Base (%) PR (%) Change
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/offsets/OffsetEncodingBackPressureTest.java parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/offsets/OffsetEncodingBackPressureUnitTest.java 39.2 39.9 +0.6
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PCRetriableException.java 36.5 36.7 +0.2
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java 34.5 34.7 +0.2
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 34.5 34.7 +0.2
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 62.9 63.0 +0.1
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java 37.0 37.1 +0.1
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 37.0 37.1 +0.1
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java 34.9 35.0 +0.1
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/AbstractRevokeUnderWorkScenario.java parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosRevokeUnderWorkIT.java 35.6 35.7 +0.1
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerOptions.java parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/ProducerManager.java 32.4 32.4 +0.1
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/ParallelEoSSStreamProcessorRebalancedTest.java parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessorTest.java 34.8 34.8 +0.1
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyPCTest.java parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorPCTest.java 71.1 71.2 +0.1
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/AbstractRevokeUnderWorkScenario.java parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosChurnStormIT.java 48.7 48.8 +0.0
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CoreBatchTest.java parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyBatchTest.java 50.7 50.8 +0.0
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CoreBatchTest.java parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorBatchTest.java 52.0 52.1 +0.0
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/PCModule.java parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/PCModuleTestEnv.java 32.2 32.2 +0.0
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/DrainingMemberRebalanceIT.java parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/BrokerPollSystemDrainTest.java 31.1 31.2 +0.0
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CoreBatchTest.java parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/VertxBatchTest.java 44.8 44.9 +0.0
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/TestConventionsArchTest.java parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/TestConventionsArchTest.java 89.6 89.6 +0.0
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/TestConventionsArchTest.java parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/TestConventionsArchTest.java 89.6 89.6 +0.0

...and 12 more

Full similarity report
parallel-consumer-core/src/main/java/io/confluent/csid/utils/Java8StreamUtils.java

📄 parallel-consumer-core/src/main/java/io/confluent/csid/utils/Java8StreamUtils.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/csid/utils/JavaUtils.java 35.35
parallel-consumer-core/src/test/java/io/confluent/csid/utils/CollectionUtils.java 33.28
parallel-consumer-core/src/main/java/io/confluent/csid/utils/JavaUtils.java

📄 parallel-consumer-core/src/main/java/io/confluent/csid/utils/JavaUtils.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/csid/utils/CollectionUtils.java 39.54
parallel-consumer-core/src/main/java/io/confluent/csid/utils/Java8StreamUtils.java 35.35
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerException.java 54.2 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java 39.9
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PCRetriableException.java 36.68
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java 36.57
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java 34.95
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java 34.72
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 34.72
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/OffsetDecodingError.java 34.65
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelEoSStreamProcessor.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelEoSStreamProcessor.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelStreamProcessor.java 61.15 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessor.java 54.83 ⚠️
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelEoSStreamProcessor.java 40.41
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelStreamProcessor.java 37.13
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContextInternal.java 32.75
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelStreamProcessor.java 31.49
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelStreamProcessor.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelStreamProcessor.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelEoSStreamProcessor.java 61.15 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessor.java 50.85 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelStreamProcessor.java 36.86
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelEoSStreamProcessor.java 32.51
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContextInternal.java 31.74
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelStreamProcessor.java 30.28
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PCRetriableException.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PCRetriableException.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java 36.68
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java 33.0
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerException.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerException.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java 54.2 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java 52.94 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java 44.56
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java 34.62
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/OffsetDecodingError.java 34.16
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalRuntimeException.java 30.85
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java 30.42
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 30.42
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerOptions.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerOptions.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/ProducerManager.java 32.41
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/PartitionState.java 30.06
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessor.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessor.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelEoSStreamProcessor.java 54.83 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelStreamProcessor.java 50.85 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelStreamProcessor.java 45.54
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContextInternal.java 33.98
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/AbstractParallelEoSStreamProcessor.java 32.95
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/TestParallelEoSStreamProcessor.java 31.13
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelStreamProcessor.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelStreamProcessor.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessor.java 45.54
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelEoSStreamProcessor.java 37.13
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelStreamProcessor.java 36.86
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelStreamProcessor.java 32.7
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContext.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContext.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/RecordContextInternal.java 35.52
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContextInternal.java 32.21
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContextInternal.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContextInternal.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessor.java 33.98
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/RecordContextInternal.java 33.15
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelEoSStreamProcessor.java 32.75
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContext.java 32.21
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelStreamProcessor.java 31.74
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/RecordContextInternal.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/RecordContextInternal.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContext.java 35.52
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PollContextInternal.java 33.15
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/AbstractParallelEoSStreamProcessor.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/AbstractParallelEoSStreamProcessor.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/BrokerPollSystem.java 33.28
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessor.java 32.95
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/BrokerPollSystem.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/BrokerPollSystem.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/AbstractParallelEoSStreamProcessor.java 33.28
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/ExternalEngine.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/ExternalEngine.java

File Similarity (%)
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/VertxParallelEoSStreamProcessor.java 39.52
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java 60.36 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerException.java 52.94 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/OffsetDecodingError.java 52.13 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/NoEncodingPossibleException.java 48.44
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java 39.9
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalRuntimeException.java 39.68
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java 36.9
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java 33.33
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 33.33
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/PCRetriableException.java 33.0
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalRuntimeException.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalRuntimeException.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java 39.68
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java 31.51
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerException.java 30.85
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/PCModule.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/PCModule.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/PCModuleTestEnv.java 32.24
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/ProducerManager.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/ProducerManager.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerOptions.java 32.41
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/ProducerManagerTest.java 30.5
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java 51.36 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java 37.08
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 37.08
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java 36.9
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java 34.95
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerException.java 34.62
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/OffsetDecodingError.java 33.33
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java 60.36 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java 51.36 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/OffsetDecodingError.java 49.42
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java 46.74
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 46.74
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/NoEncodingPossibleException.java 46.23
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerException.java 44.56
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java 36.57
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalRuntimeException.java 31.51
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/NoEncodingPossibleException.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/NoEncodingPossibleException.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java 48.44
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java 46.23
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/OffsetDecodingError.java 46.12
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java 37.44
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 37.44
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/OffsetDecodingError.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/OffsetDecodingError.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java 52.13 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java 49.42
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/NoEncodingPossibleException.java 46.12
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java 34.65
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerException.java 34.16
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java 33.33
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java 31.8
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 31.8
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java 63.05 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java 46.74
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/NoEncodingPossibleException.java 37.44
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java 37.08
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java 34.72
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java 33.33
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/OffsetDecodingError.java 31.8
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerException.java 30.42
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV2EncodingNotSupported.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/RunLengthV1EncodingNotSupported.java 63.05 ⚠️
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/EncodingNotSupportedException.java 46.74
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/NoEncodingPossibleException.java 37.44
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/BitSetEncodingNotSupportedException.java 37.08
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ExceptionInUserFunctionException.java 34.72
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/InternalException.java 33.33
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/offsets/OffsetDecodingError.java 31.8
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerException.java 30.42
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/PartitionState.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/PartitionState.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/PartitionStateManager.java 30.25
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelConsumerOptions.java 30.06
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/PartitionStateManager.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/PartitionStateManager.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/WorkManager.java 39.64
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/PartitionState.java 30.25
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/ProcessingShard.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/ProcessingShard.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/ShardManager.java 36.7
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/ShardManager.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/ShardManager.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/ProcessingShard.java 36.7
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/WorkManager.java

📄 parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/WorkManager.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/state/PartitionStateManager.java 39.64
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/BrokerIntegrationTest.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/BrokerIntegrationTest.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/state/LatestResetTailNudgeIT.java 30.48
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/DrainingMemberRebalanceIT.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/DrainingMemberRebalanceIT.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/BrokerPollSystemDrainTest.java 31.17
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/KafkaSanityTests.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/KafkaSanityTests.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/csid/utils/LoopingResumingIteratorTest.java 34.18
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/MultiInstanceHighVolumeTest.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/MultiInstanceHighVolumeTest.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/VeryLargeMessageVolumeTest.java 55.41 ⚠️
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/TransactionAndCommitModeTest.java 46.98
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/MultiInstanceRebalanceTest.java 38.65
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/MultiInstanceRebalanceTest.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/MultiInstanceRebalanceTest.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/VeryLargeMessageVolumeTest.java 44.18
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/TransactionAndCommitModeTest.java 41.14
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/MultiInstanceHighVolumeTest.java 38.65
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/RebalanceEoSDeadlockTest.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/RebalanceEoSDeadlockTest.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/RebalanceTest.java 36.36
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/RebalanceTest.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/RebalanceTest.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/RebalanceEoSDeadlockTest.java 36.36
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/TransactionAndCommitModeTest.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/TransactionAndCommitModeTest.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/VeryLargeMessageVolumeTest.java 60.71 ⚠️
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/MultiInstanceHighVolumeTest.java 46.98
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/MultiInstanceRebalanceTest.java 41.14
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/TransactionTimeoutsTest.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/TransactionTimeoutsTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/ProducerManagerTest.java 30.29
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/VeryLargeMessageVolumeTest.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/VeryLargeMessageVolumeTest.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/TransactionAndCommitModeTest.java 60.71 ⚠️
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/MultiInstanceHighVolumeTest.java 55.41 ⚠️
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/MultiInstanceRebalanceTest.java 44.18
parallel-consumer-vertx/src/test-integration/java/io/confluent/parallelconsumer/vertx/integrationTests/VertxConcurrencyIT.java 39.09
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/AbstractRevokeUnderWorkScenario.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/AbstractRevokeUnderWorkScenario.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosChurnStormIT.java 48.78
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosRevokeUnderWorkIT.java 35.66
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosChurnStormIT.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosChurnStormIT.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/AbstractRevokeUnderWorkScenario.java 48.78
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosScenarioBase.java 38.04
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosRevokeUnderWorkIT.java 30.0
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosRevokeUnderWorkCooperativeIT.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosRevokeUnderWorkCooperativeIT.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosRevokeUnderWorkIT.java 49.01
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosRevokeUnderWorkIT.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosRevokeUnderWorkIT.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosRevokeUnderWorkCooperativeIT.java 49.01
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/AbstractRevokeUnderWorkScenario.java 35.66
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosChurnStormIT.java 30.0
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosScenarioBase.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosScenarioBase.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/chaostests/ChaosChurnStormIT.java 38.04
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/state/LatestResetTailNudgeIT.java

📄 parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/state/LatestResetTailNudgeIT.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/BrokerIntegrationTest.java 30.48
parallel-consumer-core/src/test/java/io/confluent/csid/utils/CollectionUtils.java

📄 parallel-consumer-core/src/test/java/io/confluent/csid/utils/CollectionUtils.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/csid/utils/JavaUtils.java 39.54
parallel-consumer-core/src/main/java/io/confluent/csid/utils/Java8StreamUtils.java 33.28
parallel-consumer-core/src/test/java/io/confluent/csid/utils/LoopingResumingIteratorTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/csid/utils/LoopingResumingIteratorTest.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/KafkaSanityTests.java 34.18
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/BatchTestBase.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/BatchTestBase.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CoreBatchTest.java 30.63
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CheckQuarantineOwnersScriptTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CheckQuarantineOwnersScriptTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/QuarantineLaneReportScriptTest.java 44.98
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/QuarantineRegistryScriptTest.java 44.14
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CommitRejectionTestBase.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CommitRejectionTestBase.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerTest.java 32.39
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerCommitTimeoutTest.java 32.34
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CoreBatchTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CoreBatchTest.java

File Similarity (%)
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorBatchTest.java 52.06 ⚠️
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyBatchTest.java 50.76 ⚠️
parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/VertxBatchTest.java 44.86
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/BatchTestBase.java 30.63
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerCommitTimeoutTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerCommitTimeoutTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerEarlyCloseTest.java 70.17 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerTest.java 56.45 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerSaslAuthenticationTest.java 49.39
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CommitRejectionTestBase.java 32.34
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerEarlyCloseTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerEarlyCloseTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerCommitTimeoutTest.java 70.17 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerTest.java 54.88 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerSaslAuthenticationTest.java 52.84 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerSaslAuthenticationTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerSaslAuthenticationTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerEarlyCloseTest.java 52.84 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerCommitTimeoutTest.java 49.39
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerTest.java 46.71
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerCommitTimeoutTest.java 56.45 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerEarlyCloseTest.java 54.88 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/MockConsumerSaslAuthenticationTest.java 46.71
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CommitRejectionTestBase.java 32.39
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/ParallelEoSSStreamProcessorRebalancedTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/ParallelEoSSStreamProcessorRebalancedTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessorTest.java 34.8
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessorTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessorTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/ParallelEoSSStreamProcessorRebalancedTest.java 34.8
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/QuarantineLaneReportScriptTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/QuarantineLaneReportScriptTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CheckQuarantineOwnersScriptTest.java 44.98
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/QuarantineRegistryScriptTest.java 33.36
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/QuarantineRegistryScriptTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/QuarantineRegistryScriptTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CheckQuarantineOwnersScriptTest.java 44.14
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/QuarantineLaneReportScriptTest.java 33.36
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/TestConventionsArchTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/TestConventionsArchTest.java

File Similarity (%)
parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/TestConventionsArchTest.java 90.27 ⚠️
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/TestConventionsArchTest.java 89.6 ⚠️
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/TestConventionsArchTest.java 89.6 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/BrokerPollSystemDrainTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/BrokerPollSystemDrainTest.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/DrainingMemberRebalanceIT.java 31.17
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/ExceptionConstructorsTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/ExceptionConstructorsTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/InternalRuntimeExceptionTest.java 30.04
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/InternalRuntimeExceptionTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/InternalRuntimeExceptionTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/ExceptionConstructorsTest.java 30.04
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/PCModuleTestEnv.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/PCModuleTestEnv.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/PCModule.java 32.24
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/ProducerManagerTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/ProducerManagerTest.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/ProducerManager.java 30.5
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/TransactionTimeoutsTest.java 30.29
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/TestParallelEoSStreamProcessor.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/internal/TestParallelEoSStreamProcessor.java

File Similarity (%)
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelEoSStreamProcessor.java 31.13
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/offsets/OffsetEncodingBackPressureTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/offsets/OffsetEncodingBackPressureTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/offsets/OffsetEncodingBackPressureUnitTest.java 39.86
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/offsets/OffsetEncodingBackPressureUnitTest.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/offsets/OffsetEncodingBackPressureUnitTest.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/offsets/OffsetEncodingBackPressureTest.java 39.86
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/truth/CommitHistorySubject.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/truth/CommitHistorySubject.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/truth/LongPollingMockConsumerSubject.java 36.45
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/truth/LongPollingMockConsumerSubject.java

📄 parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/truth/LongPollingMockConsumerSubject.java

File Similarity (%)
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/truth/CommitHistorySubject.java 36.45
parallel-consumer-mutiny/src/main/java/io/confluent/parallelconsumer/mutiny/MutinyProcessor.java

📄 parallel-consumer-mutiny/src/main/java/io/confluent/parallelconsumer/mutiny/MutinyProcessor.java

File Similarity (%)
parallel-consumer-reactor/src/main/java/io/confluent/parallelconsumer/reactor/ReactorProcessor.java 52.18 ⚠️
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyBatchTest.java

📄 parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyBatchTest.java

File Similarity (%)
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorBatchTest.java 79.0 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CoreBatchTest.java 50.76 ⚠️
parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/VertxBatchTest.java 49.08
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyPCTest.java

📄 parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyPCTest.java

File Similarity (%)
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorPCTest.java 71.17 ⚠️
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyTest.java

📄 parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyTest.java

File Similarity (%)
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorTest.java 32.49
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyUnitTestBase.java

📄 parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyUnitTestBase.java

File Similarity (%)
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorUnitTestBase.java 32.24
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/TestConventionsArchTest.java

📄 parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/TestConventionsArchTest.java

File Similarity (%)
parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/TestConventionsArchTest.java 91.06 ⚠️
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/TestConventionsArchTest.java 90.39 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/TestConventionsArchTest.java 89.6 ⚠️
parallel-consumer-reactor/src/main/java/io/confluent/parallelconsumer/reactor/ReactorProcessor.java

📄 parallel-consumer-reactor/src/main/java/io/confluent/parallelconsumer/reactor/ReactorProcessor.java

File Similarity (%)
parallel-consumer-mutiny/src/main/java/io/confluent/parallelconsumer/mutiny/MutinyProcessor.java 52.18 ⚠️
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorBatchTest.java

📄 parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorBatchTest.java

File Similarity (%)
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyBatchTest.java 79.0 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CoreBatchTest.java 52.06 ⚠️
parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/VertxBatchTest.java 50.32 ⚠️
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorPCTest.java

📄 parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorPCTest.java

File Similarity (%)
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyPCTest.java 71.17 ⚠️
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorTest.java

📄 parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorTest.java

File Similarity (%)
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyTest.java 32.49
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorUnitTestBase.java

📄 parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorUnitTestBase.java

File Similarity (%)
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyUnitTestBase.java 32.24
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/TestConventionsArchTest.java

📄 parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/TestConventionsArchTest.java

File Similarity (%)
parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/TestConventionsArchTest.java 91.06 ⚠️
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/TestConventionsArchTest.java 90.39 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/TestConventionsArchTest.java 89.6 ⚠️
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelEoSStreamProcessor.java

📄 parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelEoSStreamProcessor.java

File Similarity (%)
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/VertxParallelEoSStreamProcessor.java 41.51
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelEoSStreamProcessor.java 40.41
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelStreamProcessor.java 39.9
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/VertxParallelStreamProcessor.java 35.45
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelStreamProcessor.java 32.51
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelStreamProcessor.java

📄 parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelStreamProcessor.java

File Similarity (%)
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelEoSStreamProcessor.java 39.9
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/VertxParallelStreamProcessor.java 39.37
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/ParallelStreamProcessor.java 32.7
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelEoSStreamProcessor.java 31.49
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/JStreamParallelStreamProcessor.java 30.28
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/VertxParallelEoSStreamProcessor.java

📄 parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/VertxParallelEoSStreamProcessor.java

File Similarity (%)
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/VertxParallelStreamProcessor.java 41.67
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelEoSStreamProcessor.java 41.51
parallel-consumer-core/src/main/java/io/confluent/parallelconsumer/internal/ExternalEngine.java 39.52
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/VertxParallelStreamProcessor.java

📄 parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/VertxParallelStreamProcessor.java

File Similarity (%)
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/VertxParallelEoSStreamProcessor.java 41.67
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelStreamProcessor.java 39.37
parallel-consumer-vertx/src/main/java/io/confluent/parallelconsumer/vertx/JStreamVertxParallelEoSStreamProcessor.java 35.45
parallel-consumer-vertx/src/test-integration/java/io/confluent/parallelconsumer/vertx/integrationTests/VertxConcurrencyIT.java

📄 parallel-consumer-vertx/src/test-integration/java/io/confluent/parallelconsumer/vertx/integrationTests/VertxConcurrencyIT.java

File Similarity (%)
parallel-consumer-core/src/test-integration/java/io/confluent/parallelconsumer/integrationTests/VeryLargeMessageVolumeTest.java 39.09
parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/TestConventionsArchTest.java

📄 parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/TestConventionsArchTest.java

File Similarity (%)
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/TestConventionsArchTest.java 91.06 ⚠️
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/TestConventionsArchTest.java 91.06 ⚠️
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/TestConventionsArchTest.java 90.27 ⚠️
parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/VertxBatchTest.java

📄 parallel-consumer-vertx/src/test/java/io/confluent/parallelconsumer/vertx/VertxBatchTest.java

File Similarity (%)
parallel-consumer-reactor/src/test/java/io/confluent/parallelconsumer/reactor/ReactorBatchTest.java 50.32 ⚠️
parallel-consumer-mutiny/src/test/java/io/confluent/parallelconsumer/mutiny/MutinyBatchTest.java 49.08
parallel-consumer-core/src/test/java/io/confluent/parallelconsumer/CoreBatchTest.java 44.86

@github-actions

github-actions Bot commented Aug 10, 2026 •

Copy link
Copy Markdown

⚠️ SpotBugs Report

466 bug(s) found (rule-level exclusions only - see docs/inflight/static-spotbugs-rule-registry.md). See the annotations on the Files Changed tab for details.

astubbs and others added 2 commits August 10, 2026 21:22
…Jackson, share the streams harness

Three CI failures on #271, all in this module.

The proof tests started the topology before replaying the fixture, and
recorded each record's produced offset only after its send was acked. A
RUNNING StreamThread can fetch and process a record inside that window, so
the probe found no offset for a record it was holding and reported one "the
test never sent". Widest on the first record, where an idle poll returns the
instant it lands - which is exactly where it failed. Producing before the
topology starts closes it by construction rather than by timing; the arms
already set auto.offset.reset=earliest, so nothing is missed. Both proof
tests had it; the shared helper now states the ordering requirement.

jackson-databind 2.16.2 carries two high-severity PolymorphicTypeValidator
bypasses, first patched at 2.18.8. The declaration is compile scope, so that
version is what reaches anyone depending on this module. Neither advisory is
reachable from Kafka's metadata serialisation here, but shipping a flagged
version to save a patch bump is not a trade worth making.

Five tests each carried their own copy of build-props/start/await-RUNNING.
Extracted to BrokerStreamsIntegrationTest; callers now add only the terms
they actually vary. The duplication with parallel-consumer-example-streams
stays: those are the two arms of a controlled comparison and sharing code
between them would destroy the experiment.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018rW6vipDQY2SKwk6kudnJn
Thirteen learnings from the parallel-consumer-streams work, split between
method and domain.

Method: chasing refuted predictions (three refutations, three real finds -
a zero test-delta that turned out to be a fact about what the measurement
could see, two product defects, and the poll-wait throttle); why an author
cannot review same-session work; control arms varying exactly one term;
choosing the statistic that states the claim rather than the one that sounds
rigorous; and why status words belong in status artefacts, never identifiers.

Domain, the part a future reader of this module actually needs: Kafka
Streams polls and processes on one thread, so poll.ms is free to stock and a
throttle to anything asynchronous; a task lifecycle callback's contract is
its call sites, not its name, and postCommit is reached after a swallowed
commit failure; a high-water mark cannot express out-of-order completion, so
the fix is a deletion; one owner per metadata field with an opaque rider for
the other; a progress signal must count work consumed, not work accepted; how
to patch a dependency at build time without vendoring it, and the four checks
that make that trustworthy; a restart assertion satisfiable by pre-crash data
proves nothing; and how to re-cut history non-interactively and prove the
tree is unchanged.

Each is grounded against the tree with file:line and cites PR numbers rather
than SHAs, since re-cutting this branch is precisely why its SHAs no longer
exist.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018rW6vipDQY2SKwk6kudnJn
@claude

claude Bot commented Aug 10, 2026 •

Copy link
Copy Markdown

Claude is reviewing this PR

Task list

  • Gather context (PR diff, changed files)
  • Understand the request — run /code-review:code-review astubbs/parallel-consumer/pull/271 --comment per repo's custom instructions
  • Run code review across the diff via the code-review skill
  • Post inline comments and summary findings

Review in progress — this comment will be updated once complete.

astubbs and others added 2 commits August 10, 2026 21:34
…he 14 tests now mean

U9 landed, so three artefacts were still telling users this module can lose
records on a crash. It cannot: commit data comes from PC's completion
tracking, so the committed offset is the frontier and a commit cannot cover
work still in flight.

The plan's shortcoming entry is retired rather than deleted, because the
interesting part is what did NOT retire with it. Those 14 StreamTaskTest
cases were predicted to go green with U9 and did not - the measured delta was
zero. They assert Kafka's encoding of the commit metadata, which KTD-S7 has
this module owning wholesale, so they detect the deliberate divergence rather
than any loss. A future reader counting the 33 failures needs those 14
accounted for, not silently dropped.

The module README said the same thing to users, in the summary they are told
to read before relying on this. Corrected, with the encoding point stated so
the remaining failures do not read as unfixed data loss.

Also corrects pcProcess's @return, which still documented the pre-fix
contract ("handed to the worker pool") above a body that returns records
consumed - the sentence the next implementor would have read. Regenerated
through the unpack-edit-regen cycle rather than hand-edited, since editing
the patch directly changes a hunk's line count without its header and breaks
git apply. 30 hunks before and after, 653 -> 657 lines.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018rW6vipDQY2SKwk6kudnJn
…e pool

dispatchAvailable's return value is the progress signal the patched
process() hands back, and stock's TaskExecutor reads a false as "no progress"
and returns the StreamThread to a blocking poll. Counting pool submissions
rather than records consumed was a real defect, fixed earlier - but nothing
in this class could have caught it, because every existing WorkPreparer here
returns a non-null Runnable and none throws. On that route "consumed" and
"dispatched" are numerically identical, so the whole suite agreed with the
broken definition.

Covers both no-pool routes separately, since they are separate branches with
separate bookkeeping - a drop completes the work, a preparation failure
records a failure - and a regression could hit one without the other.

Verified as a control arm rather than assumed: with the pre-fix definition
restored (count only pool submissions) both new tests fail; with the fix in
place both pass, and the module stays at 25 green plus Kafka's 188. A test
that passes either way would have proved nothing.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018rW6vipDQY2SKwk6kudnJn
astubbs added a commit that referenced this pull request Aug 10, 2026
…mit surface

collectCommitData, hasCommitDataOutstanding and onCommitSuccess already said
"StreamThread only" in their javadoc. This turns that comment into a check,
because the cost of getting it wrong is silent: WorkManager and its partition
state are not thread-safe, so a second thread calling in corrupts offset
bookkeeping without throwing, and the damage surfaces later as a commit covering
work that never completed.

Worth doing now, and specifically on this surface. Those three methods reach
WorkManager directly and are exactly where a caller wiring up its own commit
path - which the Connect work is about to do - would naturally call from a
commit or scheduler thread rather than the owner.

Deliberately narrow: registerRecords, dispatchAvailable and close are NOT
guarded. They are hot-path, existing callers may legitimately drive them from
elsewhere, and breaking a working suite to enforce a rule nothing currently
violates is a bad trade. The commit surface is new, so it carries no such
history.

Throws IllegalStateException naming both threads rather than using assert, since
assertions vanish without -ea and this failure class produces no symptom at the
call site.

Verified it does not false-positive: PcTaskDispatcherTest 12/12 with the guard
in place, core and streams both green. That was the result that mattered - a
guard on a shared module is only worth having if that module's own suite agrees
the rule was already true.

Lives on the Connect branch because that is where the need was found, but it
belongs to the Streams module and is a clean single-file cherry-pick for
#271.

Refs #240
Refs #255

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_015GyZmjA6k1Fi9YQWavenke
astubbs and others added 4 commits August 10, 2026 23:16
…mit surface

collectCommitData, hasCommitDataOutstanding and onCommitSuccess already said
"StreamThread only" in their javadoc. This turns that comment into a check,
because the cost of getting it wrong is silent: WorkManager and its partition
state are not thread-safe, so a second thread calling in corrupts offset
bookkeeping without throwing, and the damage surfaces later as a commit covering
work that never completed.

Worth doing now, and specifically on this surface. Those three methods reach
WorkManager directly and are exactly where a caller wiring up its own commit
path - which the Connect work is about to do - would naturally call from a
commit or scheduler thread rather than the owner.

Deliberately narrow: registerRecords, dispatchAvailable and close are NOT
guarded. They are hot-path, existing callers may legitimately drive them from
elsewhere, and breaking a working suite to enforce a rule nothing currently
violates is a bad trade. The commit surface is new, so it carries no such
history.

Throws IllegalStateException naming both threads rather than using assert, since
assertions vanish without -ea and this failure class produces no symptom at the
call site.

Verified it does not false-positive: PcTaskDispatcherTest 12/12 with the guard
in place, core and streams both green. That was the result that mattered - a
guard on a shared module is only worth having if that module's own suite agrees
the rule was already true.

Lives on the Connect branch because that is where the need was found, but it
belongs to the Streams module and is a clean single-file cherry-pick for
#271.

Refs #240
Refs #255

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_015GyZmjA6k1Fi9YQWavenke
…option weighed

Gates publication rather than development, so it is written down and parked.
Kept as a decision record, not a conclusion, because the options were argued
down over several passes and the losing ones deserve to stay lost.

The crux is the coordinates, not the fork. Publishing a whole patched jar
under our own groupId does NOT remove the collision: Maven dedupes on
groupId:artifactId, so anything pulling kafka-streams transitively puts both
jars back on the classpath. Same coordinates with a version qualifier makes
nearest-wins do the work with no action from the user, which is why Confluent
ships -ccs from its own repository rather than a renamed artifact.

That makes a self-hosted repo the enabling move rather than a compromise, and
it is the only option where a user who never reads our documentation still
ends up with a correct classpath. Central stays available later under the
our-groupId coordinates; the two are not exclusive.

Records why the two non-starters are non-starters, so they are not
re-proposed: relocation destroys the mechanism it would protect, since
shadowing works precisely because the package matches; and alpha-only,
never-transitive is a promise rather than a property.

Build feasibility checked rather than assumed - the sources jar carries all
678 files including the pre-generated message classes, so no Kafka build
system is needed. Trademark, CVE ownership and the version matrix are left
explicitly open.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018rW6vipDQY2SKwk6kudnJn
…owner-thread rebind

regen-patch.sh's own header told you to run generate-sources, which only
unpacks - the patch is applied at process-sources. Follow the comment and
target/kafka-patched is pristine, so edits land on unpatched sources and
regenerating from that tree deletes every existing hunk. Two agents followed
it and hit exactly that; only the hunk-count tripwire caught it. The script
that exists to prevent a foot-gun was documenting one. Also records the
required `.` in -pl, since selecting the leaf module alone fails the
enforcer, and corrects the same claim in apply-patch.sh with the reason the
phases are split.

Separately, records a seventh item on the task-lifecycle list: the
owner-thread guard binds at construction, which task recycling falsifies. A
reassigned task carries a stale owner and the guard then throws on a
legitimate call. Wants an explicit bind at the point a task is handed to a
thread - the same seam as the recycle leak, so one visit fixes both.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018rW6vipDQY2SKwk6kudnJn
Enables the cheap way to reconcile two branches that both regenerate
pc-streams.patch: commit the full generated Kafka Java on each branch, let git
merge it as ordinary source, re-derive the patch, drop the files again.
Merging the generated Java is trivial where merging the patch is arithmetic on
hunk headers - on #255 the same change had zero conflicts as Java and
eight conflict blocks as a patch.

Those files carry Apache's headers and must keep them. They are Apache Kafka's
code - neither fork-original nor upstream-Confluent-derived - so every rule the
gate applies is the wrong question to ask of them, and without the skip the
technique fails on its first commit, which is precisely when someone would
conclude it does not work. Redistribution stays covered by NOTICE under Apache
2.0 section 4(b), which is the right instrument for it.

Verified with a control arm rather than by inspection: an Apache-headered file
inside the excluded path leaves the count at 253 and the gate green, while the
same file outside it is checked and fails. The skip is doing the work, not the
file being innocuous.

Also revises regen-patch.sh's guidance: parallel branches are a speed win, so
reconcile at the source rather than serialising to dodge conflicts, and a
conflict here is worth having - when the two #255 branches met, the
merge exposed a test asserting StreamThread loads from the kafka-streams jar,
which the other branch had just made false.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018rW6vipDQY2SKwk6kudnJn
@astubbs

astubbs commented Aug 31, 2026

Copy link
Copy Markdown
Owner Author

Unit U10 from this study - task lifecycle, rebalance and the commit frontier - has been cut onto the
reconstructed Wagon B stack as #394, stacked on
#389.

It answers two of the review threads opened here.

"revive() is permanently broken for a task after any dirty close, not just fatal ones." The
finding was right and the blast radius was wider than "revival is not supported yet" suggested. A
revived task now builds a new dispatcher over the partitions it holds at that moment, rather than
throwing. Measured against Apache Kafka's own suite with the seam on, before and after, one term
changed: shouldRecoverFromInvalidOffsetExceptionOnRestoreAndFinishRestore green on every parameter,
no exception leaving a StreamThread where three did, and nothing that passed before regressing.

"close() is missing the owner-thread guard its siblings enforce" and "registerRecords /
dispatchAvailable lack the assertOwnerThread guard their commit-side siblings have."
Both are now
answered by documentation rather than by adding guards, and the class javadoc states the position
instead of leaving it implied: the surface is split into a guard-enforced row and a
convention-only row, the convention-only row is described as just as unsafe off-thread rather
than as safe, and adding guards to it is called out as a behaviour change worth deciding separately.
The cross-thread hazard that actually bit - DefaultStateUpdater calling maybeCheckpoint from its
own thread - is handled instead by making the outstanding-work question a genuine query that touches
neither the WorkManager nor the completion mailbox.

One correction to something this study's own draft claimed: bindToCurrentThread() has no production
caller.
In Kafka 3.9.2 a reassigned task is closed and rebuilt rather than handed across threads, so
the constructor's bind is the only one that happens. It is capability with an unstated assumption
removed, not a fix, and the javadoc and a docs/inflight/ note both say so now.

Two threads here are explicitly not answered by that PR and are named in its body as belonging to
the error-surfacing unit: "a worker's processing failure can be committed past", and the same
asynchrony seen from the recovery side - a TaskCorruptedException raised inside a processor arrives
wrapped and late, so Kafka's recovery never fires. That second one is what the dispatch default is now
waiting on, with its own measurement recorded.

@astubbs

astubbs commented Aug 31, 2026

Copy link
Copy Markdown
Owner Author

Two of this PR's open review threads now have work behind them, on a rung reconstructed from this
branch: #395, stacked on #394.

"No backpressure on the PC dispatch path." The pause is applied from the same
buffered.records.per.partition stock uses, with the occupancy read from Parallel Consumer, and the
resume mirrors the pause exactly rather than being computed from the whole assignment - so it cannot
hand back a partition Kafka paused for an offset reset, and it is provably no more eager than the
pause, which is the property the bound rests on. numBuffered() and the active-buffer-count metric
stop reporting zero.

The bound is measured with a control arm rather than asserted: two arms differing only in
pc.streams.backpressure.enabled, on the same broker, topology, data, processing cost and JVM, with
occupancy sampled from a watcher thread because the interesting quantity is a peak and a peak is
invisible to any assertion made after the run. Over a 600-record backlog with a processor slower than
the feed: 596 held with the pause off against 30 with it on, the fixed arm landing exactly on its
derived bound.

One thing worth flagging beyond the thread as written. The occupancy is derived from PC's own
incomplete-offset set rather than counted alongside it, which reverses a decision the original work
took deliberately. A count raised by predicting what PC will accept drifts the moment PC applies a rule
the prediction does not model - a re-delivered offset, for instance, which ProcessingShard drops
because it already holds a live container for it. That record is never handed out, nothing decrements
it back, and the partition is paused for good with no symptom but a topology that has gone quiet.

"A worker's processing failure can be committed past." Established as a property. A pending failure
fences commit-data collection, so nothing commits past the record that failed; it deliberately does not
fence the "is there uncommitted work" question, which is the opposite one - that is what
validateClean turns into a TaskMigratedException so the task closes dirty, and fencing it would
make a failed task look clean to close. Both directions are pinned by tests, and removing the fence was
sabotaged with the prediction stated first: exactly one case red, its control arm green.

Related, and on the same rung: a control-flow exception raised inside a processor now reaches Kafka's
TaskManager with its type intact and from the same runOnce that ran the record, which is what
Kafka's own recovery test asserts. Fixing only the type would have been unverifiable rather than
partial - with the timing untouched, no test can tell the unwrap from its absence.

Not posting into the threads themselves; this is a pointer so whoever reads them next knows where the
work is. The seam's default is still OFF, and #395 deliberately does not move
it.

@astubbs

astubbs commented Aug 31, 2026

Copy link
Copy Markdown
Owner Author

Stream time and punctuation have been extracted from this study into a reviewable rung.

Where it went, and how it relates to the open threads here:

  • The STREAM_TIME punctuation thread on this PR is answered. That thread predicted "a common,
    non-windowed pattern ... silently never fires" before any of the refusal work existed, and it was
    right. Stream time now advances on the dispatch path as a low-water mark over work in flight, so
    punctuators fire - and the public ProcessorContext.currentStreamTimeMs(), which returned a constant
    -1 on that path and is reachable from any processor with no refused DSL call in sight, returns a real
    value.
  • The handover's ranked defect 1 is decided, not deferred. Its hasUncommittedWork() || commitNeeded
    one-liner is not taken: measurement shrank the claim behind it to a bounded idle-window replay cost,
    while the fix would change commit cadence for every caller on that path. What is taken is the other
    half - both punctuation types now warn once per task at registration, WALL_CLOCK_TIME being the new
    one, and the warning is gated by a test with a seam-off control.
  • Two divergences are confirmed and pinned rather than closed: the mark can sit above what stock holds
    where dispatch order differs from timestamp order, and after a restart against a group this module
    committed, punctuators re-fire over event time an earlier run already covered.

The default does not move. One of its named triggers is closed by this work and the reasoning is
recorded where the decision lives, but the exception-surfacing trigger is still open on a sibling rung, so
the flip is left for after those rungs reconcile.

This PR remains the source study and keeps the measurements, the semantic work and the open review threads;
nothing here is closed by the extraction.

@astubbs

astubbs commented Aug 31, 2026

Copy link
Copy Markdown
Owner Author

Wagon B, unit B6 (the evidence suite) is extracted and open as #398, stacked on the task-lifecycle rung #394.

It carries the three evidence workstreams this PR still holds, re-derived on this stack rather than quoted:

  • The seam-on upstream lane. Kafka's own suite now runs twice - seam off, then seam on - and every case that passes in the first and fails in the second is classified by mechanism. The refusal envelope's contribution is derived from PcUnsupportedConstruct at run time, so the envelope can gain or lose a construct with no edit in the lane; the rest is attributed in a docs/inflight/ triage note that the rung closing a mechanism edits in the same change. An unexplained divergence fails the lane, and that was proven by sabotaging a semantic and watching it land in the unexplained pile.
  • The PC-specific broker-backed laws, with their controls - including the external stock baseline generated in parallel-consumer-example-streams, which this PR's PcSupportedEnvelopeTest already cited as the reason non-windowed aggregation is not refused.
  • The realistic keyed benchmark, re-measured here with a noise floor taken from two identical stock arms in the same run. No figure from this PR's branch forest is carried into the diff.

Two corrections to this PR's own record, both with the evidence in #398's body: the earlier seam-on triage attributed the invalid-timestamps case's seam-on failure to the known flake, and it is a systematic divergence; and bin/streams-benchmark.sh selected nothing on the reconstructed module's layout while exiting 0.

Nothing in that PR touches the seam, the patch, or any main Java in the module.

Merges `feats/ks-streams-reconciled` - the reconciliation of #395,
#396, #398 and #391 on top
of the refreshed #394 - into this branch, so that
#271 can be retargeted onto that rung and its displayed diff collapse to the
residue. This is the plan's "retarget move",
`docs/plans/2026-08-31-001-process-god-branch-decomposition-plan.md`, and it is an ordinary merge: no
history was rewritten and nothing was force-pushed. The pre-merge tip is preserved as
`origin/backup/pre-stack-merge-271`.

The merge also brings `master` forward, because the stack is current with it and this branch was
not. That is most of the volume and all but a handful of the conflicts.

## Which side won, and the check that made it safe to say so

The rule was: the rungs' refined versions win wherever both sides have a file; the spike's residue is
kept. **The second half of that is only safe if the spike contributed nothing to the shared files,
and that was measured rather than assumed.** Both trees were exported with package paths normalised
(`io.confluent.parallelconsumer` and `bz.stub.parallelconsumer` folded to one token, and
`io.confluent.csid.utils` to its new home) and compared file by file against the branch's own merge
base. The answer is clean: outside `parallel-consumer-streams` and the streams example, **every
difference this branch carries is the package rename itself plus a copyright-header line** - the same
rename `master` performed independently, which is why taking the rungs' side loses nothing. The one
substantive exception is the streams example's stock-baseline fixtures, which are already on the
stack because #398 depends on them.

- **Files on both sides: the rungs' version, 93 of them.** They carry current master, the rename, and
  every fix the five rungs landed.
- **Files only on the spike: kept, 32 of them.** The plan and result documents, the branch handover,
  the ranked next-work notes, the eight `docs/solutions/` write-ups, and two streams tests the stack
  deliberately left behind - `HeadOfLineBlockingBenchmarkTest`, which
  #398 says it is leaving on the forest, and
  `ProcessorContextConfinementTest`. That residue is what this PR still exists to review.
- **Files master deleted: they stay deleted, 12 of them** - `.semaphore/`, `.travis-archived.yml`,
  `service.yml`, two `.idea` run configurations, the uppercase `docs/` names master lowercased, and
  the `InternalRuntimeException` pair master renamed to `PCInternalRuntimeException`.
- **Two shared ledgers merged as unions, not resolved to a side**:
  `docs/inflight/pr-strategy-doc-merge-triggers.md` and `docs/inflight/release-0.6.0.0.md` keep this
  branch's paragraphs alongside master's.

**The strongest check is the one that makes a mis-paired rename impossible to hide.** Rather than
counting recorded renames, the merged tree was differenced against `feats/ks-streams-reconciled`
directly: the only differences are the 32 additions and the two union-merged ledgers. Nothing else
survived on the wrong side, in the wrong module, or under an `io/confluent` path -
`git ls-files | grep io/confluent` is empty, and the prescribed `grep -rnE 'io[\./]*conflu'` sweep
returns only prose in dated documents.

**Both sides were already renamed, which is not what the plan predicted.** It expected the spike to be
pre-rename; `origin/feats/ks-on-pc-spike` had in fact moved five commits past the worktree's copy and
carries the rename. That is the good case AGENTS.md describes - the mis-pairing surfaces as
rename/rename conflicts on the right files instead of silently applying one module's edit to another -
and it did: three `TestConventionsArchTest.java` files were offered paired across modules, and the
right file was taken into each module rather than accepting git's pairing.

## One collision the merge produced, and the reason it is a good one

`HeadOfLineBlockingBenchmarkTest` carried a private `sleep(Duration)` helper. Since
#395 hoisted the identical helper onto `BrokerStreamsIntegrationTest` -
because two arms that simulate cost differently are not comparable - the two now collide, and the
compiler said so rather than either winning silently. The local copy is deleted and the shared one
inherited, which is what the hoist was for.

## The gates, and the one that is a scoping artefact

`bin/check-all.sh` was run in this worktree. Four gates failed on first pass and **three of the four
did not exist on this branch before the merge** - they arrived with master, and every finding was on
content that predates them.

- **`check-inflight-tags`** flagged sixteen notes with no `inflight-type`, all of them this branch's
  own and all older than the tag scheme. Tagged - a first pass a reviewer should correct where it
  reads wrong. `bug-core-tests-jar-junit-parallelism-leak.md` additionally carries a `closed` state,
  because #265 removed the file it is about.
- **`check-file-refs`** flagged citations broken by the rename - the exact class
  [`docs/citations.md`](docs/citations.md) says must be repaired rather than left. Forty-seven paths
  were re-pointed to their successors, **each verified to exist before the edit was written**; the
  claims around them are untouched, which is the line that document draws. What genuinely has no repo
  path - Apache Kafka's own sources inside the published jar, a placeholder in a shell recipe, a class
  that was only ever proposed - carries a line-scoped `file-refs: N/A` with its reason. Two targets
  that moved out of the tree point at the history holding them instead.
- **`check-issue-refs`** flagged two bare `#NNN` in the handover; qualified.
- **`check-branch-self-reference`** is the artefact, and the gate's own header names this exact
  situation: mid-merge the merge base is still the old one, so master's sentences read as this
  branch's. Twenty-five of the twenty-eight files it flagged are on `feats/ks-streams-reconciled` too
  and twenty-four are byte-identical to it, so marking them would be attesting to somebody else's
  notes - which that header explicitly says not to do. Committing the merge is the documented fix.

## Verified

JDK 17, macOS. The whole reactor compiles, main and test. The reconciled branch it merges was verified
in its own right before this: seam-off oracle unchanged, module unit suite green, and all 35
broker-backed integration arms green against a real broker.
…oes not

Follow-up to the stack merge, splitting `check-branch-self-reference`'s findings into the ones this
branch may answer and the ones it may not.

**Answered**: the Kafka Streams handover's branch table and PR list, which name every branch and PR
outright rather than saying "this branch" - a dated record whose whole value is that it says what was
true on 2026-08-13 - and the stall note's citation of the PR its sighting came from. Both carry a
`post-merge: checked` attestation with the reasoning, which is what that marker is for. One state line
this merge itself wrote said "this branch" and now names the spike instead.

**Not answered, and deliberately**: nine findings across five notes that are byte-identical to their
copies on `feats/ks-streams-reconciled` and therefore on `master`. Marking those would be attesting to
somebody else's note, which the gate's own header says not to do. They are visible here only because
this branch's merge base with `master` is old, so the gate's phrase arm scopes over master's own
files - the same artefact its header describes, one step further out.
…uite

`ParallelEoSStreamProcessorTest.inFlightMessagesCommittedIfProcessedDuringShutdown[3]` failed once in
the full `parallel-consumer-core` unit run on this branch - `assertCommits` expected the record that
completed during shutdown and found nothing committed.

**It is the load-tightness family's signature, and that is established rather than asserted.** Green
three times out of three in isolation. Green in a full core unit run on
`feats/ks-streams-reconciled` minutes earlier in the same session - and `parallel-consumer-core` is
**byte-identical between the two branches**, so there is no code difference available to explain the
difference in outcome. Fast-failing assertion, heavy contention, passes alone: that is the row this
family is defined by.

Recorded rather than re-run past. Nothing was loosened, retried or serialised.
…he clean third run

The full unit suite was run three times on this tree. Run one failed
`inFlightMessagesCommittedIfProcessedDuringShutdown[3]`; run two failed a **different** case in the
same class, `processInKeyOrder[3]`, and did so in its own **preamble** sanity check, before the
behaviour under test was reached; run three was **entirely green** - every module, including core's
whole unit suite and the streams module's own suite and upstream oracle.

Two different cases, on an unchanged tree, in a module this branch does not modify, one of which
already has a diagnosis on another branch showing the flake reddens unmodified `master` harder than
the branch it was seen on. Both rows go to the load-tightness family, whose signature is exactly
this: a fast-failing assertion under contention that passes alone.

The diagnosis is cited by the command that reads it rather than copied, because the note lives on a
branch this one does not carry.
@astubbs
astubbs changed the base branch from master to feats/ks-streams-reconciled September 1, 2026 04:46
astubbs added a commit that referenced this pull request Sep 1, 2026
#271 carries eleven unresolved review threads. Nine of them either close
against a named rung of this stack or are answered as a documented contract. **Two do not**, and a
thread is the worst place for a live defect to sit: it disappears with the PR.

Both are the same shape. The patched `commitNeeded()` short-circuits to `pcAwareCommitNeeded()`
whenever a dispatcher is present, and that helper asks one question - does PC hold uncommitted work.
Stock asks two more, and both of them are real: an explicit `context().commit()` sets
`commitRequested`, which nothing on this path reads, and a poll batch of nothing but control markers
advances the consumer position through a sweep this path cannot run because `consumedOffsets` is
gone by design.

**Read against the reconciled tree rather than against the state the threads were written on**, which
matters because four rungs have changed this area since - the refusal envelope, the task lifecycle,
error surfacing and stream time. `grep -n 'commitRequested' parallel-consumer-streams/src/main/patch/pc-streams.patch`
is the check that the first one still stands: it finds the declaration and the javadoc calling it
"the one sanctioned cross-thread commit-state write", and no read.

The note also says what settling them is NOT: a one-line `|| commitRequested` is the same shape as
the `hasUncommittedWork() || commitNeeded` candidate the stream-time rung measured and rejected,
because `validateClean()` reads the same answer and would start throwing on a clean close after a
punctuate-only interval.
…this PR

Merges `feats/ks-streams-reconciled` again, now that it carries
`docs/inflight/core-streams-two-commit-signals-the-pc-path-cannot-see.md` - the two review threads on
this PR that no rung of the stack answers, written down where the code they are about lives rather
than left in a thread that disappears when the PR does.

Also brings the two attestation commits that took `bin/check-all.sh` to zero failures on that branch.
Clean merge; the flake ledger auto-merged as a union of both sides' sightings.
@astubbs

astubbs commented Sep 1, 2026

Copy link
Copy Markdown
Owner Author

This PR has been retargeted onto the top of the Wagon B stack, and its diff is now the study's
record rather than its code.

feats/ks-streams-reconciled was merged into this branch as an ordinary merge, and this PR's
base changed from master to that branch. No history was rewritten and nothing was force-pushed.
The body has been rewritten to the residue framing.

The displayed diff collapsed from 158 files (+35,562) to 43 (+7,880), and what is left is 41
documents plus two tests: the plan and result write-ups, the branch handover, the ranked next-work
notes, eight docs/solutions/ learnings, HeadOfLineBlockingBenchmarkTest (which
#398 says outright it is leaving on the forest) and
ProcessorContextConfinementTest.

The stack it now sits on

#379 → #388 → #389 →
#394 → {#395, #396,
#398}, with #391 on the seam rung - all reconciled
onto feats/ks-streams-reconciled, and this PR on top. They merge bottom-up; this one merges
last.

Before the reconciliation, #379's post-cut work - the prepare-deps warm
for Kafka's sources and test-sources classifier jars - was cascaded down the spine through
#388, #389 and #394 by
ordinary merges, closing the Maven Central coin flip that had reddened whole lanes on two of those
PRs with zero tests run.

The dispatch default: measured on the reconciled tree, and it stays OFF

Every rung declined to move it on its own evidence and reserved the decision for the merged module.
The measurement was taken there.

The acceptance test ran red first, which is what makes it a test. The two divergence attributions
the sibling rungs owned - stream-time-never-advances and exception-type-lost-in-the-worker - were
deleted before the lane ran, so a fix the merge had lost would match no mechanism and be reported
as unexplained rather than quietly keep its old label. It exited 1 with three unexplained
divergences.

A control arm then told "the merge lost a fix" apart from "the mechanism only partly closed": the
same seam-on suite at #396's own tip, in a detached worktree, so the one
term that differs is the reconciliation itself. Result:
shouldReinitializeRevivedTasksInAnyState fails on all three parameters there and passes on two
here; shouldRecordBufferedRecords fails there and passes here; the two StreamThreadTest
punctuation cases pass in both. Nothing was lost. The residue changed mechanism rather than
surviving: shouldPunctuateOnceStreamTimeAfterGap went from 0 punctuations to 6 of 7, and
shouldRespectPunctuateCancellationStreamTime now fails a different assertion.

So the fourth named trigger is that STREAM_TIME punctuation lags stock by the work still in
flight
- the low-water mark doing its job, and a divergence in the direction the module's own
documentation does not state. Closing it needs no code fix, only the lagging direction recorded to
the same standard as the overtaking one. With the two mechanisms re-attributed, the lane exits 0 with
every divergence explained and no stale entry. The three sites that carry the decision were
re-pointed together, and the reservation note was deleted - its own "delete when" was the default
being re-decided against a fresh measurement, whether or not it moved.

The guards and lanes

  • Full reactor test: BUILD SUCCESS, all 13 modules, zero failures, including core's whole unit
    suite and the streams module's own suite plus Kafka's upstream oracle.
  • All 35 broker-backed streams integration arms green against a real broker on the base branch -
    the rebalance arm, both memory-bound arms, the five punctuation arms, both proof arms, the
    commit-frontier crash-restart law, the shadowed stock control and the wake-on-work shutdown arm.
  • The seam-on evidence lane exits 0 with a clean control arm.
  • bin/check-all.sh passes every runnable gate on the base branch. On this branch it exits 1 on
    check-branch-self-reference alone, with nine findings across five notes that are byte-identical
    to their copies on the base branch and therefore master's - marking those would be attesting to
    somebody else's note, which that gate's header says not to do. Three other gates failed on first
    pass and all three did not exist on this branch before the merge; every finding was on content
    older than the gate, and they are fixed - citations broken by the rename re-pointed to verified
    successors, untagged notes tagged, bare issue references qualified.
  • Two flakes met and neither chased, both in a module this branch does not modify: three
    consecutive full unit runs gave two different failures in ParallelEoSStreamProcessorTest and one
    entirely green run. Both are recorded in the load-tightness ledger.

The ten findings, mapped - and NOT replied to

The body carries a table mapping each unresolved review thread to the rung that changed the code
under it. No reply has been posted to any thread; they are yours to review.

Eight close against a named rung or are answered as a documented contract. Two do not, and both are
still open
: an explicit context().commit() is set and never read on the PC path, and a poll batch
of nothing but control markers advances the consumer position unseen. Both were re-read against the
reconciled tree rather than against the state they were written on, and are recorded in
docs/inflight/core-streams-two-commit-signals-the-pc-path-cannot-see.md with the reason a one-line
|| commitRequested is not the fix.

The backup

The pre-merge tip of this branch is preserved on the remote as backup/pre-stack-merge-271
(8c3bcecb1). Nothing about this operation needs it, but it is there.

The ks-streams-* forest branches are untouched.

astubbs added a commit that referenced this pull request Sep 1, 2026
`check-branch-self-reference` flagged the two places this note cites #271.
It is right to ask and the answer is that the citation is the point: the note exists precisely
because a review thread stops being readable when its PR closes, while a landed PR is a permanent
link. The sentences name the PR outright rather than saying "this PR", so they read the same
afterwards. Attested with `checked-begin`/`checked-end` rather than reworded.
astubbs added a commit that referenced this pull request Sep 1, 2026
… only beside its siblings

The `Integration Tests` lane on #271 went red on
`PostCommitCheckpointGapTest.pcPathAlsoRefreshesTheCheckpointUnderLoad`. Diagnosed rather than
re-run, and the first reading was wrong.

**The value is `-4`, which is a sentinel and not a small number** - the checkpoint file exists and its
changelog offset was never populated. That is a different failure from a bound missed under load, and
it does not move with more time, so "the CI runner is slow" was never available as an explanation.

**A control arm ruled out the reconciliation, which was the hypothesis worth killing first**: a
commit fence and a checkpoint expectation meeting for the first time is precisely what a
reconciliation surfaces. The same class, run alone at #396's own tip in a
detached worktree, fails 3 of 3 with the identical sentinel. The merge is not the cause.

What is left is an ordering dependency: the class passes in a single-fork local suite run, where the
siblings that ran before it warmed the JVM, and fails alone - 3 of 3 here, 3 of 3 on its own branch -
and on CI at `forkCount=4`, where a forked class does not reliably get that neighbour.

**That matters more than the red build.** The class's own assertion message says it REFUTES a sibling
class's inference that `postCommit` never runs on the PC path. A refutation that holds only when
another class ran first in the same JVM is weaker than the one written down, and the green suite is
what hid it - the shape this module's handover names as "several tests here passed while proving
nothing".

Nothing was widened, shortened or re-run into green. The note names what to look for and says
explicitly not to widen the bound, because the bound is not what is failing.
@astubbs astubbs added the 0.6.0.0 Targeted at the 0.6.0.0 release label Sep 1, 2026
@astubbs astubbs added this to the 0.6.0.0 milestone Sep 1, 2026
astubbs added a commit that referenced this pull request Sep 2, 2026
… the baseline is a lookup

Two corrections from Antony, both to the tracking-gap detector, and the first is a
false positive it produced on its very first real run.

A PR EXPLAINS MORE THAN ITS OWN HEAD BRANCH. The detector asked whether a branch had a
PR of its own, a branch note, or a mention in a note on the baseline - and answered
"tracked nowhere" for `origin/feats/ks-streams-reconciled`. It was documented all
along: #271 BASES on that branch and names it in its body,
which the detector could not see because it only read notes and head refs.

That is the cry-wolf failure - a detector wrong about a documented branch is one nobody
reads by its third report. It now treats two further signals as tracking, reported
separately because they differ in strength: a PR whose BASE is this branch explains it
exactly, by definition; a PR body that merely names it is textual and weaker.
`prsByBranch` carries `baseRefName` and `body` for this, since it already makes the
call.

BASELINING IS A TIMESTAMP LOOKUP, NOT STORED STATE - Antony's design, and better than
either option offered. The moment `bin/inflight.mjs` first appears on the baseline is
the moment tracking became expected, and git can answer that at any time:

  git log <baseline> --diff-filter=A --format=%ct -- bin/inflight.mjs | tail -1

A branch predating it was cut when nothing asked, so it reports as backlog rather than
a new gap - which is what stops the detector opening with sixty findings and being
ignored. One cut afterwards has no excuse. The grandfathered set shrinks on its own as
those branches land or die, and there is no snapshot file to rot, which was the whole
objection to the alternatives: docs/todo-index.md is this repository's cautionary tale
for a committed generated index.

AN UNKNOWN MOMENT NEVER GRANDFATHERS. If the tool has not reached the baseline yet -
true right now, since this PR has not merged - every gap still reports loudly. Silencing
everything on a lookup that returned nothing is the worst available default, and it is
the shape of every silent-pass defect this PR has already fixed twice. Its check
asserts the unknown case explicitly, and the mutant that flips it to grandfather-on-
unknown goes red.

Also records the ks-streams workstream's actual shape in
docs/inflight/branch-ks-streams-workstream.md, which existed as a signpost that named no
branches at all - so it said the workstream exists without saying how to reach it. It
now names all eight rungs with their PRs, the reconciled branch that integrates them,
and #271 above it. Found by running the tool.

26 checks, 52 assertions, green in the two-ref no-origin/master CI shape.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01UTX8obQMsjs9kq2rkpU5cZ
astubbs added a commit that referenced this pull request Sep 3, 2026
…on the gate's own configuration (#432)

The learning from #421's rejected first candidate. Running the chaos
suite under forkCount=2 went red on ChaosChurnStormIT's INSTANCE_STALL probe,
and the obvious reading was contention from the change. Two seeded replays
on the same runner class, prediction written first, said otherwise: the
two-fork arm passed and the one-fork arm - the gate exactly as it runs today
- failed on a different probe. The seed owned the failure, not the change.

The doc scopes itself to what docs/investigating.md's general control-arm
rule does not say: the control is the gate's own unchanged configuration on
the runner class that fired, not a quieter machine; one sample is not a rate
and node bin/inflight.mjs codecov test is where the recorded rate lives,
with the caveat that dispatched measurement runs never reach it; record the
seed and both arms before any re-run; a result no prediction row covered is
the finding. It cites the two siblings that exist only on the #271
branch by title, since master has no path for them yet.

CONCEPTS.md gains a flagged ambiguity: "shard" is the engine's ordering unit,
and the CI chaos lane now borrows the word for its split jobs.

Co-authored-by: Claude Fable 5.1 (1M context) <noreply@anthropic.com>
@astubbs astubbs modified the milestones: 0.6.0.0, Streams on PC Sep 7, 2026
@astubbs astubbs removed the 0.6.0.0 Targeted at the 0.6.0.0 release label Sep 17, 2026
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.

1 participant