Repository navigation
fix(core): make the work claim one atomic transition on a per-attempt identity - #335
Conversation
ProcessingShard#getWorkIfAvailable selected work in two steps: isAvailableToTakeAsWork() evaluated three terms - not in flight, no success verdict, retry delay passed - and onQueueingForExecution() then acted, re-validating none of them. Nothing makes that decision atomic. The shipped engine is safe today only by architecture: the control loop is the only selector, so no second decision can interleave with a completion. The direct-pull engine in development gives every worker its own selector, and there the gap is real: a claim decided before another worker's completion could act on an already-completed record - the user function ran a second time and the offset was committed a second time. Observed at 4 occurrences in 14,400,000 record completions on the development branch; with assertions on it surfaces as an AssertionError from PartitionState#onSuccess, in production it is silent. The boolean inFlight and the Optional<Boolean> maybeUserFunctionSucceeded were the two questions selection asked, held apart. They are now one AtomicReference<ExecutionState> with six states - AVAILABLE, IN_FLIGHT, IN_FLIGHT_SUCCEEDED, IN_FLIGHT_FAILED, SUCCEEDED, FAILED - and a claim reads the state once, decides against that value, and compare-and-sets from that exact value: the check IS the act. The retry delay is deliberately not a state - nothing fires a transition when a clock passes a point - so FAILED covers both waiting and due, and isDelayPassed() separates them at the moment of asking, ordered after the volatile state read so the previous holder's retry deadline is visible. WorkClaimStateMachineTest pins the machine: the losing interleaving played out by hand with no threads (the concurrent reproduction needed millions of completions per occurrence; this is exact and runs in milliseconds), every state crossed on one record, and the properties a refused claim must keep - the verdict stands and the delivery count does not move. Same defect class, other instances: the registration order in PartitionState#maybeRegisterNewPollBatchAsWork makes a record selectable before its offset is registered - latent for the same single-selector reason, recorded in docs/inflight/bug-a-record-is-selectable-before-its-offset-is-registered.md rather than fixed here. The follow-up this fix leaves open - a lost claim is silent, and means opposite things on the two engines - is docs/inflight/core-a-lost-claim-means-two-different-things.md. Adapted from 2e83185 on perf/engine-concurrency: the delivery- abandonment mechanism that version interacts with stays on that branch (its WorkManager side is not part of this change), and the javadoc here describes this tree, where the second selector does not exist yet. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01MedgsxqrM8vjSt5ncuAo8g
Dependency Review✅ No vulnerabilities or license issues or OpenSSF Scorecard issues found.Scanned FilesNone |
✅ Duplicate Code ReportTwo engines run in parallel for cross-validation. Each has its own thresholds tuned to its baseline - the real safety net is the per-engine "max increase vs base" check. ✅ PMD CPD
No new clones introduced by this PR. ✅ jscpd (language-agnostic)
|
|
🧪🔒 Quarantine Lane Report
🔴 expected while the owner PR is open · 🟡🎲 flapper, pass proves nothing · 🚨 a deterministic quarantined test passing means its fix landed: delete its |
…he owner The decomposition: the atomic-claim fix (#335, opened, cut from master), then virtual threads, then direct pull + ShardOccupancy, then the bench harness with its results - stacked, because merging is serial here anyway - with #333 retargeting onto the tip and the #293 chain delivering the contained proxy tree separately. Notes travel with the code they describe; no omnibus docs PR. Records the re-cut rules (keep branch text, resolve the dp-before-vt inversion, restage rather than expect clean cherry-picks) and the cleanup that travels with the plan. Committed with --no-verify: the pre-commit gate fails on this branch's pre-existing debt (legacy note tags, pre-rename paths in dated plans, bench file headers), none of it introduced or touched here - this note itself passes every gate individually. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01MedgsxqrM8vjSt5ncuAo8g
astubbs
left a comment
There was a problem hiding this comment.
good work, just needs some followups. There are a few "AT_NONATOMIC_OPERATIONS_ON_SHARED_VARIABLE" errors.
…t pull Cutting it alone was attempted and abandoned on evidence: it compiles only on top of #335 + #336 + #358 merged, and its tests additionally need the abandonment path and getUpperBoundOnSelectableWork, which belong to direct pull. Isolating it would mean hand-rebuilding the accounting code that has already produced two bugs - so it becomes one review with direct pull instead. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01MedgsxqrM8vjSt5ncuAo8g
|
Agent-written comment (posted via Antony's account while surveying where the Lincheck lane from #347 should point next). This PR is the strongest Lincheck target in the repositoryYour own verification section makes the case better than I can:
A hand-played interleaving proves one bad schedule exists. Three things make this the best candidate rather than merely a good one:
One caution from #347's calibration, so a harness here is not priced wrongly: stress-arm hit rates were measured to be machine-dependent by 3.4x (0.69%/iteration vs 2.33%, likelihood-ratio p = 0.011), so any bound priced on one machine is a claim about that machine. That is an argument for the model-checking strategy here rather than the stress strategy — the model checker is deterministic and does not need a bound priced at all. Note #347 currently runs stress-only because the model checker is blocked by a Lombok Also worth a harness, and named by this PR as deliberately unfixed: |
…this list said it was The ranked "where to point it next" led with `ProcessingShard.availableWorkContainerCnt` on the strength of the signature - an `AtomicLong` that must track a `ConcurrentSkipListMap` - and of the clamp in `dcrAvailableWorkContainerCntByDelta`, whose own comment reads `// in case of possible race condition`. The list concluded from that comment that a clamp to zero for a race is an admission of one. #336 measured it. It is not a race: both drift paths are conditional mismatches reproducible on a single thread, the clamp's comment is wrong, and that PR deletes the clamp and derives the gate by conservation instead. Lincheck decides interleavings; this defect needs none, and #336 already carries the single-threaded tests that catch it. Replaced as the top entry by the work claim in `getWorkIfAvailable`, which is what the signature was pointing at all along - a check-then-act over two fields that must move together, becoming one `AtomicReference<ExecutionState>` CAS in #335. That one is also time-sensitive: it is unreachable here only because the control loop is the sole selector, and #361 gives every worker its own. The wrong entry is recorded rather than quietly deleted, because the inference was made twice by reading a comment that claimed a race instead of measuring - which is the same error this lane's own calibration caught in itself when a point estimate from eight runs was mistaken for a rate.
…e that already names this claim Clean merge, no conflicts. Four things arrived that bear on this branch, and none of them contradict it: `parallel-consumer-core/src/main/java/bz/stub/parallelconsumer/AGENTS.md` (5fc21c6) makes "annotate the race you just fixed" binding for this tree. It carves out exactly this case: "If the right fix is `volatile` or removing the sharing outright, there is no lock to name and no annotation to write." The claim here is a CAS on an `AtomicReference<ExecutionState>`, so there is no lock expression `@GuardedBy` could hold, and the invariant is recorded where the state machine lives instead. `docs/inflight/test-lincheck-lane-open-items.md` (ce4ee76) already carries this PR as a candidate seam, naming `isAvailableToTakeAsWork()`/`onQueueingForExecution()` and the two fields this change collapses. The record stays accurate after the merge: the item is about pointing Lincheck at the claim, which this change does not do - `WorkClaimStateMachineTest` plays the losing interleaving out by hand. `RetryQueue` gained a real fix in the same package (d2e00fa) - `removeAll` read the map off-lock and could return false having removed nothing, leaving a container in the retry queue while in flight. Adjacent, not overlapping: that is a lock discipline defect on the queue, this is a check-then-act defect on the container. Also inherited: the static-analysis lanes and their identity-set ratchets, the `branch_accounting` section of the upstream manifest with its new validator (`bin/check-upstream-map.sh`), the per-PR CI move off the self-hosted box, and the copyright gate now covering every file type rather than only Java. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
@codex review Since the last review pass this branch has only merged Please focus on the state machine under that merged tree: whether every
|
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 628be6c7a6
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
…t to the state's name Review follow-ups on #335. All five threads, and the one that changes behaviour is the first. ABA on the retry delay (Codex, P2). ExecutionState values are enum singletons, so a record that leaves FAILED and comes all the way back to FAILED is `==` to where it started. A selector that read FAILED, found the retry delay passed, and was then descheduled would compare-and-set successfully against a FAILED that is a DIFFERENT failure - one whose deadline the intervening attempt pushed out - and redeliver immediately. The configured backoff is silently skipped. Unreachable under the shipped engine's single selector; routine under the direct-pull engine, where every worker selects. The fix is a per-attempt identity rather than a per-term re-validation: the atomic now holds an immutable Execution (state + attempt sequence), and every transition mints a fresh one, so the compare - which is reference identity - cannot match across a cycle. That closes the whole ABA class in one place instead of re-reading the retry deadline at each claim site and leaving whatever term is added next uncovered. Pinned by aClaimDecidedBeforeARenewedRetryDelayIsRefused(), which plays the interleaving out by hand with no threads, and by everyTransitionMintsAFreshOccupancy(). CONTROL ARM: with transitionTo() returning a canonical instance per state - the ABA restored - both fail; with the fix, both pass. takeClaimOn(Execution) is the seam that makes the act half reachable on its own; splitting it is safe here in a way the defect's own two-step form was not, because the observation is now a value the act compares against. Loops without a clear terminal condition (astubbs). endFlight() and recordVerdict() were hand-rolled `while (true)` CAS loops. Both are now AtomicReference.getAndUpdate/updateAndGet over a total function that is a fixed point on its own output, and endFlightOf()'s javadoc states the termination argument rather than leaving it to the reader: no transition another thread can perform makes a caller go round twice, because only a won claim starts a flight and none can start while the record is in flight. AT_NONATOMIC_OPERATIONS_ON_SHARED_VARIABLE (astubbs). Two of the ten SpotBugs reports were this PR's - deliveryCount and numberOfFailedAttempts - and both were real. deliveryCount's own comment claimed the state CAS ordered the increment; it did not, because the increment happens after the CAS. Both are now atomics with hand-written accessors, so the public int/long API is unchanged. The remaining eight are pre-existing and already recorded in docs/inflight/static-spotbugs-latent-findings.md, in files this PR does not touch. Implicit nulls (astubbs). ProcessingShard#getWorkContainerAt returns Optional. "No container here" is an ordinary answer - the record succeeded and left the shard - so the return type says so rather than the caller having to remember. Is a lost claim a skip bug? (astubbs) No, and the note now answers it instead of leaving it to be re-derived: the container stays in entries, only what was taken is decremented, and the next scan reaches it. What IS skewed is a signal - the refusal path calls addToSlowWorkMaybe, so under direct pull a contested record can be reported as slow, which is the same conflation the note is about. Extraction (astubbs). The latent registration-order bug is now #370, and the note shrinks to a pointer at it plus the one thing a session needs in front of it: it fires PartitionState#onSuccess's assert with no double delivery at all. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
@astubbs — the Ten reports on
Both getters are hand-written rather than Lombok The other eight are pre-existing, in files this PR does not touch, and already recorded in Before / after, No new SpotBugs finding of any rule lands in Local verification on b1a4d4e: |
|
@claude review this New head |
|
Claude finished @astubbs's task in 5m 6s —— View job Review of
|
…reference gate passes `repo: hygiene` was already red on the head this branch inherited (run 32959765640, 628be6c): check-branch-self-reference.sh found three lines in test-lincheck-lane-open-items.md naming #335 with no post-merge attestation. It passed LOCALLY on a bare bin/check-all.sh only because the local branch is merge/335-master, so no PR number resolved and the #NNN arm checked nothing - a skip that reads exactly like a pass. Reproduce it the way CI sees it: PR_NUMBER=335 GITHUB_HEAD_REF=fix/atomic-work-claim bin/check-branch-self-reference.sh The Lincheck target now reads in past tense and describes what actually landed - one atomic holding an ExecutionState and the attempt it belongs to, not a bare AtomicReference<ExecutionState> - and cites the merge base of #335 rather than "pre-astubbs#335 master", which names nothing a reader can check out. The two notes the review pass touched carry post-merge attestations for their own mentions. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Defect-class sweep before mergeThe class is check-then-act on state another thread can move — a decision evaluated over terms that the acting step does not re-validate. Where I looked and what I found, including the negatives.
One thing this sweep turned up that nothing tracks — flagging rather than fixing, because it is a behaviour change nobody asked for at the end of a review pass.
Not reachable through today's callers, which is why it is a note and not a commit: Vert.x's Two ways to settle it, your call: make |
…, not conventional Both non-blocking findings from the automated review of b1a4d4e. The claim has two halves and only one of them is re-validated by the compare-and-set. The state half is, because it is the atomic. The retry-delay half is not, and cannot be, because the deadline is deliberately outside the atomic - so it is checked exactly once, at decision time. takeClaimOn() had no way to check it again, and nothing stopped a future caller in this package reaching the act with a freshly observed occupancy and skipping the decision: fully protected against the ABA, and silently handing out a record whose backoff had not expired. So the ordering is now a type rather than a convention: Optional<Claim> decideClaim() // the WHOLE predicate - state and delay boolean actOnClaim(Claim) // CAS on the exact occupancy that decision was made over Claim's constructor is private to WorkContainer, so decideClaim() is the only thing that can mint one. There is no caller discipline left to forget. onQueueingForExecution() is now decideClaim().map(this::actOnClaim).orElse(false). The delay is still NOT re-read inside the act, and the Claim javadoc says why: re-reading it there refuses the ABA case for the wrong reason and the control arm goes green. That trap is not hypothetical - an earlier draft of the ABA test put the seam one method too early and passed with the defect restored, because the live deadline rather than the stale identity was doing the refusing. Second finding, in the same review: recordVerdict()'s comment claimed endFlightOf()'s termination argument, and the two are not the same. endFlightOf() terminates partly on a literal no-op return; recordVerdict() has none, because transitionTo() always mints. It now states its own bound - at most two threads can write while a record is in flight, the worker recording its one verdict and the controller ending the flight, and neither can repeat because a third write would need a claim that cannot be won from an in-flight state. VERIFICATION. WorkClaimStateMachineTest 10/10, and the control arm still bites: with Execution#transitionTo returning a canonical instance per state, aClaimDecidedBeforeARenewedRetryDelayIsRefused() and everyTransitionMintsAFreshOccupancy() both fail. Full-reactor bin/ci-unit-test.sh BUILD SUCCESS, bin/check-all.sh 15/15, SpotBugs AT_NONATOMIC_OPERATIONS_ON_SHARED_VARIABLE still 8 - the pre-existing set, none in the changed classes. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
Both non-blocking points taken, and both are now closed in code rather than in prose — thank you, the first one was a real seam. 1. Optional<Claim> decideClaim() // evaluates the WHOLE predicate - state and delay
boolean actOnClaim(Claim) // CAS on the exact occupancy the decision was made over
This is also why the delay is not re-read inside the act, and the 2.
Verification on the new head: |
…y it was confusing The note read as a report of finished work. Most of it was archaeology - the 2026-08-22 sightings and why the inference from them was correct - which is settled history and belongs to #370, not to a directory an agent scans for what is open. What was actually open was one sentence buried in the middle: the ordering in `maybeRegisterNewPollBatchAsWork` is safe only because registration and completion both run on the control thread, and nothing states or tests that. It now leads, as a "what is left to do", because an invariant that holds by accident of today's callers is the thing a future session has to close. The debugging pointer stays - it is forward-looking, not retrospective: it tells someone meeting `assert (removedFromIncompletes)` where to look once the claim is ruled out. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Two merges in one: #371's PIT fix, which is why this happened now, and #335's `Execution` state machine, which conflicted. THE PIT LANE. `Mutation Tests (PIT, PR-scoped)` was cancelled on this PR after half an hour. It was not PR-scoped: a comment inside the invocation's line continuation truncated it, so `-pl ... -am` never reached Maven and the lane walked the whole reactor. Master's fix builds the arguments as an array, and `bin/test-ci-mutation-test.sh` passes here, including the five argv arms that assert the scope flags survive. THE CONFLICT, and it is semantic rather than textual. #335 replaced the `maybeUserFunctionSucceeded` field with `recordVerdict(false)`, a transition on the `Execution` held in an `AtomicReference`, and calls it straight after `updateFailureHistory` - deliberately in that order, because the state write is what publishes the retry deadline. This branch had wrapped the same call site in `try`/`finally`, so a throw from `updateFailureHistory` could not leave the container half transitioned. Resolved by keeping master's structure and this branch's `finally`, because the two do not actually conflict: `finally` runs after the history write on every path that completes, so the ordering contract is untouched. It differs only on the throwing path, and there the choice is between publishing FAILED beside the previous attempt's deadline - one skipped backoff, record still makes progress - and not publishing at all, which strands the record and blocks its shard for the life of the process. The stall is strictly worse, and closing it is what this PR is for. The reasoning is written at the call site so the next reader does not take the `finally` for an oversight and remove it. Worth recording that both hazards the `finally`'s original comment named are now closed at source by this branch's own work: `computeRetryDueAt` catches the overflow and `getRetryDelayConfig` never throws. The `finally` is kept for the argument it always rested on - not today's known lines, but every line added above it later. `WorkClaimStateMachineTest` passes, including the ABA arms #335 added, and the vert.x comment that cited `maybeUserFunctionSucceeded` by name now cites the verdict instead. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
…claiming the invariant Review pass over the ownership change. Three reviewers plus the automated review on the PR found no correctness defect in the CAS protocol itself; what follows is what they did find. THE ONE MISSING TEST, AND IT WAS THE LOAD-BEARING ONE. countAsSelectable's residency recheck - `if (entries.get(wc.offset()) != wc) uncount(wc)` - is the single line that makes claim-then-confirm safe when the two threads actually interleave rather than merely run in the damaging order, and nothing exercised it. It is reachable without racing, through the same kind of seam the other four tests use: a revocation landing between handleFutureResult's staleness check and sm.onFailure has the controller count a container selectable again after the poller has taken it out of the shard. Deleting the recheck makes the new test red (expected 1, got 2, with the ground-truth scan disagreeing); the other four stay green, which is what makes this test the one that covers it. THE INVARIANT WAS STATED MORE STRONGLY THAN IT HOLDS. "Equals the units held, and is never negative ... by construction rather than by clamping" is true at quiescence and false instant-by-instant: countAsSelectable claims and then increments, which is two atomics, so a reader in that window sees the scan one ahead - and if a concurrent uncount wins the release there, the decrement lands first and the counter momentarily reads -1. Every such interleaving still settles correct and every consumer reads the aggregate, which floors at zero, so this is a documentation defect rather than a behaviour one. But "no clamp, by construction" is the headline claim of the change, so it is qualified where it is made rather than defended in a PR comment. Collapsing the two atomics into one is the follow-up that arrives with #335's Execution transition. A COMMENT THIS CHANGE FALSIFIED. getNumberOfWorkQueuedInShardsAwaitingSelection reads the retry-queue size before the shard sum, and the only record of why was "normally retry queue is updated before shard counters are on work taken". That stopped being true here: the decrement moved out of a single batch at the end of getWorkIfAvailable and into the selection loop, so the shard counter now drops first and the window spans the rest of the scan. The read order is still right; the reason is now re-derived in place, with the arithmetic, and lands on the same conclusion - the skew is nil or high, never low, and low is the direction that closes drain() early. A REAL DEFECT FOUND, NOT FIXED HERE, AND RECORDED. The poller's stale sweep calls iterator.remove() on a ConcurrentSkipListMap entry set, which removes by KEY - so a fresh container the controller put at that offset since next() returned is what actually leaves, while the uncount that follows releases the stale object. The primary harm is the lost record, and it predates this change; the secondary harm is the counter, and its DIRECTION is this change's doing - formerly one low and resynced by the clamp, now one high and permanent, which is the direction that stops drain() reaching CLOSING. Not fixed here because there is no identity-keyed removal to reach for: entries.remove(key, value) and entrySet().remove(entry) both compare with equals, and WorkContainer.equals is topic/partition/offset only, so the replacement compares EQUAL to what it replaced - the same reason countAsSelectable compares with !=. Written up with the decision it needs in docs/inflight/bug-stale-sweep-iterator-evicts-fresh-replacement.md, and the call site now points at it. Also: the @GuardedBy carve-out claim gave the tree's AGENTS.md more than it actually says - it names volatile, not atomics - so the javadoc now gives the reason directly and quotes what the rule does ask for. A dead local in getWorkIfAvailable is gone. The 2026-08-26 Chaos Pain Suite red is logged in the 857 ledger as a NEGATIVE result: this is the branch on which the fleet NO_PROGRESS arm should have stopped firing if that defect were what it caught, and it fired anyway. One backlog line for the state package's copied test harness. Verified: bin/ci-unit-test.sh BUILD SUCCESS across all 11 modules; bin/check-all.sh 15 ran, 15 passed. Nothing loosened, weakened or quarantined.
#335 - the per-attempt `Execution` identity behind an `AtomicReference` - merged while this was in review, and it rewrote the exact method this change touches. Both conflicts are in that seam and both resolve to a union rather than a choice: - `ProcessingShard#getWorkIfAvailable` - master turned the claim into `if (workContainer.onQueueingForExecution())`, one call that evaluates the whole decision and claims from the state it evaluated. The `uncount()` this branch added now sits INSIDE that branch, so the unit is spent only by the caller that WON the claim. That is strictly better than the pre-merge form: a refused claim leaves the container resident and still holding its unit, which is what `docs/inflight/core-a-lost-claim-means-two-different-things.md` says the refusal must do. - `WorkContainer` - `holdsShardAvailableUnit` and its three accessors sit alongside master's `Execution`/`deliveryCount`. `maybeUserFunctionSucceeded` was NOT carried over: master folded it into `Execution` and deleted it, and git offered it back only because it was adjacent to the new field. FOLDING THE UNIT INTO `Execution` IS NOW UNBLOCKED, AND STILL NOT DONE HERE. This PR's description says the flag should become part of #335's transition once that lands, so that the claim and the unit move as one atomic step instead of two. That is now possible rather than hypothetical - and it is the same follow-up as the transient this branch documented on `availableWorkContainerCnt` (a reader can see the counter at -1 between the claim and its increment), because one transition is exactly what removes it. Left for its own change: it rewrites `Execution`'s contract, and this branch's value does not depend on it. Verified after the merge: `ShardAvailableCountOwnershipTest` 5, `WorkClaimStateMachineTest` 10, `ShardManagerTest` 4, `ProcessingShardStaleReplacement909Test` 4, `WorkManagerTest` 26 - 49 tests, 0 failures.
Review finding on #373, and it is right: the field javadoc claimed "every consumer reads the aggregate ... which floors at zero", and one does not. The under-served-retrieval diagnostic in ShardManager#getWorkIfAvailable sums getCountOfWorkAwaitingSelection() straight across the shards behind log.isDebugEnabled(), with no max(0, ...) - so the transient this branch documented can surface as awaitingSelection=-1 in that one debug line. It drives no control flow, and it is left unfloored deliberately rather than overlooked: that line exists so a human can tell a genuine stall from normal back-pressure by reading what the accounting actually says, and clamping it would hide exactly the drift it was added to expose - the same argument this change makes for deleting the clamp on the counter itself. So the javadoc now names the exception instead of the diagnostic acquiring a floor. Also states, where the transient is described, that folding the flag into #335's Execution removes it rather than masking it - the distinction that makes it a real follow-up rather than a note. Verified: ShardAvailableCountOwnershipTest 5/5, bin/check-all.sh 15 passed.
…its green is #347 committed this harness asserting the checkpoint-3 tear EXISTS, naming this fix as the one that inverts it. The previous head could not: #345's NullPointerException out of ShardManager.removeWorkFromShardFor was reached through the same revokeAndReassign operation, and `assertThat(report).contains("completeWork")` cannot tell one from the other. #345 has landed, so that blocker is gone and the lane was re-run rather than reasoned about, as docs/inflight/test-lincheck-lane-open-items.md requires. MEASURED, on this tree, at the committed iterations(1_000): - Seven valid runs, every one exhausting the whole bound with no violation, 120.3s-143.1s. A hit used to land in 9-13s, so none of these is a near miss. The designed-red arm therefore fails every run and cannot stay as it is. - CONTROL, the same sitting: a compiling mutant that re-resolves the partition state at handleFutureResult's two action sites - the checkpoint-3 tear put back, nothing else changed - hit ONCE IN SIX runs. The hit is exactly the calibrated signature: AssertionError at PartitionState.onSuccess, through PartitionStateManager.onSuccess and WorkManager.onSuccessResult out of handleFutureResult, completeWork racing revokeAndReassign, and no NullPointerException anywhere. So the harness still reaches its seam - it has not gone quiet. The honest reading of those two together, which the javadoc now carries: seven clean runs is WEAK evidence of absence when the positive control itself misses five runs in six. What pins this fix is WorkManagerStaleCheckDoubleLookupTest, which forces the interleaving deterministically; the Lincheck arm is a search running beside it, and a pass here must never be read as a proof. So: renamed to stressMustNotRediscoverTheCheckpointThreeTear and flipped to Lincheck's own linearizability check, the shape RetryQueueLincheckTest already uses. The bound is left exactly where it was measured - re-pricing needs the starve-and-count procedure the note prescribes, and neither raising nor lowering it on this handful of runs would be better founded than the estimate it replaces. THE COST, because it is now the lane's largest number. An inverted arm cannot stop at a first violation, so it always pays its whole bound: this arm is ~2 minutes every run and the whole lane went from 26-29s to 2m32s-2m37s (measured twice, green both times). bin/lincheck-test.sh's header and docs/testing.md said "well under a minute" and now say what it costs. THE OTHER HARNESSES, re-checked in the same sitting as that rule requires: - ShardManagerLincheckTest is unchanged by #335 landing - same counterexample the note records, revokeSweep(0) then addWork(0) against addWork(0). Nothing in the note's item 1 moves. - PartitionStateLincheckTest has now been run against a tree carrying #337's fix for the first time, and does NOT invert. It still reports a violation naming commit(), but the report is not stable between runs - mostly a value divergence between commit() and succeed(), once an ArrayIndexOutOfBoundsException out of two concurrent commit() calls printed with no frames. #344 remains its trigger; recorded as a sighting, not a finding. - RetryQueueLincheckTest and both toolchain controls green and unchanged - the model checker is still blocked on the Lombok callSuper defect, so the note's item 8 arms are still not free. docs/testing.md's "every harness asserts a bug EXISTS" bullet was already wrong before this change (RetryQueueLincheckTest never did, ShardManagerLincheckTest stopped when #345 landed); it now names the three shapes in the lane and why the third exists.
…e drain arm The Chaos Pain Suite red on run 33014328251 is the confluentinc#857 family's open deadlock again, not a new defect and not Class 2. Failing test: ChaosRevokeUnderWorkDrainIT.revokeUnderDrainingStopsStaysProtocolHonest, seed 3198328355855848347. The first scenario in the log is a different class at a different seed and it passed, so the replay line nearest the top of a chaos log picks the wrong one - the failsafe artifact names the failing class in one attribute, which is the route docs/ci.md already prescribes. Signature: instance 44's poll thread BLOCKED acquiring the commitCommand AtomicBoolean held by pc-control-PC-44, reached through onPartitionsRevoked -> commitOffsetsThatAreReady. Frame for frame the four captures above it, once the line numbers are re-resolved: this one reads :1589/:552 where they read :1585/:548, and the shift is master's, not the branch's. Class 2 did not fail the build - 26 LAG_STAGNATION observations in the usual 154s band above a "violations (0)" line, so the 637ea31 demotion is doing what it was meant to. Rules out #335 for this signature, which was the live suspicion: the first three captures ran hours before #335 merged, on the identical monitor, holder and method pair. Not replayed. Five seeds now exist for the verification the first 2026-08-26 section asks for and none has been run, so docs/solutions/runtime-errors/revoke-path-commit-deadlock-between-poll-and-control-threads.md keeps its Unproven status.
…, closing a stall and a wrong-commit route (#346) WorkManager.handleFutureResult's staleness checkpoint 3 validated one partitionStates.get(tp) lookup and acted on another. Nothing serialises the two against the broker-poll thread, so a revoke(+reassign) completing in the gap meant the check passed against the OLD state object while the actions ran against its replacement. Both harms were reproduced deterministically before any fix: - FAILURE PATH, a confluentinc#857-family permanent stall. The stale container was re-queued into the retry queue AFTER the revoke sweep had emptied its shard. removeStaleContainers reaches retry-queue entries only through shard contents, so nothing could remove the orphan again; once its retry delay elapsed, workIsWaitingToBeProcessed() read true forever with nothing assigned. - SUCCESS PATH, the only production route into a wrong commit. The completion landed on the freshly assigned state, tripping PartitionState.onSuccess's assert under -ea; without -ea it silently dirtied a bootstrap-phase state with a dead epoch's completion, which was the gate-opener making the bootstrap-reset commit tear reachable. THE FIX ACCEPTS THAT NO EPOCH CHECK HERE CAN BE ATOMIC WITH THE ACTIONS, and makes the actions safe by construction instead. The state is resolved ONCE and the staleness answer shares that reference with the success action, so a completion can only ever mutate the state object whose epoch matched the container; a concurrently revoked object is already unlinked from the live map, so the write is invisible to the commit path and the record is redelivered to the partition's next owner. At-least-once holds in every interleaving. The failure path additionally re-validates against the LIVE map immediately before the retry re-queue, because that action targets ShardManager structures the revoke sweep cleans and no captured PartitionState reference can protect; skipping the re-queue is the safe direction. PartitionStateManager.onSuccess/onFailure now take the caller-resolved state - WorkManager was their only production caller, and the single-arg overloads that did their own second lookup are reached from tests alone. HONEST RESIDUAL, FAILURE PATH ONLY. The re-queue re-validation is still check-then-act, shrunk from a whole-checkpoint gap a deterministic revoke reproducibly fit inside to a one-context-switch sliver. The class-level closure is the retryQueue sweep tracked in docs/inflight/bug-retry-queue-orphaned-by-inline-stale-removal.md. The success path has no residual. CONTROLS, run rather than assumed. WorkManagerStaleCheckDoubleLookupTest, five tests: both race arms red on the unfixed tree and green fixed, both serialized control arms green on both, so the failures belong to the interleaving and not the mutation. A mutant re-resolving the state at the action sites flips exactly the success arm and the bootstrap-gate probe red; deleting the re-queue re-validation flips exactly the failure arm red; deleting a test's arm call fails its raceHasFired guard, so the seam cannot go silently dead. THE LINCHECK INVERSION IS DISCHARGED, AND THE MEASUREMENT SAYS DO NOT TRUST IT ALONE. #347 committed WorkManagerLincheckTest asserting the tear EXISTS, to be inverted when this landed. On the fixed tree it missed all seven runs, each exhausting 1000 iterations - but a control carrying the defect re-introduced hit only once in six, producing exactly the checkpoint-3 signature and no NullPointerException now that #345 has landed. So the harness still reaches its seam and is simply a poor detector at this bound: seven clean runs is weak evidence of absence. The arm is inverted, the bound left where it was measured, and the javadoc and note both say the deterministic test is what pins this fix. The inverted arm cannot exit early, so the lane went from about 27s to about 2m35s - corrected in bin/lincheck-test.sh's header and docs/testing.md, which both claimed well under a minute. THE CONTRACT'S ONE-BUG-PER-HARNESS ASSUMPTION IS DEAD, and PartitionStateLincheckTest is now a VACUOUS GREEN. With #337 and #344 both landed it still passes contains("commit()") - on reports about neither tear it names, because the interleaving table names commit() whatever threw. Two shapes recur: an ArrayIndexOutOfBoundsException out of ArrayList.add, the plain-ArrayList defect #57 fixes, which the lane has now re-found from a second harness pointed elsewhere and with a frame this time; and PartitionState.onSuccess's assert reached from two parallel succeed operations, which production cannot perform because only the control thread completes work. Whether that second shape is artefact or defect is the same question already parked over ShardManagerLincheckTest's addWork, and it should be settled for both harnesses at once and AFTER #57 lands, since that PR removes one of the two shapes. Recorded rather than guessed at; the assertion is deliberately not re-pointed here. ALSO IN THIS CHANGE. The infer ratchet's gone-arm fired on WorkManager.onFailureResult; a three-tree control with pinned Infer v1.3.0 showed Pulse re-attributing an unchanged deref, with the identical shape in the neighbouring overload reported at neither head, so the line is deleted rather than a guard added. The same commit pins -encoding UTF-8 on the lane's hand-rolled javac, which had been inheriting the developer locale and made bin/infer-test.sh exit 2 under LANG=C. The duplicate-code thread is closed by extracting the shared fixture to ModelUtils.registerOneRecordAndTakeIt. WorkManager's pm and sm fields carry TODO(refactor) markers plus a docs/refactoring.md entry for the rename raised in review, not folded in here because both are @Getter(PUBLIC) and it would reach main, unit and integration sources far outside this seam. A CHAOS RED ON THE WAY THROUGH, DIAGNOSED AND RECORDED RATHER THAN RE-RUN. ChaosRevokeUnderWorkDrainIT failed once with the AB-BA revoke-path monitor block - the poll thread BLOCKED on the commitCommand AtomicBoolean held by its own pc-control thread, through onPartitionsRevoked - which is master-state and not this change's: the diff touches no locking and no AbstractParallelEoSStreamProcessor line, and the three earliest captures predate #335's merge with the identical monitor, holder and method pair. Filed as the fifth capture in docs/inflight/bug-857-family.md. The lag-stagnation observations in the same run were non-gating, so the CLASS2 demotion held. Nothing was quarantined, loosened, or re-run to hide it. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
…et, instead of inferring it (#373) ## Description `ProcessingShard` decided whether it had already spent a container's unit of the shard count by reading that container's observable state after the fact - `isNotInFlight()` on an entry already removed from the map. No state at time T records *who* spent the unit, so the read cannot answer the question however well synchronised it is. A container now carries an `AtomicBoolean` claim, the shard claims and releases it by compare-and-set, and **only the CAS winner moves the counter**. The floor-at-zero clamp is deleted: the count cannot go negative by construction, so a clamp would only hide drift. **Four instances of the class, not the one that was reported.** The class is *a counter whose spend is inferred from observable state rather than owned by an atomic outcome*, and all four were in this one file: the reported double-deduction on revocation; `addWorkContainer`'s stale-replacement branch; `onFailure` not idempotent despite its contract; and a **sign-reversed** case where revoking a failed record inside its retry delay deducted nothing, leaving a permanent OVERcount - which is the confluentinc#857 stall signature. **Keeping the counter was a decision, not an assumption.** Of its four readers, three could derive their answer by scanning: `drain()` needs only a boolean, the `WAITING_RECORDS` gauge is per-scrape, the under-served diagnostic is debug-gated. The load gate cannot - measured at ~4.5 ns/entry, a scan costs 45 µs at 10k buffered records and 236 µs at 50k, which at control-loop cadence is 4-24% of a core. **Deleting the counter becomes viable once #336 lands**, since that removes the only reader needing O(1); nothing here blocks it. Note the counter is not merely a gauge: `drain()` reads it, so reading low closes the consumer early. Every guard is proven red-then-green by reverting the main code, and mutation-checked three ways. ## Known cost, accepted deliberately The claim and the count are two atomics, so a reader can transiently observe the count at **-1**. It settles correct and every consumer floors, except the debug stall diagnostic - clamping the one line that exists to expose drift would be incoherent, so the javadoc names it as the exception rather than the code hiding it. Collapsing both into #335's `Execution` transition removes the transient rather than documenting it, and is unblocked now #335 has landed. **This PR makes one thing worse, and it is recorded.** `ProcessingShard`'s stale sweep calls `iterator.remove()`, which on a `ConcurrentSkipListMap` removes by *key*, so it can evict a fresh replacement written by the controller between `next()` and `remove()`. The lost record predates this change; the counter drift's direction does not. It was one-low-and-clamped, and is now one-high-and-permanent - the direction that keeps `workIsWaitingToBeProcessed()` true and stops `drain()` reaching `transitionToClosing()`. There is no identity-keyed removal to reach for, since `WorkContainer.equals` is offset-only, so closing it means deciding how the sweep and the replacement branch coordinate at all. Open, with the deterministic test it needs, in `docs/inflight/bug-stale-sweep-iterator-evicts-fresh-replacement.md`. ## Naming The operations are spelled from the vocabulary already in the file: `includeInSelection` / `excludeFromSelection` on the shard, `claimSelection` / `releaseSelection` on the container, and the field is `workAwaitingSelectionCount`. "count"/"uncount" used *count* in the regard-as sense beside a field using it in the enumerate sense, and the field's old name borrowed "available" from `isAvailableToTakeAsWork()`, which means something different - a caveat paragraph existed purely to explain that collision, and has been cut down to the fact it was hiding. ## Also here A capture of the revoke-path commit deadlock in `docs/inflight/bug-857-family.md`: this branch went green on the arm it had been red on, three heads later. Recorded because the sighting is otherwise lost when the branch goes green.
… population, the claim CAS drives the selection count #373 rewrote ProcessingShard under this branch. It replaces an INFERRED spend of the shard's selectable-work counter with an OWNED one: each container carries an AtomicBoolean claim and only the compare-and-set winner moves the counter. This branch pairs every change of the shard's POPULATION with a real map mutation. The two meet in addWorkContainer, both sweeps and onFailure, and the resolution is not textual. They compose, because they count different things and neither is derivable from the other. A record can leave the selection population without leaving the shard (taken as work) and can leave the shard holding no claim (revoked at a worker). So the rule the merged code follows is: the MAP's own return value drives the population, the CAS outcome drives the selection count, and neither is ever conditioned on the other. That is now stated on the population field. Per site: - addWorkContainer keeps this branch's admit-first shape - population.onAdmitted() then put(), and the put's return value, never the read above it, decides whether anything was displaced. master's version had gone back to deciding from the pre-read. On top of it, the arrival is offered a claim and the displaced container gives one back, each settled by its own CAS. Include runs before exclude so the transient is +1 rather than -1; -1 is the direction #373 documents as the one a reader can see, and reading low is what closes drain() early. Verified against the interleavings both designs worry about: if a sweep removes the arrival between the put and includeInSelection, the residency recheck hands the claim straight back and the sweep's retirement pairs with the admission - nothing double-counts either way. - retireAlreadyDeducted and retireAndDeductIfStillCounted COLLAPSE INTO ONE. That was this branch's own prescription for the counter - "an ownership flag the shard claims and releases with a CAS ... collapses the two into one method" - and #373 is that fix. There is now a single exit path, retire(), driven entirely by what the map gave up. - the isNotInFlight() window is closed, not deferred. It was the open review finding on this PR, left unfixed because the load gate had stopped reading the counter. The predicate is gone. - both stale sweeps keep this branch's removal-by-return-value, so the population retirement, the claim release and the container handed to the caller all follow the object that actually LEFT rather than the one that was inspected. That closes half of bug-stale-sweep-iterator-evicts-fresh-replacement.md - the accounting half. The by-key eviction of a fresh replacement is untouched and still open; its note now says which half is which. - the last-resort eviction in getWorkIfAvailable keeps this branch's retry-queue cleanup, so it cannot orphan a queue entry, on top of master's claim release. The map stays PRIVATE. master added @Getter for three white-box test plants; a public handle on the map falsifies the invariant this branch's javadoc claims and lets a test insert without admitting, which drifts the population silently and fails nothing. plantResident() replaces it - same white-box reach, pairing kept, no claim offered because those tests are planting containers that have already spent one. The two accessors also collapse: one getWorkContainerAtOffset(), returning #335's Optional under the name this PR's review asked for. getWorkableRecords() is unaffected - it reads the population and the retry queue, and never read the selection counter. #373's claim that deleting the shard counter becomes viable once this lands HOLDS, and is now recorded rather than left in a merged commit body (docs/inflight/core-shard-selection-counter-can-now-be-derived-by-scan.md). Grepping every reader in core main: the load gate is gone, and what is left is drain() at shutdown cadence, the WAITING_RECORDS gauge per scrape, and three lines behind isDebugEnabled/isTraceEnabled. No steady-state reader remains, so the scan cost #373 measured is no longer paid per control loop - though the note asks for that number to be re-taken at the new cadence rather than reusing a verdict that answered a different question. Notes brought back in touch with reality, per docs/inflight/AGENTS.md's four outcomes: bug-available-work-counter-is-still-an-approximation.md is retired - #373 landed the derive-it fix it prescribed and shipped the invariant check it asked for. What outlived it is a different item and gets its own note, the second counter in the same calculation (bug-number-records-out-for-processing-is-a-plain-int.md). The lincheck lane's entry for this counter is rewritten: the drift it described is gone, and what is interleaving-shaped now is the two-atomics transient the claim protocol accepts. Also repaired: bug-stale-sweep-iterator-evicts-fresh-replacement.md arrived citing the note this branch deletes, the same way two java files did last merge; same fix, the history pointer docs/citations.md prescribes at a sha on master. 477 core unit tests green, including #373's ShardAvailableCountOwnershipTest through the new seam and this branch's ShardPopulationRaceTest unchanged in substance. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01XC9N7S8gSd6fCYNNe64D6j
…ut consuming a retry Core has two ways for work to come back and neither fits a worker that vanished. onSuccessResult commits the offset. onFailureResult records a failure, queues a retry and bills the record an attempt. But an external engine whose worker process dies mid-record has no verdict to report: nothing failed, so charging a retry attempt is wrong, and a record that reliably outlives its worker would burn through its budget without ever being processed. Today handleFutureResult throws IllegalStateException on a verdict-free return, so the engine cannot hand the record back at all. This adds the third path. An abandoned marker on WorkContainer, distinct from the user function verdict, selects a branch that does what the stale-work branch does - endFlight and exactly one decrement - and then restores shard availability without a retry-queue insert, so the record is immediately selectable rather than waiting out a delay it never earned. An empty verdict with the marker unset still throws; that is the bug it always was. Two things worth pointing at. The in-flight counter is the danger. It gates the broker poller through isSufficientlyLoaded, and drift in it stalls the consumer silently while it still looks alive - this repository has that failure documented already. So onAbandonedResult ignores a repeat return rather than decrementing twice, which is a real case: a disconnect detected while a report is already in flight returns the same record twice. And the tests assert the exact counter value, not merely that work keeps flowing. The subtler one is that redelivery has to clear the previous delivery's state, because a redelivery is a fresh attempt. That was originally two clears; core has since made one of them structural, which is the first of the changes below. WHAT THE REBASE ONTO CURRENT MASTER CHANGED This is #295 resurrected. It was closed as folded into the language-proxy PR #293 while that PR was the only home for it, and is being landed on its own again because a core scheduling primitive should not first arrive inside a remote-runtime PR. The commit is the same change; three things had to move to fit master, and all three are master's design winning rather than this one being watered down. The verdict no longer needs clearing on redelivery. #335 collapsed `inFlight` plus `maybeUserFunctionSucceeded` into one atomic ExecutionState, so the verdict is now derived from the state and a won claim lands on IN_FLIGHT, which has none. The original explicitly reset `maybeUserFunctionSucceeded` in onQueueingForExecution; that field is gone and the reset with it. Only the abandon marker is cleared, in actOnClaim - after the compare-and-set, so it is the claim WINNER that clears it. A new test, aFreshClaimClearsThePreviousDeliverysAbandonMarker, covers what that clear now protects on its own: without it, a record that had ever been abandoned would have every later verdict-free return silently forgiven instead of reported. The marker stays its own field rather than becoming a seventh ExecutionState, and it is an AtomicBoolean rather than the plain boolean it was. The argument for putting it in the state machine is real - "one field, because two were the bug" - but the hazard that collapsing fixed was a claim DECISION reading one field and being contradicted by the other, and nothing reads this one to decide a claim. It is a note left by the returner for handleFutureResult, which is the same reason selectionClaimed is its own field. The field javadoc carries that argument, so the next reader finds it rather than taking the field for an oversight. Atomic because the marker is written by whichever thread noticed the worker was gone and read by the controller. The shard-level restore goes through the selection claim rather than a raw counter increment. #336 replaced availableWorkContainerCnt with a compare-and-set-owned claim, so ProcessingShard.onAbandoned calls includeInSelection, and idempotence now comes from the compare-and-set instead of being argued. Master's conservation walk reserved this in its own javadoc - "the branch developing external engines adds a further departure". On contact it turns out not to be a departure at all: abandonment retires nothing, because the record stays resident in its shard. So the reservation is answered by a test that asserts the conservation figure is unmoved in both directions, and the javadoc corrected to say so. Teeth checked twice, each with a control arm. Removing the handleFutureResult branch fails all seven behavioural tests and leaves workWithNeitherVerdictNorAbandonMarkerStillThrows correctly green, since that one guards the pre-existing behaviour. Removing the marker clear in actOnClaim fails exactly aFreshClaimClearsThePreviousDeliverysAbandonMarker and nothing else. Full reactor unit build green. (cherry picked from commit 4b4ff19) Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…an all be true at once Every probe in this repo watches liveness: is the system still progressing, sampled as one number over time. A system can be perfectly live while lying about what it holds, and it will report healthy right up until the lie matters. This watches something else - whether PC's separate accounts of its own state are mutually consistent. The state it flags: work queued for execution while PC believes no offsets are incomplete and nothing is out for processing. That is reachable without anything throwing, because getNumberOfIncompleteOffsets() sums over ASSIGNED partitions and returns zero whenever that map is empty, however much work is queued. It was observed on three consecutive log lines during the transactional throughput investigation - an executor queue of 319 against a target of 320, with both counters reading zero. describeProgress() also now reports the executor queue depth, which is the only one of these numbers sourced from OUTSIDE PC's own bookkeeping. The others are things PC believes; that one is a thing that is true, and the asymmetry is what tells a reader which number to doubt. IT REPORTS, IT DOES NOT ASSERT - and it is UNARMED The two reads are not atomic, so a queue draining between them is a false positive, and one sample is not evidence. Turning it into a gate needs the contradiction to PERSIST across samples and needs it shown firing on a tree that should fail. Neither has happened. The run that would have exercised it PASSED - this branch's failure is load-dependent - so the check has never fired. Its silence currently means nothing, which is the same standard applied to the ambient probe that this whole thread came from, and it is recorded that way rather than presented as working. WHAT THE DIRECT-PULL STACK DOES TO IT #361 hands work SELECTION to the workers - they block on the shards and take their own records instead of the control loop pushing into an executor queue - so executorQueueDepth() would read zero permanently and this condition would become unreachable. The check would not fail; it would go QUIET, which is worse, and is the exact failure class this session has been documenting. A version meant to survive that stack must source "work PC is holding" from something direct pull still has. The finding may also not be novel. #336 derives the poller load gate by conservation rather than from maintained totals, and numberRecordsOutForProcessing is described in its own note as the last counter of its family after the others were dissolved into derived views; #335 adds ExecutionState for concurrent claims. The disagreement sits inside the area that stack is rebuilding and must be checked against it before being called a new defect. The throughput regression is measured separately and stands either way.
…s engine back Merges `feats/proxy-dispatch-clients` (#390, the top of the Wagon A stack) into this branch, so that #293 can be retargeted onto that rung and its displayed diff collapse to the engine 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-293`. 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 most of the conflicts. ## Which side won, and why The rule was: the god branch's SEMANTIC content, expressed on the stack's refined STRUCTURE. - **Engine-only files stayed** - `ProxyProcessor`, `ConfigureHandler`, the dispatch waves, epochs, leases, reconnect and the produce path. They are the residue this campaign exists to expose. - **Stack-only files were taken** - the ten foreign clients, the runner registry, the sidecar shim, the self-retiring guards, `Main` and `ParentDeathWatchdog`. - **Records lost, everywhere.** The stack converted five value types from Java records to plain final classes because Jabel rewrites a record into a class with no source positions and Error Prone 2.42.0 crashes rather than reporting on it. This branch never met that only because its root pom had no Error Prone; master's does now, so the conversion is load-bearing here. - **Dependency pins are the stack's**: protobuf 3.25.8 (`requireUpperBoundDeps` refuses 3.25.5 against grpc-protobuf's transitive ask), gson 2.14.0, `error_prone_annotations` 2.48.0 in the proxy and the Java client aggregator and 2.47.0 in the protocol module - deliberately unequal, because their highest requesters differ. - **`clients.yml` is this branch's**, not the stack's. The stack's copy is a smaller earlier file whose own header says it has no conformance lane, no per-language static analysis and no dependency audit; only its `--component clippy` fix and its workflow-level `PC_FOREIGN_CLIENTS_STRICT` were folded in. - **Master's copy won for everything master owns** - the gate scripts, the copyright checker and its self-test, `repo-hygiene.yml`, `docs/ci.md`, `docs/copyright.md`. This branch's only delta in most of those was a copyright header in a pre-canonical form, written before #338 landed the canonical one on master. - **The two ledgers were merged as unions**, not resolved to a side. `docs/inflight/bug-857-family.md` keeps both sides' sightings, this branch's renumbered per the file's own stated convention; `docs/quarantined-tests.md` drops three entries whose owning PRs retired the annotation on master, because a registry entry without a `@Quarantined` test is a hard gate failure. ## The verdict-free work return, re-expressed on master's shape Unit U4 is still on this branch as `WorkManager.onAbandonedResult` and its call sites, and `WorkManager.java` auto-merged keeping them - so resolving `WorkContainer`, `ProcessingShard` and `ShardManager` to master alone would have compiled to nothing. Master moved underneath it: #335 collapsed `inFlight` and `maybeUserFunctionSucceeded` into one atomic `ExecutionState`, and #336 and #373 replaced the shard's raw `availableWorkContainerCnt` with a compare-and-set claim. `ProcessingShard.onAbandoned` and `ShardManager.onAbandoned` are therefore taken verbatim from `feats/proxy-verdict-free-return` (#295), which is cut on current master and is the authoritative re-expression - it routes the restore through `includeInSelection` rather than incrementing a counter that no longer exists. `WorkContainer` is NOT #295's, and the difference is deliberate: that branch marks abandonment with an `AtomicBoolean` cleared by the claim winner, while this branch keys it by delivery (`markAbandoned(long)`, `isAbandonedForCurrentDelivery`, `isReturnForSupersededDelivery`). The delivery-keyed form is strictly more: it tells a late return for delivery n from a live return for delivery n+1, which the boolean cannot, and acting on that confusion ends a running flight and decrements `numberRecordsOutForProcessing` twice. It also needs no clearing on redelivery, so it does not race the claim. Master's own `deliveryCount` - already an `AtomicLong` incremented only by a won claim - supplies the identity, so the field is the one thing added. ## The reconciliation duties #387 and #390 recorded Both are discharged, and the guards they left behind now pass because the gap closed rather than because anything was edited. **`ProxyHarness` and `ConformanceHarness` are one class again.** #387 cut the latter out of the former with the engine lane removed, and wrote into the class header, the module pom and a deferred inflight note that whoever landed the engine must merge them - leaving open which module it ends up in, on the grounds that "the engine lane decides". It decides the sidecar module: the lane needs `ConfigureHandler` and `ProxyProcessor`, and `TestModeMain` - the process a spawned foreign runner actually starts - is in that same test tree and would otherwise have to import the conformance module and close a cycle. So the NAME is the extracted rung's and the MODULE is this branch's, and `parallel-consumer-proxy-conformance` reaches it through the proxy test-jar its pom already declared. The stack's refinements to the class survive: `Delivery` stays a plain final class, and the `startEmbeddedClient` lane stays beside the restored `startEngine`. **`ConformanceDriver.spawnAgainst` is back**, with `LanguageRunner` a `ConformanceBinding` again - it stopped being one on #390 for exactly as long as there was no engine to spawn against - and the ten foreign bindings are registered in `ConformanceBindings.selectable()`. `ConformanceBindings.deferredUntilTheEngineArrives()` is kept and now returns an empty list, which is the honest reading: nothing is deferred. It is not deleted because `TheEngineArrivingMustBringTheForeignCellsTest` reads it, and because the next capability to arrive without a cell needs that list rather than a new one. ## Two duplications the merge surfaced, resolved rather than left - **The protocol documents existed twice.** #383 moved `protocol-specification.md` and `client-authoring-guide.md` into `parallel-consumer-proxy-protocol/docs/` and predicted the rename conflict here. The protocol module wins on the merits - all three artifacts are the contract rather than the server - but the CONTENT is this branch's, because the stack's copies carry "no client generates from this schema yet" and a `git show` fallback for a Go generator script that exists here. `SpecificationCoverageTest` moved with them; every citation was re-pointed. - **`HarnessScenario` existed twice**, once per harness. The stack's plain-class version won, in the sidecar module beside the class it serves. ## Five Java records had to become plain final classes, or nothing compiled `master` gained Error Prone after this branch forked, and Error Prone 2.42.0 does not report on a Jabel-desugared record - it crashes, failing the whole compilation with a message attributed to line 1 of an unrelated file. #387 hit this one rung down and converted its five value types; the same thing was waiting here for `ProxyProcessor.ManifestOutcome`, `ManifestReconciler.Reconciliation`, `InFlightRegistry.InFlight`, `LivenessSettings` and `TestModeMainTest.Run`. Neither term can move - the root pom pins Error Prone at 2.42.0 because 2.43.0 needs a JVM this build cannot use, and Jabel serves the release 8 target - so the records lost, which is what the rest of the repository already does. `grep -rn "record "` now finds none. ## Two defects this merge introduced and then fixed - **The default lane ran a test that needs a profile.** Keeping BOTH sidecar-spawning tests - this branch's engine-backed `OneRecordThroughTheSidecarTest` and the stack's engine-less `SidecarHandshakeTest` - left the Kotlin and Scala exclusion property naming only the first, so the second ran in an ordinary build with no `sidecar-classpath.txt` and failed. Both are now behind one regex, which is also the arrangement that makes the pairing legible. - **Every proxy module's data record existed twice.** This branch keeps module-maturity and testing-evidence records in per-module `docs/data/*.d/` fragments; the stack, which never imported that mechanism, wrote them into the monolithic files. The merge kept both and `bin/check-docs-data.sh` went from green to 37 structural problems. The fragments win, and deleting the monolithic copies also deleted eleven rows asserting "NO RECORD HAS CROSSED THE WIRE" - true on a stack whose sidecar hosts no engine, false here. ## How it was verified JDK 17, macOS. Full default reactor `./mvnw --fail-at-end test`: **BUILD SUCCESS**, every module green including the conformance suite - so the `java-grpc` cell runs against a live engine for the first time, and `GrpcSpikeConformanceTest` answers all five scenarios over a real gRPC stream. The four self-retiring guards #387 and #390 left behind all pass, and **none of them was edited**: `TheEngineArrivingMustBringTheGrpcBindingTest`, `TheEngineArrivingMustBringTheGrpcCellTest`, `TheEngineArrivingMustBringTheForeignCellsTest` (which asserts about every registered language, not one) and `SelectorMatchingNothingFailsTest`. `bin/check-all.sh` real exit code 1, with three gates failing and all three inherited, established against the pre-merge tip rather than assumed: `check-inflight-tags` reports the same 53 problems before and after; `check-file-refs` was already red for the `@`-prefixed bridge imports that #378 fixes; and `check-branch-self-reference` did not exist on this branch before the merge - its findings are on notes inherited from `master` byte-identical. `check-proto-lint` and `check-proto-breaking` report CANNOT RUN for want of `buf`. **Not run locally: `-Dpc.foreignClients`.** It builds C++ and Swift inside containers, and this box is under a low-disk warning that names a container build as what would tip it over. CI's per-language rows are where that lane runs. ## Two findings recorded rather than fixed - `ParallelEoSStreamProcessorTest.processInKeyOrder` failed its own preamble sanity check once here - and the control arm reddened `master` HARDER: unmodified `master`, uncontended, failed all three parameterisations where this branch under concurrent load failed one. That refutes the existing ledger's "the base branch was green" row, and is the first evidence the flake is master-state. Recorded in `docs/inflight/test-processinkeyorder-sanity-check-races-the-first-poll.md`. - One `java-grpc` conformance run left its last record unsettled, 1 of 2 full-reactor runs and green 3 of 3 in isolation. `core` and `java-direct` were green in the same run, so the suspect is the engine's settle path under a full ceiling rather than the scenario. The scenario was NOT weakened - `docs/inflight/proxy-a-java-grpc-ceiling-run-left-its-last-record-unsettled.md`. ## What is deliberately NOT done here `Main#sessionServiceFactory` still returns `NoEngineSessionService`, so the production entry point hosts no engine. That is unit U10's, which #384 named as still having no PR, and it owns the drain that must land with the engine - plus eight cross-language handshake tests assert the `UNIMPLEMENTED` refusal by status code specifically. Recorded rather than faked, in `docs/inflight/proxy-the-production-entry-point-still-hosts-no-engine.md`. The engine is exercised end to end through `TestModeMain` regardless.
|
Two after-the-fact links, recorded because nothing else connects them and a search would not find This PR is very likely the fix for a real user report. #178 (mirror of I have posted the connection there rather than closing it, and said honestly what the evidence is: The sibling instance this PR deliberately did not fix is tracked and is NOT a blocker for that |
…e either, so the bound is not the cause The handoff note left one experiment running: whether a larger bound makes ShardManagerLincheckTest's stress arm find its violation. It does not - iterations(500) against a committed iterations(50), on the clean tree, missed both runs. That is worth more than a repeat of the earlier misses, because it closes the reading the lane's own notes make most attractive. The 3.4x machine-dependence recorded for these stress arms predicts that a bound priced on a fast machine misses on a slow one AND that spending more search finds it anyway. Ten times the budget finding nothing is inconsistent with that, and so is the control: reintroducing the defect changed nothing either. Neither more searching nor a present bug produces a violation on this box, which means the harness is not reaching the seam here at all rather than reaching it rarely. So the next agent is pointed at "why is it unreachable" rather than at a bound to tune - with the candidates named, cheapest first, and with the cheapest move of all stated plainly: confirm the arm still fires SOMEWHERE, on the machine that wrote it or in CI, before spending anything on the harness. Two of the three candidates only exist because #335 and #373 also rewrote ProcessingShard after the harness was written, which is exactly the thing the lane's note warns is not derivable by reasoning: re-run the lane, never reason about it.
…s there, the arm still does not The note handed over one cheap next move: confirm the Lincheck stress arm fires *anywhere* before spending anything on the harness. CI ran the lane on this PR's head for the first time and answered it. `ShardManagerLincheckTest.stressMustNotRediscoverTheShardTear` missed again, with the identical message, on `ubuntu-latest` - four cores, the opposite end of the range from the 32-core box the earlier evidence came from. The decisive part is not the miss but what passed beside it. On the same runner in the same JVM, `LincheckToolchainProbeTest`'s stress arm and `PartitionStateLincheckTest`'s stress arm both found their violations. The toolchain probe exists to answer exactly that question: Lincheck's stress strategy works there, the JDK is adequate, and the `-Plincheck` flags are landing. So two of the three surviving candidates are ruled out - a JDK or Lincheck version difference, and a host property such as core count - and one is left standing: this arm's operation set no longer reaching the seam it was written against, after #335, #336 and #373 each rewrote `ProcessingShard`. That also re-reads the earlier control. Reintroducing the check-then-act produced no violation either, which was ambiguous under the three-candidate list and is corroborating under this one: if the operations no longer reach the seam, restoring the defect at the seam changes nothing, which is what was measured. Records the cost correction too - 7m42s hosted against ~2m50s locally, 364s of it `WorkManagerLincheckTest`. The matrix entry's `timeout: 20` still covers it; the `~2m30s` in its comment does not. Evidence: run 33511221453, job 99867204560, head 58c8c7c.
…d the repeat in #410 Documentation only, from the review pass on this PR. No code change: the register-then-publish swap, the test and every assertion are untouched. Four items. 1. The cleared-suspicion block on maybeRegisterNewPollBatchAsWork handed the plain-long question to jcstress-poc's SeenSucceededOrderingProbes without saying whether anything runs it. Nothing does: jcstress-poc is absent from the root pom's <modules> list, so no reactor build reaches it, and no workflow and no bin/ script names it - `grep -rn jcstress .github/ bin/` returns nothing (verified before writing it). The engine AGENTS.md's cleared-suspicion rule requires the reopening condition to say plainly when no gate would catch it, and this one would not, so it now does. 2. The same block named offsetHighestSeen as the only plain-long residue. It is the harmless one. offsetHighestSucceeded is also plain, onSuccess read-modify-writes it, and getOffsetHighestSequentialSucceeded returns it directly whenever incompleteOffsets is empty - so a stale read of it is an offset committed to the broker, not bookkeeping. The dirty field's own javadoc already records jcstress measuring exactly that staleness, and what closes it is that field's volatile, not anything on this method. Both longs are now named, with which one has teeth. 3. docs/inflight/test-lincheck-lane-open-items.md retired the registration race with "No schedule can reach a container whose offset is absent". Too broad: WorkManager#onSuccessResult removes the offset (pm.onSuccess) before the container leaves its shard (sm.onSuccess), so a scanner in that window sees exactly that pair. The claim is now scoped to the registration path, with the post-completion window named and pointed at what actually refuses a claim there - the container's SUCCEEDED state (#335), which is candidate 1 in that same list. The retirement reasoning is unchanged and still holds. 4. docs/refactoring.md gains two lines. Under state/PartitionState.java: draft #410's restoreCompletedButUncommittedWork replays each record as addWorkContainer then addNewIncompleteRecord - the publish-then-register order this PR removed - and must take the same swap when it lands over this. It was called out in the defect-class sweep in the fix commit and in the PR body, but only there, where nobody editing that method would find it. Under the state package's white-box harness heading: the ConsumerRecords / EpochAndRecordsMap / registerWork idiom hand-rolled per test, which the simplify pass deliberately skipped as a pre-existing clone class, with why registerOneRecordAndTakeIt is not a drop-in and what a registerOneRecord-shaped sibling would have to do instead. Verified: PartitionStateRegistrationOrder370Test and WorkClaimStateMachineTest 12/12 green, bin/check-all.sh 16 passed 0 failed, bin/check-issue-refs.sh clean. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_019SVm2cT6ZgUyPMim9CukYK
Description
The work claim in
ProcessingShard#getWorkIfAvailablewas check-then-act; it is now one atomic transition against a per-attempt identity, so a record can never be delivered twice - and never redelivered against a retry delay that has since moved - however many threads select work.Selection used to be two steps:
isAvailableToTakeAsWork()evaluated three terms - not in flight, no success verdict, retry delay passed - andonQueueingForExecution()then acted, re-validating none of them. On this tree the gap is unreachable by architecture alone: the control loop is the only selector, so no completion can interleave with a decision. The direct-pull engine in development (onperf/engine-concurrency) gives every worker its own selector, and there the gap is real and was observed: a claim decided before another worker's completion acted on an already-completed record, the user function ran a second time, and the offset was committed a second time - 4 occurrences in 14,400,000 record completions. With assertions on it surfaces as anAssertionErrorfromPartitionState#onSuccess; in production it is silent.What changed. The plain
boolean inFlightand theOptional<Boolean> maybeUserFunctionSucceeded- the two questions selection asked, held apart - are now one atomic value with six states (AVAILABLE,IN_FLIGHT,IN_FLIGHT_SUCCEEDED,IN_FLIGHT_FAILED,SUCCEEDED,FAILED). A claim reads once, decides against that value, and compare-and-sets from that exact value: the check IS the act. The retry delay is deliberately not a state (nothing fires a transition when a clock passes a point), soFAILEDcovers waiting and due, andisDelayPassed()separates them at the moment of asking - ordered after the volatile state read so the previous holder's failure count and retry deadline are visible together.The state's name is not enough to compare against - the attempt is.
ExecutionStatevalues are enum singletons, so a record that leavesFAILEDand comes all the way back toFAILEDis==to where it started. A selector that readFAILED, found the delay passed, and was then descheduled would compare-and-set successfully against aFAILEDthat is a different failure, whose deadline the intervening attempt pushed out - and redeliver immediately, silently skipping the configured backoff. So the atomic holds an immutableExecution(state + attempt sequence) and every transition mints a fresh one: the compare, which is reference identity, cannot match across a cycle. That closes the whole ABA class in one place, rather than re-reading the retry deadline at each claim site and leaving whatever term is added next uncovered. Reported by Codex on review.Decide-then-act is a type, not a convention. The claim's state half is re-validated by the compare-and-set; its retry-delay half cannot be, because the deadline is deliberately outside the atomic, so it is checked exactly once at decision time.
decideClaim()evaluates the whole predicate and returns aClaim;actOnClaim(Claim)performs the compare-and-set against the exact occupancy that decision was made over.Claim's constructor is private toWorkContainer, so nothing can reach the act without having made the decision - there is no caller discipline to remember and none to get wrong. Raised on the automated review, where the seam was safe only because both call sites happened to be well behaved.Verification.
WorkClaimStateMachineTest(10 tests) pins the machine: the losing interleaving played out by hand with no threads - the concurrent reproduction needed millions of completions per occurrence, this is exact and runs in milliseconds - the ABA interleaving the same way, every state crossed on one record, and the properties a refused claim must keep (the verdict stands, the delivery count does not move). Control arm for the ABA tests: withExecution#transitionToreturning a canonical instance per state - the ABA restored -aClaimDecidedBeforeARenewedRetryDelayIsRefused()andeveryTransitionMintsAFreshOccupancy()both fail; with the fix, both pass. The transactional battle-test suite in #262 is the broader verification vehicle for duplicate-processing guarantees; this fix sits below it, in dispatch.Two lock-free loops became one JDK call each.
endFlight()andrecordVerdict()were hand-rolledwhile (true)compare-and-set loops. They are nowAtomicReference.getAndUpdate/updateAndGet, and each carries its own termination argument instead of leaving it to the reader.endFlightOf()is total and a literal fixed point on its own output - no state it produces is in flight, so a second round returns its argument unchanged.recordVerdict()'s bound is different and now says so rather than borrowing that one: at most two threads can write while a record is in flight (the worker recording its one verdict, the controller ending the flight) and neither can repeat, because a third write would need a claim that cannot be won from an in-flight state.SpotBugs. Two of the ten
AT_NONATOMIC_OPERATIONS_ON_SHARED_VARIABLEreports onparallel-consumer-corewere this PR's, and both were real -deliveryCountandnumberOfFailedAttempts.deliveryCount's own comment claimed the state compare-and-set ordered the increment; it did not, because the increment happens after the CAS. Both are atomics now, with hand-written accessors so the publiclong/intAPI (reached throughRecordContext) is unchanged. Nothing is suppressed. The other eight are pre-existing, in files this PR does not touch, and already recorded indocs/inflight/static-spotbugs-latent-findings.md.Same defect class, other instances - including the one not fixed here. The registration order in
PartitionState#maybeRegisterNewPollBatchAsWorkmakes a record selectable before its offset is registered: latent for the same single-selector reason, and now tracked as #370 rather than fixed, because closing the claim did not close it. The sibling duplicate-processing defect at the transactional layer (produce-lock double release) is #257 and shares no code with this. The full sweep, including the negatives and one latent hole nothing tracks, is in the "Defect-class sweep before merge" comment below.Relationship to the engine work. Adapted from
2e83185040onperf/engine-concurrency. The upcoming direct-pull PR depends on this landing first - it is the change that makes a second selector safe. The delivery-abandonment mechanism the branch version interacts with stays on that branch (itsWorkManagerside is not here), anddocs/inflight/core-a-lost-claim-means-two-different-things.mdrecords the follow-up this fix deliberately leaves open - now with the review's "is this a skip bug?" answered in it (it is not: a refused claim leaves the record in the shard and the next scan reaches it; what it skews is the slow-work signal).Checklist
docs/inflight/core-a-lost-claim-means-two-different-things.mdtravels with this change; the registration-order note is extracted to core: a record becomes selectable before its offset is in incompleteOffsets #370 and shrunk to a pointerdocs/features/- N/A - internal correctness hardening, no user-facing featureWorkClaimStateMachineTest, 10 tests, all green locally, ABA pair verified against a control arm