Repository navigation
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
…sured Every figure this project has published is throughput. Nothing measures how long a record actually spends inside Parallel Consumer, and no existing meter can be made to answer it: `user.function.processing.time` starts when the handler starts, and `WorkContainer#timeTakenAsWorkMs` starts when a delivery is claimed. Both begin AFTER the interval that matters - the time a record sits in the shards, fetched but not yet picked up. That interval is the one operational blind spot the metric set has. An instance holding 543 records against a configured maxConcurrency of 5,000 is ambiguous today between "the poller has not fetched more" and "they are fetched and queued behind other work", and the two call for opposite responses. Under a backlog, residence time is Little's law applied to PC's own buffers - buffered depth over throughput - so it is exactly the quantity that separates them, and it is what makes the cost of `DynamicLoadFactor`'s over-buffering visible as latency rather than as an argument. WHAT IS MEASURED. From the record leaving the consumer - the instant its `WorkContainer` is constructed from a poll batch - to the control loop finishing with that delivery. Produce time and time waiting on the broker are excluded: they are the environment, not the engine, and `ConsumerRecord#timestamp()` is there for anyone who wants them in. The end instant is the controller's rather than the worker's `succeededAt`, because the hop from the worker's mailbox back to the control loop is PC's own queueing and belongs inside the number. FAILURES AND ABANDONMENTS ARE SAMPLED TOO, which is the surprising choice and the one that makes the number honest. Arrival is stamped once, at construction, so a record that fails twice before succeeding carries both attempts and both retry delays in its final sample. Sampling only successes would make the distribution look best exactly when the engine is doing worst - which is the case this measurement exists to expose. It also means the sample count exceeds the record count whenever anything is retrying, so the denominator for a rate is `pc.processed.records`, never this timer's own count. A record whose partition is revoked mid-flight records nothing: its result is discarded and it will be redelivered elsewhere, so counting it would measure a rebalance rather than the engine. ONE TIMER, NOT ONE PER PARTITION, unlike the `processed.records` and `failed.records` counters beside it. A percentile histogram is a much heavier meter than a counter, and percentiles cannot be merged after the fact - so per-partition timers could never be combined into a consumer-wide p99, while an operator narrowing a bad p99 to a partition still has the per-partition counters and offset gauges. A per-partition breakdown, if wanted, has to be an additional meter rather than a re-tagging of this one. Flagged as a judgement call, not a settled one. THE TEST WAS BROKEN BEFORE IT WAS TRUSTED. Its sample count is asserted exactly (records + 1), not "greater than zero", because the plausible wrong implementation - sampling only successes - produces exactly the record count and sails through a loose assertion. Sabotaged by removing the failure-path sample: the test fails, `expected 1001L`. Its other assertion is that the maximum clears the one-second default retry delay, which is the single observable that a time-in-flight measure could not produce. Release-Note: New Micrometer timer `pc.record.residence.time` reports how long each record spends inside Parallel Consumer - from leaving the consumer to Parallel Consumer finishing with it, including client-side queueing and every retry, and excluding produce time and broker wait. Percentiles are published alongside the mean, and failures and abandonments are sampled as well as successes, so a retrying workload is visible rather than hidden. It is the only meter that distinguishes records not yet fetched from records fetched and queued.
Stamping the arrival instant in WorkContainer's constructor made it the first thing that constructor reads off the module, and WorkManagerTest built its containers from `mock(PCModuleTestEnv.class)` - a mock of a class that is already a test double, whose working MutableClock the mock replaced with null. The test errored where it had previously passed, which is a real regression in the test's double rather than in the test's subject: TreeMap ordering was never in question. Uses the class's own PCModuleTestEnv field instead, which needs no stubbing at all. Every sibling that mocks PCModule for this purpose stubs clock() by hand - ShardManagerTest#retryQueueOrdering and its two neighbours - so the unstubbed one was the outlier. DEFECT-CLASS SWEEP: four sites across every module's sources construct a WorkContainer from a mocked PCModule. Three already stub clock(); this was the only one that did not, and it is fixed. None outside parallel-consumer-core. Also files the PARTITION starvation note under `bug-` rather than `perf-`, since the prefix names the AREA and defects in product code belong under bug- where a reader scanning for them will look; and states, up front, the symptom rather than the mechanism - you configure maxConcurrency 24 and get 2 to 6 - plus that messageBufferSize is a workaround and not the fix, because a user cannot be expected to derive `max.poll.records x partitions` and the real fix (a prefetch target expressed in shard coverage) is not built. Full unit suite: every Java module passes. The one red module is the cross-language conformance prebuild, which fails on this machine's toolchains - go.mod uses a `tool` block the installed Go does not parse, and the Rust runner cannot execute the cached protoc (permission denied). Neither is reachable from anything on this branch.
Dependency Review✅ No vulnerabilities or license issues or OpenSSF Scorecard issues found.Scanned FilesNone |
astubbs
added a commit
that referenced
this pull request
Aug 25, 2026
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01MedgsxqrM8vjSt5ncuAo8g
✅ 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)
No new clones introduced by this PR. Powered by astubbs/duplicate-code-cross-check |
🧪🔒 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 |
|
astubbs
added a commit
that referenced
this pull request
Aug 25, 2026
…cut) Two ledger commits from the base - the wagon table with what is cut and what is next. No code; taken as-is to keep this branch current with its base. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_019HwDcZdfv1Ne4u3g6ThrG9
astubbs
marked this pull request as draft
August 25, 2026 13:58
astubbs
added a commit
that referenced
this pull request
Aug 31, 2026
…eatures, filed and bounded Fifth bit of the follow-up Codex strategy conversation (2026-08-29/30): the deliberate move away from adaptive-concurrency territory, mining what PC uniquely knows where Kafka semantics meet application execution. Eight new notes, each kept to the idea, its prior-art hook, and the caveat that decides whether it is honest: - core-record-semantic-tracing.md - the per-record "why did it wait" timeline; #359 stamps only the endpoints, and the per-gate stamps are the opportunity model's own instrumentation. - core-ordering-profiler.md - ordering tax, ordering-scope discovery, per-record critical path: one architectural profiler in three stages, near-free on the shard state #361 added. - core-retry-economics.md - blast-radius quarantine, per-function retry amplification, and storm detection (#333's failure-fraction inhibitor promoted to an operator warning). - core-capacity-fingerprinting.md - persist what the controller learns: warm starts, empirical workload models, semantic regression detection; open question is where the fingerprint lives. - perf-workload-replay-simulator.md - sanitised trace capture feeding the #362 harness as a capacity-planning simulator over real workloads. - release-certified-execution-semantics.md - publish #293's conformance matrix as a certification claim; a table that cannot show a cross certifies nothing. - core-function-manifest.md - split the idea on the positioning line: adopt the language-neutral manifest, leave "pc deploy" to platforms (embedded-not-cluster). - core-scheduler-canarying.md - A/B a scheduler over 1% of ordering domains; stratify or the comparison reports the sampling. Amendments: web-control-plane gains true lag as the fifth instrument (broker lag vs effective lag); the facades note gains the migration advisor with its honesty bound (observation mode sees keys and poll cadence, not handler time); the research program gains broker portability as question 6 with the reproduce-at-own-operating-point trap named; and the opportunity model gains the standing task the conversation closed on - inventory the boundary knowledge on purpose. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016PFA3CXkGc2m32cBJMAWTK
4 of 6 tasks
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
depends on #335 - the atomic work claim
Description
Record how long a record spends inside Parallel Consumer, end to end — the first latency measurement this project has had.
Every metric until now has described throughput or queue depth. Neither answers the question an operator actually has: how long does a record take to get through? A run that finishes in the same wall-clock time can have put one record in a hundred through a queue ten times as long, and nothing in
PCMetricswould have moved.pc.record.residence.timeis aTimer: the clock starts when theWorkContaineris built from a poll batch, and stops when the control loop finishes with that delivery — success, failure, or abandonment, including every retry in between. That last part is deliberate: a record that failed twice and succeeded on the third attempt took as long as it took, and a measure that restarted per attempt would flatter the engine exactly where users feel the pain.What it deliberately excludes: produce time and broker wait. Those are the environment, not the engine. (The bench harness carries a separate end-to-end measure that does charge the wait, precisely because residence cannot see coordinated omission — but that belongs to the harness, not to the library's own meter.)
Why it is worth a metric rather than a benchmark number: it is the only quantity that distinguishes an engine holding N records because it fetched N from one holding N because they are queued behind a slow key. Throughput cannot tell those apart, and it is the difference between "tune your concurrency" and "your key distribution is the ceiling".
Verification.
PCMetricsTestcovers the timer's registration and that it records on the completion path.WorkManagerTestandWorkClaimStateMachineTestgreen alongside (36 tests total locally) — the second matters because the arrival instant is stamped inWorkContainer's constructor, which #335 also touches.One test fix travels with this change, and is not incidental: stamping the arrival instant made the constructor read the module's clock for the first time, which surfaced a latent defect in a test double —
WorkManagerTestbuilt containers frommock(PCModuleTestEnv.class), a mock of a class that is already a test double, whose workingMutableClockthe mock replaced with null. The fix uses the class's ownPCModuleTestEnvinstead, which needs no stubbing. The test that errored (treeMapOrderingCorrect) was never testing the clock; it was testing TreeMap ordering, and its double had been wrong all along.Provenance. Extracted from the engine-performance campaign on
perf/engine-concurrency. Stacked on #335 because both touchWorkContainer; the branch's verdict-free abandonment path (which also records residence) belongs to the proxy work and is deliberately excluded here.Checklist
PCMetricsDefcarries the metric's contract in javadocdocs/features/-docs/features/record-residence-time.yamlPCMetricsTest, plus the test-double repair described above