Repository navigation
Conversation
…rdered concurrency Requirements-only plan for running Parallel Consumer in a sidecar that an application's worker processes talk to over loopback, so a runtime other than the JVM gets concurrency beyond partition count while keeping per-key order. Python ships as the flagship client library, Go as a falsification test written from the specification alone. The framing that shaped it: Share Groups (KIP-932, GA on Kafka 4.2) already give any runtime per-message ack and partition-decoupled scaling, so the claim is deliberately narrow. What survives that baseline is what README.adoc already names - key-level ordering with concurrency beyond partition count, and no processing clock. Everything else in the contract is held to that bar. Nine decisions carry their rejected alternative inline. Three are worth naming: - A sidecar rather than a shared multi-tenant server, on strategy grounds: the approach is "nobody's permission to deploy", and a shared server forfeits it while competing head-on with Kafka REST Proxy. - One sidecar serving many worker processes rather than one per process. PC needs no shard-to-worker affinity for this - ProcessingShard already keeps at most one record per shard in flight under ordered modes - and a protocol where one connection owns all the work cannot be extended to several later. - Demand-driven flow control, with credit riding on the acknowledgement. The push-with-advertised-capacity model it replaced needed a trusted capacity number on an unauthenticated surface, an aggregate in-flight target recomputed on every join and leave, and an unenforceable rule that advertised capacity exceed real concurrency. Choosing the pull shape deleted two requirements outright rather than patching them; R30 and R34 are gaps, not renumbered. Two review rounds with six independent reviewers each ran against it, and the codebase claims were verified separately. That surfaced two things the plan now records rather than assumes: PC has no verdict-free way to return work, so returning a record when a worker's connection drops needs a core-side path that leaves the failed-attempt count untouched; and core's scheduled dead-letter queue (#149) has a one-line done-condition and no design, so the proxy adopts its destination semantics but defines its own triggers. Three blockers travel with it, all of them numbers or evidence nobody has yet: the Go client's effort budget, the latency multiple the first success criterion is judged against, and whether the users who asked for a Python client need key ordering at all or the parallel consumption Share Groups now supply. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…ew module has to register The plan is the contract; this is the branch note for what no command can answer. Three things worth having written down before anyone builds the module: The data records - a module-maturity row, its testing-evidence entry, and a feature record - land in the PR that lands the module, never earlier. The corpus already carries that scar: two feature records were removed rather than shipped because their modules were not in pom.xml, so their Maven coordinates could not resolve. check-docs-data.sh validates the schema but does not cross-check the module list against pom.xml, so a missing row is silent. Registration is six places, and two of them are the same file with different separators - the space-separated and comma-separated duplicate-detector lists in maven.yml. Missing one is the documented failure mode. And release.target is 8: the build compiles Java 17 source to Java 8 bytecode, so modern networking APIs are invisible to a wire-protocol module. The mutiny pom already models the override and records why it has to be deliberate - at the wrong target that module compiled happily and failed at runtime, because nothing in the build detects it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard. |
Dependency ReviewThe following issues were found:
|
✅ 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
|
| PR | Base | Change | |
|---|---|---|---|
| Clones | 123 | 88 | :face_with_raised_eyebrow: +35 |
| Duplicated lines | 2018 | 1257 | :face_with_raised_eyebrow: +761 |
| Duplication | 1.08% | 1.02% | 🫤 +0.06% |
| Rule | Limit | Status |
|---|---|---|
| Max duplication | 2% | ✅ Pass (1.08%) |
| Max increase vs base | +0.1% | ✅ Pass (+0.06%) |
⚠️ 36 new clones introduced
- 8 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-java/parallel-consumer-proxy-client-java-grpc/src/main/java/bz/stub/parallelconsumer/client/grpc/GrpcParallelConsumerClient.java:10<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-java/parallel-consumer-proxy-client-java-grpc/src/test/java/bz/stub/parallelconsumer/client/grpc/SessionEndTest.java:7 - 9 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-java/parallel-consumer-proxy-client-java-direct/src/test/java/bz/stub/parallelconsumer/client/direct/DirectParallelConsumerClientTest.java:9<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-java/parallel-consumer-proxy-client-java-direct/src/test/java/bz/stub/parallelconsumer/client/direct/DirectSpikeConformanceTest.java:10 - 7 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-java/parallel-consumer-proxy-client-java-api/src/test/java/bz/stub/parallelconsumer/client/conformance/SpikeConformanceTest.java:106<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-java/parallel-consumer-proxy-client-java-api/src/test/java/bz/stub/parallelconsumer/client/conformance/SpikeConformanceTest.java:77 - 23 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-java/parallel-consumer-proxy-client-java-api/src/test/java/bz/stub/parallelconsumer/client/conformance/SpikeConformanceTest.java:192<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-java/parallel-consumer-proxy-client-java-api/src/test/java/bz/stub/parallelconsumer/client/conformance/SpikeConformanceTest.java:90 - 16 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-scala/src/test/scala/bz/stub/parallelconsumer/client/scaladsl/OneRecordThroughTheSidecarTest.scala:159<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-scala/src/test/scala/bz/stub/parallelconsumer/client/scaladsl/SidecarHandshakeTest.scala:96 - 13 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-kotlin/src/test/kotlin/bz/stub/parallelconsumer/client/coroutines/OneRecordThroughTheSidecarTest.kt:94<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-kotlin/src/test/kotlin/bz/stub/parallelconsumer/client/coroutines/SidecarHandshakeTest.kt:92 - 11 lines:
parallel-consumer-proxy-protocol/src/test/java/bz/stub/parallelconsumer/proxy/protocol/ProxyProtocolRoundTripTest.java:74<->parallel-consumer-proxy-protocol/src/test/java/bz/stub/parallelconsumer/proxy/protocol/ProxyProtocolRoundTripTest.java:59 - 8 lines:
parallel-consumer-proxy-protocol/src/test/java/bz/stub/parallelconsumer/proxy/protocol/GoldenSessionFixture.java:10<->parallel-consumer-proxy-protocol/src/test/java/bz/stub/parallelconsumer/proxy/protocol/ProxyProtocolRoundTripTest.java:10 - 9 lines:
parallel-consumer-proxy-protocol/src/test/java/bz/stub/parallelconsumer/proxy/protocol/GoldenSessionFixture.java:21<->parallel-consumer-proxy-protocol/src/test/java/bz/stub/parallelconsumer/proxy/protocol/ProxyProtocolRoundTripTest.java:20 - 16 lines:
parallel-consumer-proxy-protocol/src/test/java/bz/stub/parallelconsumer/proxy/protocol/GoldenSessionFixture.java:132<->parallel-consumer-proxy-protocol/src/test/java/bz/stub/parallelconsumer/proxy/protocol/GoldenSessionFixture.java:101 - 8 lines:
parallel-consumer-proxy/src/test/java/bz/stub/parallelconsumer/proxy/config/ConfigureHandlerTest.java:444<->parallel-consumer-proxy/src/test/java/bz/stub/parallelconsumer/proxy/config/ConfigureHandlerTest.java:145 - 10 lines:
parallel-consumer-proxy/src/test/java/bz/stub/parallelconsumer/proxy/config/ConfigureHandlerTest.java:495<->parallel-consumer-proxy/src/test/java/bz/stub/parallelconsumer/proxy/config/ConfigureHandlerTest.java:461 - 17 lines:
parallel-consumer-proxy-protocol/src/main/proto/parallelconsumer/proxy/v1/proxy.proto:230<->parallel-consumer-proxy-protocol/src/main/proto/parallelconsumer/proxy/v1/proxy.proto:153 - 14 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-typescript/src/generated/google/protobuf/duration.ts:97<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-typescript/src/generated/google/protobuf/timestamp.ts:127 - 30 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-typescript/src/generated/google/protobuf/duration.ts:113<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-typescript/src/generated/google/protobuf/timestamp.ts:143 - 12 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-typescript/src/generated/google/protobuf/duration.ts:149<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-typescript/src/generated/google/protobuf/timestamp.ts:179 - 31 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-typescript/src/generated/google/protobuf/duration.ts:164<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-typescript/src/generated/google/protobuf/timestamp.ts:194 - 45 lines:
parallel-consumer-proxy-clients/parallel-consumer-proxy-client-java/parallel-consumer-proxy-client-java-direct/pom.xml:47<->parallel-consumer-proxy-clients/parallel-consumer-proxy-client-java/parallel-consumer-proxy-client-java-grpc/pom.xml:76 - 11 lines:
docs/data/testing-evidence.d/parallel-consumer-proxy-client-cpp.yaml:43<->docs/data/testing-evidence.d/parallel-consumer-proxy-client-swift.yaml:63 - 8 lines:
docs/data/module-maturity.d/parallel-consumer-proxy-client-ruby.yaml:16<->docs/data/module-maturity.d/parallel-consumer-proxy-client-rust.yaml:14
...and 16 more
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 |
…orrect two premises it was resting on Adds the Planning Contract, twelve implementation units in three phases, a verification contract and a definition of done. The Product Contract is unchanged - no requirement, actor, flow or acceptance example was added, removed, renumbered or reworded. Two things the plan asserted turned out to be false, and both were load-bearing. It recorded that a streaming RPC framework is materially harder to build under GraalVM native image than plain Vert.x, and flagged the claim as unverified. Verified: it is false. gRPC's reachability metadata tracks its current release and covers bidirectional streaming in a native test that runs in CI, while Vert.x removed its own metadata in 2026 and the shared repository's coverage lags its release and never exercises an HTTP or WebSocket server. Native image does not discriminate between the candidates; toolchain weight against free HTTP/2 flow control does, so the transport is now settled by a spike whose tiebreaker is the effort figure the Go client already had to record. It also said later languages would arrive as community clients proving themselves against a conformance suite. Nobody decided that - it was an unattributed sentence that then became the main argument against gRPC, since every citation against gRPC is really evidence about external implementers declining a heavy dependency. Client libraries are generated and live in this repository, arriving by pull request, so the build owns the toolchain and CI covers every client. No conformance suite is scoped. The units carry what the research found rather than leaving it to be rediscovered: the module must override release.target because the build compiles to Java 8 bytecode and nothing detects the mistake; the verdict-free work return is a core change mirroring the existing stale-work branch, and the in-flight counter returning to baseline is its acceptance criterion because drift there stalls the consumer silently; the integration-test package must be named integrationTests or failsafe never runs it; and the module lands depending on nothing from feats/web-gui, because the duplicate-code gate has less headroom than a verbatim lift would consume. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…degen and the lifecycle a unit that builds them Six review lenses read the half of this plan that two prior rounds never reached: the Planning Contract, the twelve units, the verification contract and the Definition of Done. The Product Contract had been reviewed twice; none of this had been reviewed at all. The sharpest finding is a defect, not a gap. KTD6 said to override the per-pass target hook with the outstanding credit total. But calculateQuantityToRequest computes `target - numberRecordsOutForProcessing`, and ExternalEngine's own override returns an absolute figure - so credit, which is already net of the records a worker holds, gets the in-flight count subtracted from it a second time. A worker of concurrency four holding four records that reports one back leaves credit 1 and in-flight 3: delta is -2 and the sidecar never hands out another record. Two reviewers found this independently against the source. The hook has to return credit plus in-flight. The pull design is untouched and correct; only the arithmetic mapping it onto core's hook was wrong. Two units were missing outright. U6 and U8 were told to generate their transport from a specification that no unit produced - R13 belonged only to the spike, whose code is discarded by design. And R11, R12, R24 and R25 belonged to no unit at all, so nothing started, supervised or drained the sidecar that U8 and U9 both need running. Those are now U2 and U7, and every requirement and acceptance example is covered by some unit - asserted, not asserted-to-be. The spike was deciding on the wrong evidence. It could have picked a transport R29 disqualifies, since raw framing declares no authority to reject on, and it recorded no native-image result although R25 says that path wins any conflict. Both are now pass/fail gates ahead of the effort figure, and the primary criterion is whether a generator emits the client transport at all - that is the slope seven languages ride on, where a one-off hand-build measures only the intercept. Its tiebreaker cited R16, which is measured two phases later and cannot inform it. Batch size and R3 turned out to be incompatible in kind: core gives one completion verdict per batch, R3 requires per-record outcomes. R9 drops the option. Maximum concurrency no longer bounds concurrency under the credit model but still gates the broker poller, so it had quietly become a buffer ceiling; whether to withdraw it entirely is left open rather than settled here. Three controls the Product Contract requires were listed against a unit that never implemented them - credential redaction, failure-reason sanitisation on the retry path, and validating that a result report matches a record the connection actually holds. The last one matters most: each spurious report drives one in-flight decrement, and drift in that counter is this repository's documented silent-stall signature. The clients had no CI gate anywhere, which is the one thing KTD2's in-repo model rests on. ASM3 is settled - users confirmed the demand is for key-ordered concurrency, not the parallel consumption Share Groups supply - so the four places still calling it the highest-risk open assumption, including a Definition of Done item that could not be closed as written, now say so. The stale transport paragraph repeating the falsified native-image claim is gone; KTD1 already records that verification. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…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 maybeUserFunctionSucceeded, 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. maybeUserFunctionSucceeded was never reset, so a record that failed, came back, and was then abandoned still carried its old false verdict, took the failure path and earned a retry delay - the precise case a worker killed on its second attempt hits. onQueueingForExecution now clears both the verdict and the marker, because a redelivery is a fresh attempt. Verified by removing the new branch and confirming all six behavioural tests fail, so they are not passing for an unrelated reason. Full core suite green at 329 tests. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…uplicate cannot end a live one (confluentinc#154) Review of the previous commit found the duplicate-return guard does not hold, and verified it by running probes rather than by inspection. The guard was `wc.isNotInFlight()`, on the theory that a record which has already been returned is no longer in flight. But an abandoned record is immediately re-selectable by design, and controlLoop drains the mailbox and re-selects work in the same iteration. So a duplicate arriving any later than that single drain finds the record in flight again and sails straight through: it ends a live delivery's flight, decrements the in-flight counter a second time, and leaves an availability increment orphaned on a shard whose entry has since completed. Observed: outForProcessing at -1, and awaitingSelection reporting three records when one of the three had already succeeded. That is the silent-stall signature - isRecordsAwaitingProcessing stays true, so drain never reaches CLOSING. The previous commit's own test constructed the one interleaving where the guard works - no re-selection and no re-mark - and so certified it. A boolean on shared mutable state cannot express "this return is stale", because the container has no delivery identity: a late return for delivery n is indistinguishable from a live one for delivery n+1. So deliveries are now counted, and the abandon marker records which delivery it was raised for. A return whose marker names a delivery that has already ended is ignored outright. Callers capture the delivery at dispatch, which is what makes the marker meaningful - reading it at return time would just relabel a stale return as live. That check runs before the revoked-partition branch, not inside the abandoned path, because that branch decrements unconditionally and was the second way to drift the counter - reachable exactly when a worker fleet dies during a rebalance, which are correlated events. onAbandonedResult is no longer public: it skips both checks, and an increment landing on a revoked shard is never swept. Engines go through handleFutureResult, which is documented control-thread-only, since the counter is a plain int. Also collapses two byte-identical shard methods into markAvailableAgain - ProcessingShard cannot distinguish a failure from an abandonment, since the retry-queue insert that actually differs lives in ShardManager - and registers "verdict" in CONCEPTS.md, whose in-flight entry described only two resolutions. Verified by removing the new check and confirming all three new tests fail, including the revoked-partition one. Full core suite green at 331. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…wo gates that decision still owes (confluentinc#154) The spike KTD1 scheduled was to choose between gRPC bidirectional streaming and protobuf over hand-rolled framing, with recorded engineer effort as the tiebreaker. That tiebreaker cannot be measured by the party that would run it here, and the deciding evidence turned out to be qualitative and already in hand. With the native-image objection falsified, the comparison is one-sided. HTTP/2 gives per-stream flow control for free, which is R23's credit mechanism living in the transport instead of in our code. `:authority` is a declarable connection authority, so R29's allowlist has something to reject on - raw framing would have needed a connect frame invented for the purpose. And a .proto is the machine-readable schema R13 asks for, with mature generators for both target languages, which is what KTD2's in-repo generated clients ride on. Two things the spike would have proven empirically are still owed, so U1 survives as those two gates rather than as a bake-off: that the server can reject on a declared authority, and that the hand-out loop builds as a native image. R38 forbids fixing a transport mistake by removal once clients ship, so the second is unrecoverable if deferred. R1 moves to U5, which owns the credit mechanism that discharges it. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…hanges) Directory move only, so every path is 100% similar and git's exact-rename detection cannot fail on it. The content edits follow in the next commit. Generated by bin/rename-packages.sh.
Text edits only. No file moves in this commit, so it cannot dilute the rename detection in its parent. Generated by bin/rename-packages.sh.
…nt that fails silently (confluentinc#154) Both gates KTD1 owed are cleared, on a throwaway probe rather than desk research: gRPC 1.73.0, protobuf-java 3.25.5, GraalVM CE 25.0.2. The probe is discarded, so these outcomes are its only output, and U2 is unblocked to author a schema against gRPC. R29 holds. A ServerInterceptor reads the connection's declared :authority from ServerCall.getAuthority() and rejects an unlisted one with PERMISSION_DENIED. What R29 actually needs is that the rejection precedes application handling, so that is what was measured: closing the call in interceptCall and returning a no-op listener means the service method never runs, and across a rejected connection both the service-invocation and application-message counters were unchanged. Asserting on counters rather than on the status code is the difference between proving the ordering and assuming it. A connection declaring no authority is accepted, which is the default R29 asks for. R25 holds. The bidirectional-streaming hand-out loop builds --no-fallback and completes the full credit/record/outcome cycle as a 45MB binary under Substrate VM. The hints are the part worth keeping. protobuf-java ships no native-image metadata at all; grpc-netty-shaded ships its own; and the GraalVM reachability repository contributes exactly one entry, for java.time.Instant, resolving gRPC 1.73.0 to its 1.69.0 directory because it has no entry for the current release. Direct generated-API use - the entire hand-out path - needs no protobuf hints, which is why the first native build passed with no hand-written config. Descriptor-driven reflection is the trap. getDescriptorForType, getField, TextFormat, JsonFormat and DynamicMessage all reach GeneratedMessageV3.FieldAccessorTable, which reflects on the generated accessors. Unregistered, the build stays green and the binary runs; the call fails only when that path is first exercised. The probe only found it because it was extended to touch the descriptor path deliberately - the two gates as specified would have passed over it in silence, and U7 would have inherited it. Registering each generated message class and its $Builder with allDeclaredMethods and allDeclaredFields fixes it, verified by rebuilding and re-running rather than by assertion. That is the same shape as the release.target trap already recorded here: compiles happily, fails at runtime, build cannot detect it. So U2 owes a deliberate decision about whether the schema's consumers may touch the descriptor path at all. Also records why native-image needs a C toolchain - it shells out to gcc rather than linking anything itself, so its absence surfaces at the link step and reads like a fault in the code being compiled - and corrects the blocked-on section, which still listed ASM3 as open after the plan settled it. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…everywhere a module has to be registered (confluentinc#154) U3. An empty module that builds, tests and satisfies every gate, so U2 has a pom to hang the proto and codegen plugin on. No transport dependency and no schema yet - that is U2's unit, and adding either here would put gRPC in the tree before the module itself is known to be sound. release.target is overridden to 17. The project-wide target is 8: the build compiles Java 17 source to Java 8 bytecode via Jabel, so the modern networking API surface a wire-protocol module needs is invisible to it. The override has to be deliberate because the failure it prevents is not a compile error - the release level constrains the platform API, not what javac reads off the classpath, so a module set too low compiles happily and ships bytecode claiming a JVM it cannot run on. mutiny records the same trap after paying for it. commons-lang3 is declared explicitly at test scope rather than leaned on transitively: core's test-jar reaches for it and test scope is not transitive, so inheriting it today would break the moment core stopped using it. Registered in all three places a new module has to be, since missing one is the failure mode here and two are easy to miss: the root pom's <modules> before parallel-consumer-examples, which stays last; and both duplicate-code detector lists in maven.yml, which take the same value with different separators - one space-separated, one comma-separated. The three YAML data records a user-visible module owes are deliberately not written here. U14 owns them, and the Definition of Done lands them in the same PR as the last unit, so the rule that they accompany the module is satisfied by construction rather than by writing them early against a module that carries no content yet. The arch test uses bz.stub.parallelconsumer.proxy, not the io.confluent path the plan's Files line still names - that line predates the package rename now on this branch. Verified: a full reactor build including this module passes, with the copyright check and the reactor-convergence enforcer both green, and this module's arch test runs 3/3. The plan's literal gate - the same build with tests - cannot pass on this branch for a reason that predates this commit and is unrelated to it: core's own TestConventionsArchTest fails on FilteredTestContainerSlf4jLogConsumer, reproduced on a clean tree with none of this change present. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…d record what the ideation settled (confluentinc#154) The credit-based interaction model is superseded. The sidecar registers as PC's user function and returns a future completed when a worker reports - an ExternalEngine, the same seam vertx and reactor already use. The plan on disk has not been rewritten yet and still reads as though the ledger were live, so the inflight note carries the current position and says where the two disagree. The frame that came out of it is that the wrapper IS the layer. Every language gets the same client wrapper and to the user it looks native; Java's is that same layer with one fewer hop, sitting directly on the engine with no protobuf underneath. One client model with a missing layer in the Java case, not two architectures. Two decisions settled deliberately. Workers never produce to Kafka directly - output travels back through the engine and the sidecar produces, which keeps exactly-once entirely on the JVM side behind one producer and one epoch check. And fencing is Kafka's own EoS model borrowed rather than invented: each delivery carries an epoch, a dead client is fenced and the epoch increments when the record is handed on. PC already implements the mechanism internally, so what the protocol adds is making the epoch explicit on the wire. The note states the boundary plainly - this fences reports and Kafka-side effects, and cannot fence a worker's external ones, which is true of any at-least-once system. Compiling PC to a native shared library is recorded as a dead end rather than dropped silently. It was an agent's extrapolation during ideation, not a proposal, and it was ranked first and misattributed before being corrected. Both things that undercut it on its own terms are written down so the next session does not rediscover them: the native-image gate this branch cleared produced an executable rather than a --shared export, which differs in entry-point surface, isolate and thread-attach semantics, GC coexistence with a foreign runtime and callbacks re-entering from foreign threads, none of which was tested; and Temporal's Go and Java SDKs are independent implementations rather than bindings over its shared core, so the precedent usually cited is thinnest exactly where the idea was boldest. The ideation artefact keeps the full record - six survivors across five axes, and a rejection table with a reason for every cut, including the verdicts overruled because a fresh-context verifier scored candidates against the stale plan. Its strongest finding is not an idea but a defect, and it needs fixing regardless of which model wins: R27 returns a closed connection's records immediately, which keeps the sidecar's bookkeeping consistent but does nothing to stop two physical workers running the same key's user code concurrently during a transient disconnect. Every host-side invariant reads green while the guarantee the product is sold on is violated in fact. Two notes are opened for work that outlives this branch: a standing comparison against the systems that already solved these problems, and the observation that a PC-backed sink connector delegating over this same boundary would let people write Connect connectors in a foreign language. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
… which requirements they move (confluentinc#154) Completes the ideation. Three of the six survivors are now decisions rather than options, and the artefact marks them so a reader can tell what is still a fork. Supervision resolves to policy in the sidecar, mechanism in the client. The app starts the sidecar; the sidecar decides how many executors and when, and says so in a message; the client library - which already holds the user's function - spawns them with its own language's mechanism. The engine never learns what the function is. That beats both options it replaces. Against client-supervises-sidecar it deletes the bind-race election, because one app process starts one sidecar and there is nothing to elect, and it deletes the detached process group, which existed only so the sidecar could outlive whichever worker won the race. Against sidecar-spawns-workers it avoids the sidecar needing to know how to start the user's program - deploy-time configuration that fails late - and avoids forcing the user function to be an importable name rather than a closure, since a sidecar can only spawn a command. Configuration is code, delivered connect-time over the protocol: no config files, no environment variables, no shell. Generated stubs are not a translation layer. Max concurrency therefore keeps the meaning it already has - set by the app, sent to the engine, used as the in-flight ceiling, which is exactly what maxConcurrency times batchSize already means to an ExternalEngine. No core change. Two earlier proposals to withdraw the option were both wrong. App shutdown shuts the sidecar down: drain within the bounded timeout, commit, leave the group. One backstop remains, because on SIGKILL nothing runs - the sidecar must watch its parent and exit itself, or a JVM leaks while still holding group membership and partitions do not rebalance until the session timeout. Four requirements move, and the artefact tabulates them rather than leaving the drift implicit: R9 (options now arrive over the protocol), R33 (the no-connection grace period is unnecessary once the app owns the sidecar's lifetime), R35 (credentials may travel the protocol, on the stronger argument that the sidecar takes exactly one connection from the process that spawned it and does nothing until configured, so nothing sits on disk or in /proc), and KTD7 (superseded outright). Also corrects a claim made during the run: KTD7 was reported as a session-settled user-directed decision. It carries no such label - only KTD1 and KTD2 do among the KTDs - so superseding it reverses nothing the user directed. Still open, and marked as such: the reconnect manifest, which is a correctness gap rather than a preference; wave dispatch; and the liveness lease. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
… and delete the flow control it had rebuilt (confluentinc#154) Replaces the credit-based plan. The sidecar registers as PC's user function and returns a future completed when a worker reports, which is the ExternalEngine seam vertx and reactor already use. Clients are one wrapper layer in every language, with Java's being that same layer minus the protobuf hop. The first draft of this plan had reconstructed flow control under a new name. ExecutorCountPolicy derived a count from observed report concurrency and pushed it back whenever the observation moved - a closed feedback loop between proxy and client capacity, which is the deleted ledger with different vocabulary. The Definition of Done could not catch it, because that guard was a grep over five words and this named none of them. The count is now a pure function of connect-time configuration sent once, the guard is restated as the property rather than the words, and the general shape is recorded as a dead end: no proxy-side quantity derived from observed client behaviour is fed back. The same failure had been exported across the boundary. Deleting host-side flow control turns the gap between the in-flight target and the executor count into a queue inside the client library, and nothing specified it - so ten authors would each have invented depth, ordering, lease semantics and shutdown behaviour. It is now decided once, carried in the specification and the conformance suite, and the spike dispatches two records to one executor so the rule is exercised before the freeze rather than discovered ten ways at wave one. Three things would have made the fan-out unsafe. The engine units did not precede the reference sign-off, so nine agents could have inherited a signed-off client tested against an engine that could not answer its messages - and a red job would then have been ambiguous between a wrong client and an unbuilt feature, which is the noise the seeding rule exists to remove. The module-maturity and testing-evidence records are single files that every client must append to, which falsified the promise that agents never edit a shared file; they become per-module fragments merged at check time. And the verification commands did not select their modules: -pl takes a path or :artifactId, every client module is nested, and selecting an aggregator builds no children - so the spike's own gate would have passed green without compiling either transport. Reproduced by running maven, not by reading it. Security is deferred to post-v6 but recorded rather than dropped. The credential posture was reversed on the argument that the proxy takes one connection from the process that spawned it; over loopback TCP it can verify first, never parent, and the slot is handed back on every reconnect with nothing re-authenticating the peer. The Kafka property map arrives unallowlisted and Kafka instantiates class-valued properties reflectively. Nothing says how a client locates the binary it hands credentials to. The available mitigation - a nonce down the pipe the parent-death watchdog already establishes - is named and declined for now. Exactly-once is deferred with it: ExternalEngine throws when transactional commit mode is set, so it is unreachable today, and core modification is now sanctioned to reach it rather than worked around. R7 stops citing a dead letter queue core does not have. Adds the demo surface - a container per language, real broker by default and mock for hosted, one shape across all of them - the per-language idiomatic review, and a unit for STRATEGY.md, which this work has changed as much as it changed the plan. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Brings the branch up from 20 behind, which it needed for the arch-test fix in 71a306c - core's TestConventionsArchTest was failing on FilteredTestContainerSlf4jLogConsumer, and that failure blocked every unit on this branch whose verification is a full build. All 22 conflicts came from master's test refactoring in #265 landing on files this branch's package rename had moved. bin/rename-packages.sh is byte-identical between the two, so the rename itself was not the source and re-running it would have been a no-op. Six files were accepted as deleted, each confirmed absent from master's entire tree first rather than assumed: two .idea run configurations, SampleTestingFailsafePluginInclusionCore, JavaEnvTest, StringTestUtils and the core test junit-platform.properties, the last of these removed by #265 so the core tests jar stops configuring other modules. Every conflicted Java file took master's version, on evidence rather than convention. Comparing the two sides directly showed master had replaced sleep-based timing with explicit clock advancement and split MockConsumerTest into a base-class family, while this branch carried only the pre-refactor upstream code with renamed packages. A first measurement suggested this branch held substantial unique content; that was an artefact of diffing against an explicit post-rename path, where git cannot see the rename source and reports the whole file as added. README.adoc is generated, so it was not resolved by hand: master's template gained the upgrade-from-upstream guide and its rename-packages freeze markers, this branch had only the one summary line that guide supersedes, so both the template and its generated output take master's version. Verified after resolution: the verdict-free work return and its tests survive, the proxy module and this branch's plan, ideation and inflight notes survive, master's clock refactor survives, and no Java source references the old package.
…ualifier attached to its claim (confluentinc#154) U34. The language-proxy work changed the strategy, not only the plan, and a plan is the wrong place to keep that. The register is the point. This whole direction is an experiment, and the track says so in its first line: one claim has evidence behind it, the rest are things being tried, and there is no third category. Nothing here predicts it works, commits the product to becoming a general Kafka client, or claims a market. The proven part is architectural. The wrapper is the layer and Java is the degenerate case of it - one client model in every language, Java's having one fewer hop because it sits directly on the engine. From that follows the one line worth stating: our currency costs a version bump and librdkafka's costs a reimplementation, structurally and permanently, because the boundary sits at the process edge rather than the language edge. The qualifier travels in the same breath, because separating them would be claiming the experiment's outcome as its premise: Parallel Consumer is not current with Kafka today. The architecture can be current, the product is not yet, and only the first of those is proven. Everything else is labelled as what it is. The addressable segment is an observation about fit rather than a sizing - a possibility for users who need to be more current than librdkafka is, who skew sophisticated, which is also the population most willing to run a sidecar. The Share Groups comparison carries its cost in the same sentence as its advantage. Wrapping the core client APIs is a staged possibility with admin first and plain consumer last or never, not a plan. And whether this is Parallel Consumer for other languages or the Kafka client for other languages is recorded as earmarked for investigation with a cheap probe attached, not as a direction. The client-side sub-broker framing is unchanged and still correct; this is a track inside it. Nothing in the argument for putting the queue in the client was ever about the JVM - but every implementation of it has been, which is a limit of the library rather than of the idea. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…it per module (confluentinc#154) bin/check-docs-data.sh validated each data file against the schema but never noticed a module the data failed to mention: a module in a <modules> list with no maturity row passed clean, so the corpus claimed a completeness it did not have. The reverse direction has already cost the repo - two feature records were removed rather than shipped because their published Maven coordinates could not resolve. Both directions are now enforced, and the forward walk descends into nested aggregators' <modules> so a module declared only in the examples tree - or tomorrow's client tree - cannot slip through one level down. The corpus also becomes fragment-readable before it has to survive eleven concurrent client waves. docs/data/module-maturity.d/<artifact>.yaml and docs/data/testing-evidence.d/<artifact>.yaml carry one module's record each, merged into the root file's corpus by the checker; the root files keep only the repo-level preamble. The filename must match the artifact - the filename is the ownership claim - so two waves cannot write the same module's data, and the same module keyed in both a fragment and the root file fails loudly rather than one copy silently winning a YAML list auto-merge. Deferral lives inside a module's own fragment as `deferred: {reason, lifted_by}`, never on a shared list every deferring wave would edit, and every current deferral is named in the output on green runs too, so a deferral cannot quietly become an omission. Seven deferral fragments make the current tree green honestly: parallel-consumer-proxy (mid-build on this branch) and the five example submodules (sample code, never published). A maturity row's evidence_id must now resolve to a module_evidence id in the merged testing-evidence corpus; the plan believed this check already existed, but nothing enforced it. Its `feature` path was in contrast already resolved by the prose-reference walker, so that check gains only a pinning fixture. module_evidence artifacts are deliberately NOT reverse-checked against the reactor: the streams-alpha and connect-alpha entries describe planned modules on purpose, as the evidence targets of the staged rows in docs/data/staging/, and the maturity corpus is the surface that claims a module exists. Packaging decision, which the per-wave client records depend on: a client's user-facing package is represented on the maturity row, as optional package_ecosystem and package_coordinate fields - not in feature records. The row is one-per-module and answers "can I rely on this and how do I get it"; feature records are many-per-module, so an install coordinate repeated across them drifts. `artifact` keeps meaning the Maven coordinate, which every client has. The self-test gains thirteen cases covering each new failure path plus the passing direction (fragments from independent files all read, deferrals named). All were verified red against the pre-change checker, except the feature-path case, which was already caught; the old checker also stays green on the new tree, so branches carrying the old gate are unaffected by the new data. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…shared-file leaks (confluentinc#154) A six-persona document review of the language-proxy plan, aimed at the defect class a forty-edit corrections pass introduces: corrections that are themselves wrong, new text nobody has reviewed, and promises the edits quietly broke. Three corrections were factually wrong and are fixed against evidence rather than assumption. The branch is no longer behind master - origin/master was merged in at 3c66084 - so the rebase precondition would have sent an implementer rebasing across a merge commit to obtain a fix (71a306c) the branch already carries; the constraint is now a currency check instead of a point-in-time count. KTD11 declined Unix domain sockets because grpc-netty-shaded "bundles no epoll transport"; the 1.73.0 jar bundles it, native libraries included, so the decision now stands on its real grounds - macOS parity and the v1 auth scope - before the post-v6 security work rules out the cheapest mitigation on a false premise. And U34's currency-asymmetry line flipped the boundary the landed STRATEGY.md (5ac6b6d) deliberately states; the unit is marked landed so nobody "corrects" the shipped claim backwards. The Definition of Done's flow-control gate is scoped to capacity and count instructions, because applied literally it failed behaviours the Product Contract mandates - attempt counts, the echoed epoch, manifest Drop orders - and a gate that can only pass with unstated exemptions is a gate applied loosely, which is how the credit ledger reached implementation-ready once already. The property itself is unchanged: no capacity or count instruction may derive from observed client behaviour. The shared-file seeding promise had leaked back in three places - U15 and U28 writing the branch document concurrently, seven wave agents editing the guide at sync time, and CI rows that could only flip from skipped by editing U23's solely-owned workflow. Each now has one named writer. R71 and R16 gained the formal unit ownership their prose already performed, R18's warning now states what the surface accepts since R48 reversed R35, lease/window precedence on connection loss is stated so the reconnect machinery is not dead code, and demo containers are barred from the host Docker socket before ten authors each decide otherwise. Four findings need product decisions and are recorded as open questions rather than decided: the hosted gallery's substrate and owner, the executor-count function and its fork-safety ordering, lease expiry on a stalled-but-alive client, and whether breadth-at-launch gates the flagship's release. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SaBMWMJGbbpHQybuZUPCwa
…nfluentinc#154) Create everything the multi-language client work shares, before any client exists to conflict over it: after this commit, a language wave only ever adds files inside its own module directory, its own data fragments, and its own inflight note - so ten concurrent agents merge cleanly by construction rather than by discipline. The seed, in one commit deliberately - a partially-landed seed creates the conflicts it exists to prevent while looking done: - Root pom gains two module lines (protocol module, clients aggregator), before parallel-consumer-examples which stays last. Fifteen new modules total: the protocol skeleton, the clients aggregator, the nested Java aggregator, the three Java client modules (api, direct, grpc), and nine non-JVM client wrappers (python, go, swift, rust, dotnet, typescript, kotlin, ruby, cpp), each with its language-native build skeleton beside its pom. - Foreign toolchains are opt-in: the nine wrappers drive their language's own build/test tooling through exec-maven-plugin, wired once in the clients aggregator inside a profile inactive unless -Dpc.foreignClients is passed. An ordinary bin/build.sh builds all skeletons with no foreign tool installed - verified on a machine missing cargo, swift, dotnet, bundle, cmake and ruby - and a failing foreign command propagates its non-zero exit as a Maven build failure (verified with a /bin/false stub, and organically: `go test ./...` exits 1 on an empty module). One trap found and guarded: the root pom's exec-maven-plugin declaration (copyright check, inherited=false) still leaves a plugin DECLARATION in every child's effective model, so the profile's managed executions bound to every module in the clients tree, aggregators included. A pc.foreign.skip property (default true, overridden false only by the eight wrapper modules) plus deliberately-failing placeholder commands close that: if the guard ever regresses, the build goes red rather than a no-op reading as a passing foreign build. - The direct transport's transport-independence is enforced NOW, before any code exists to violate it: bannedDependencies on com.google.protobuf:*, io.grpc:* and the protocol module in parallel-consumer-proxy-client-java-direct. Verified red by temporarily adding protobuf-java: the failure names the banned artifact and the reason. - Both duplicate-code detector lists gain the three Java client src directories, and none of the nine non-JVM clients - the jobs take include lists with no exclusion parameter, and ten clients implementing one architecture in one intended shape is not accidental duplication. The comment beside each list records that the omission is deliberate. - Every new module carries its own module-maturity and testing-evidence fragment with a recorded deferral (reason + what lifts it), one file per module per corpus, so the wave that lifts one edits nothing anyone else owns. bin/check-docs-data.sh is green with all thirty present and goes red if one is deleted (verified). One inflight note per language under docs/inflight/clients/, the recorded exception to that directory's no-subdirectories rule. Not in this commit, by the plan's own split: the CI language matrix (clients.yml) and its skip-with-reason rows, the engine-side harness, and the frozen .proto - the three units that complete the seed before the fan-out releases anyone. The protocol module's content (proxy.proto, codegen) is likewise its own unit; only the skeleton and registration land here. Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SaBMWMJGbbpHQybuZUPCwa
…'s own fragment (confluentinc#154) Adding a client language to CI is now a matrix row, not a new job: clients.yml runs one job keyed per non-Java language (go, python, swift, rust, dotnet, typescript, kotlin, ruby, cpp - the Java clients ride maven.yml's existing lanes), with toolchain setup, Maven-driven build and test (-Dpc.foreignClients, safe here because the row is exactly where the toolchain is guaranteed present), and a language-native static-analysis scanner, all parameterised by the matrix entry. fail-fast is off, so one language's breakage cannot hide another's, and each row carries its own 30-minute timeout in a workflow of its own rather than adding eleven rows to a maven.yml already close to its limits. The skip is derived, never edited in. Every row's first step reads its module's own docs/data/module-maturity.d/<module>.yaml; while the fragment carries the deferred: block the docs-data gate defines, the row quotes the stated reason into the job summary and skips every later step - green-as-skipped, never red, so a red row during the fan-out always means a real failure in a started module. A language wave flips its row on by deleting that block from the one file it already owns, so nobody touches clients.yml after this commit and its sole-writer claim survives the fan-out. A missing fragment is a hard error rather than a green light, because absence means the docs-data corpus is broken, not that the module started. Negative control verified on this tree: all nine fragments currently carry deferred:, so every row skips. Scanners are the language's own: go vet, ruff, swift-format (bundled with the Swift 6 toolchain), clippy, dotnet format analyzers, eslint from the module's pinned devDependencies, detekt, rubocop, cppcheck. Caching follows maven.yml's discipline - explicit cache/restore plus rotating-key cache/save per language, restore-only for the shared Maven cache, and none of the setup actions' built-in immutable-key caches. GitHub-owned actions stay on the repo's one-tag-everywhere rule so bin/check-action-versions.sh and grouped Dependabot bumps keep covering them; the two non-GitHub actions (ruby/setup-ruby, swift-actions/setup-swift) are pinned to commit SHAs, and toolchain versions are exact pins in the matrix that fail loudly at setup if they stop resolving. Nothing uses pull_request_target. Verified: workflow YAML parses; bin/check-action-versions.sh green; bin/check-issue-refs.sh green; the gate logic simulated locally against all nine fragments (skip) and with a deferred: block removed (run). Upstream-Issue: confluentinc#154 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SaBMWMJGbbpHQybuZUPCwa
Three findings, all pre-existing on this branch and all invisible until the merge brought the checks that see them. A CITATION THAT NAMED A FILE THAT NEVER EXISTED. `check-file-refs.sh` could not see same-directory links until 34b5d74; with that fixed it flagged `branch-language-proxy.md`, cited from the archaeology note as where the Python client is being built. No ref in the repository has ever contained that path - `git log --diff-filter=D --all` finds no deletion, because there was nothing to delete. It was a plausible filename written from an adjacent fact. The handoff note this paragraph was rewritten from says it correctly and without a link: a Python client is being written on `feats/proxy-requirements` (#293, #242). Restored to that wording, which names a branch and two PRs a reader can check, rather than a note that has to exist for the sentence to be true. This is the third instance on this branch of the failure the branch was opened to document - a confident record built from adjacent evidence - and the second I committed myself. The other dangling reference was an ordinary rename left behind: `next-upstream-coverage-completeness.md` lost its `next-` prefix when this branch renamed the notes, and the link did not follow. TWO COUNTS OF THE SAME THING, DISAGREEING FIVE LINES APART. f7e58d7 widened "never write down what a command can answer" to name counts explicitly. The archaeology note asserted "upstream has 41 branches" in one bullet and "34 of upstream's 42 branches" in the next section. The 42 is this branch's own dated, verified measurement and matches `git ls-remote --heads upstream` today; the 41 was inherited, unverified, and had no date to warn the reader it was older. The standing total is now the shape plus the command that yields it. The dated measurement stays, because a measurement with its date and method on it is a finding, not a tracker - which is the distinction merge-checklist.md now asks a human to judge rather than a grep to enforce. THE LABELS AXIS APPLIED. Two notes here are concurrency-shaped by the label's own definition rather than by mentioning the word: core-178 is literally "same record, two threads at once", and core-139 is the missing thread-safety contract on the public API. The other triage notes are not, and are left alone - forcing a poor fit is what the closed set exists to prevent. Package rename: both sides were already on `bz.stub.*`, so the merge-order rule in AGENTS.md had nothing to do. `bin/check-all.sh` passes all fourteen tree gates. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…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
…ng model, PC the execution model Seventh bit of the follow-up Codex strategy conversation (2026-08-29/30), a short one. The engine-thesis note's gap 2 (polyglot Streams absent from STRATEGY.md) gains the follow-up's two-line positioning and why it is literal rather than marketing: #334's handles-not-IDL design assembles a real StreamsBuilder, so a Python or Rust application gets Kafka Streams itself - not a Streams-inspired API - with an execution model Streams does not currently provide. The follow-up's novelty search ("cannot find an existing implementation") is recorded as a lead, not a survey: one external model's single search does not license a public "first" claim without the prior-art sweep. The content series gains the recursion arc as a narrative piece: the machinery built to eliminate waiting inside one consumer becomes the runtime that gives every language Streams, the engine that gives Streams key-level concurrency, and the scheduler a Rust workload joins. Also fixes a bare #293 the previous commit left in the ecosystem-adapters diagram, caught by the issue-refs gate on this pass. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016PFA3CXkGc2m32cBJMAWTK
|
Hygiene extraction out of this PR: #378 - the write-time testing harness ( It carries one thing this branch does not have: the bridges import across directories, which no bridge had done before, and Same class as #338 and #339, and the same purpose: this PR should be reviewable as the remote execution model, not as repository machinery. |
…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>
The verdict-free work return has been extracted back out, as #295Unit U4 - letting work return to scheduling with no verdict, without consuming a retry (#242) - was folded into this PR when this PR was the only home for it. It now has its own home again: #295 is reopened, so its review history is intact, and its branch has been re-cut onto current It lands first, and the Wagon A stack depends on it. A core scheduling primitive should not first arrive inside a remote-runtime PR, and it was always dependency-free by design - so it can be reviewed by core owners without the proxy module in the way. Rungs of the stack that need the verdict-free path should carry What differs from the copy on this branchThis branch still carries the original commit, pre-package-rename. Reviewers comparing the two should know the extracted version is not a byte-for-byte replay -
Nothing to do on this branchNo history surgery here. This branch keeps its copy until #295 merges and Extraction A1 of the god-branch decomposition plan. |
… entry, and a gate that keeps it true The client modules bring go.mod, Cargo.toml, package.json, pyproject.toml, Gemfile, a .csproj and Package.swift into the tree. Until this commit `.github/dependabot.yml` declared `maven` and `github-actions` and nothing else, and a missing ecosystem is silent by construction: no error, no warning, no PR, indistinguishable from an ecosystem with nothing to update. So the entries land in the same change as the manifests, and bin/check-dependabot-coverage.sh (with its self-test) makes the invariant a build failure rather than a convention. It checks both directions, because each catches what the other cannot: a manifest no entry covers, and an entry whose directory does not exist or holds no manifest of the ecosystem it claims - the second being coverage that reads as present and is not. Verified red before green: run against this tree before the entries existed it exited 2 naming all seven modules. WHAT THEY WATCH TODAY IS ALMOST NOTHING, and declaring them anyway is the point. These modules compile a program that prints a line; most of the manifests declare no dependency at all. An empty ecosystem costs nothing to declare and means the coverage is already there on the day a client acquires a real dependency - which is the day it matters and the day nobody would think to come back. C++ IS COVERED BY NOTHING and is named rather than omitted, in the config and in the script's UNCOVERABLE list, which prints on green runs too. Its wording is adjusted from the source branch's: there the module has a CMakeLists.txt resolving through find_package against system packages; here it has no manifest at all, and both readings end at the same place - no version for any updater to bump. The entry is KEPT while the module has nothing, because dropping it would make the absence silent again, which is the failure the whole script is about. Also warms scala-maven-plugin's compiler and zinc bridge in the prepare-deps job. go-offline resolves declared dependencies and plugin JARs but not what a plugin fetches for itself at execution time, and the Scala compiler is fetched that way - measured on #293, where a cold fetch burned 4m01s on scala-compiler and then killed a lane on the compiler-bridge sources jar. Warmed by compiling the module rather than by naming coordinates, because the bridge version lives in the plugin's own zinc dependency and no pom property carries it, so a literal would rot silently the next time the plugin moves. `-am` keeps it to three projects - verified. Both brought over at the request of the #293 hygiene sweep (#378), which left them for whichever rung created the thing they check. The rest of that handover is judged out of scope and named in this PR's description. Extracted from #293.
|
Extraction A2 is cut: /pull/380 - the polyglot build scaffolding, the bottom rung of the Wagon A stack in What it takes out of this PR: the eleven language client module directories, reduced to toolchain detection, one compiled source file, a running program and one deterministic fixture line - plus the root pom's Three places it deliberately differs from this branch, so neither side reads as an oversight and nobody "fixes" one back:
What it leaves for this branch and the rungs above it: the protocol module and the frozen Nothing on this branch has been changed. The diff here shrinks when /pull/380 merges and |
…PR that measured them STRATEGY.md is a claims document nothing tests, so its contract is that work which falsifies a claim carries the update. This PR's measurements falsified two, and the updates were stranded on `perf/engine-concurrency` (#363) when the reviewable children were cut out of it. They belong here: this is the branch the numbers came from, and it is the only child that carries the notes they cite. WHAT THE DELTA CHANGES. The Share Groups position is rewritten around what survived. The throughput comparison is WITHDRAWN IN BOTH DIRECTIONS - not softened, withdrawn: the 2.5x against us did not reproduce at its own operating point (66,524 / 29,709 / 11,225 msg/s across three brokers in one day, a 5.9x range, while every PC arm held within 3%), the obvious explanation was tested with a negative control and refuted, and the 14.5% figure that was its counterweight goes with it because it is the same arm and the same variance. What replaces them is structural and was always the real argument: a share consumer cannot poll while records are unacknowledged, so an honest share processor is batch-synchronous, while PC holds records from many polls outstanding - which is what the offset encoding buys. Per-key ordering and retry semantics stay capability differences rather than benchmark results. The end-to-end latency success metric stops saying "not measured today" and names `pc.record.residence.time`, with the boundary it draws (client-side queueing and retries in, produce time and broker wait out) and the gap this harness cannot close: it drains a pre-produced backlog, so at 100% utilisation the number measures the backlog rather than the engine. Recording the gap beside the metric is the point - a success metric that quietly cannot be read at saturation is worse than one that admits it. The two CSVs are the async arms' arrival-capacity and tail-skew data, produced by the barrier-wiring fix this PR carries and stranded with everything else. RECONCILIATION STILL PENDING, and deliberately not attempted here. The `## Other runtimes` section this delta edits does not exist on master - it arrives from the proxy chain (`5ac6b6dc5` on `feats/proxy-requirements`, #293), which is extending the same section in parallel. Applying only this branch's delta means whichever lands second will have a real conflict in that section to resolve by hand. That is the correct outcome and not a mistake to tidy up: reconciling now would mean this PR silently editing the proxy chain's claims, or the proxy chain's, ours. Both sides own their own claims; the merge owns the reconciliation. Gates: check-docs-data and check-issue-refs clean. check-file-refs and check-inflight-tags stay at this branch's pre-existing 25 and 8 respectively - nothing added here appears in either failure list, verified by name. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
Extraction A3 is cut: #383 - the frozen v1 protocol contract, the second What it takes out of this PR: Its central claim is this PR's claim about the freeze, re-evidenced. The body here says the buf Four places it deliberately differs from this branch, so neither side reads as an oversight and
Two smaller things it left behind on purpose. The module's What it leaves for the rungs above: the foreign generated bindings, the proxy module and its gRPC Nothing on this branch has been changed. The diff here shrinks when #383 |
|
Extraction A5 is cut: #384 - the minimal proxy module and server shell, What it takes out of this PR: the The three transport classes went together because they are not separable in the code, not because The entry point is not on this branch, so it came from Four places it deliberately differs, so neither side reads as an oversight and nobody "fixes" one
One thing worth knowing if you re-derive the same cut here: the close is asserted against the
What it leaves for the rungs above, and for this PR: connect-time configuration and the options Nothing on this branch has been changed. The diff here shrinks when #384 |
|
Extraction A6 is out: #385 - the sidecar shell as a GraalVM native executable. Extraction A6 of Its source was It also re-parents the story #340 tells. Native distribution is proven Nothing about the native image is left in this PR's residue. |
|
Extraction A7 - the Java reference client - is out as #386. It takes this branch's The direct transport went with it, which the extraction plan was not sure of. It compiles against What was deliberately left here: the conformance framework - the transport-parameterised scenario Three things the extraction found, all recorded in that PR rather than only here:
Four smaller deltas are named one by one in #386's body, including the |
|
Extraction A8 is out: #387 - the conformance framework and its core The rung above #386, whose "deliberately NOT in this rung" list named What it takes out of this PR: the transport-parameterised suite ( What it leaves here, and this is the rung's one genuine judgement call: the Also left here: the ten foreign clients and everything that spawns them (the runner registry, the One reconciliation this leaves you, because a merge would resolve it silently. Three things the extraction found, all recorded in that PR rather than only here:
Nothing on this branch has been changed. The diff here shrinks when #387 |
|
Extraction A9 is out: /pull/390 - the ten foreign dispatch-only clients (Kotlin, Scala, Go, Python, TypeScript, Rust, Ruby, C#, C++, Swift), their build and test wiring at the pinned toolchains, and the foreign runner registry in the conformance module. It is the top rung of the Wagon A stack from Taken from this PR by path-scoped checkout rather than rewritten. The one thing that could not come across is the engine-dependent half: the sidecar on that stack hosts no engine and answers every session It also carries four defects found while wiring it up that apply here too: the Rust build script preferring the local Maven repository's |
…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.
|
This PR has been retargeted onto the top of the Wagon A stack, and its diff is now the engine residue. The plan's "retarget move" - The stack it now sits on#380 → #383 → #384 → They merge bottom-up; this one merges last. The body has been rewritten to the residue framing: the The two guards, and what they say nowBoth were left behind deliberately by the rungs below as self-retiring checks - they assert the
None of them was edited to go green. They pass because The reconciliation #387 asked for, and where it landedThat rung wrote into the class header, the module pom and a deferred inflight note that It decides the sidecar module's test tree, keeping the extracted rung's name: the lane needs One duty deliberately NOT discharged
The backupThe pre-merge tip of this branch is preserved on the remote as |
… engine back Forward-merges `feats/proxy-requirements` (#293) into this branch after that branch took the Wagon A stack (merge 8a7dc3e) and was retargeted onto `feats/proxy-dispatch-clients`. This is the descendant-chain half of 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, nothing was force-pushed, and this PR's base is unchanged. The pre-merge tip is preserved as `origin/backup/pre-forward-merge-328`. Most of the volume is content this branch was simply behind on: the ten foreign clients, the conformance framework, the protocol module's documents, current `master`, and the five records that became plain final classes because Error Prone crashes on a Jabel-desugared record. Only four files conflicted. ## The one that mattered: two `Main`s, and only one of them can be the sidecar `Main` was an add/add conflict, and the two sides are the same class written for opposite worlds. - This branch's is **unit U10**: it hosts `ConfigureHandler`, drains records held in a foreign process on shutdown, accepts `--socket` to listen on a Unix domain socket, and reserves exit 3 for a drain that timed out. `ReferenceDemo` - the whole point of this rung, and the 27x headline - spawns it through `SidecarProcess`, so a `Main` without an engine is a demo that cannot run. - The stack's is #384's engine-less shell: `NoEngineSessionService` behind a named `sessionServiceFactory()` seam, `EXIT_BIND_FAILED`, and a package-private `run` overload that lets a test name a port so the bind failure is reachable. The result is this branch's semantics on the stack's structure. `Main` keeps the seam, the bind exit code and the testable port; `sessionServiceFactory()` now returns the engine; the drain runs against whatever the session actually reached. Both usage texts, both sets of exit codes and both transports survive. ## Landing the engine there was not free, and the rung below said exactly what it would cost `docs/inflight/proxy-the-production-entry-point-still-hosts-no-engine.md` was written for this moment. It named the substitution as one call site, and then named the three things that had to happen with it - all three are done here, and the note is deleted rather than left describing a world that no longer exists: - **`NoEngineSessionService` is demoted to a test fixture, not deleted**, which is the outcome the note predicted. Eight cross-language `SidecarHandshakeTest`s (Kotlin, Scala, Go, Python, TypeScript, Rust, Ruby, C#) spawn a sidecar and assert the refusal arrives as `UNIMPLEMENTED` *specifically* - `PERMISSION_DENIED` and `RESOURCE_EXHAUSTED` are what the two interceptors raise before the service method runs, so the status code is the assertion rather than "it failed". An engine behind `Main` makes their subject vanish. - **`NoEngineMain` is the entry point they now spawn.** It is in the proxy module's test tree, so it ships in the test jar every client's sidecar classpath already includes for `TestModeMain` - nothing about how any of them spawns changed except the class name. It reimplements no lifecycle: it calls `Main`'s own seam with the no-engine supplier, so the binary those tests drive is the production lifecycle with one substitution, not a second copy that can drift. - **Nine call sites were re-pointed**: the eight foreign constants plus the Java client's own handshake test, which runs in-JVM and so calls `NoEngineMain.run` rather than spawning. The prose around each was corrected too - six of them described the thing they spawn as "the production `Main`", which is now the wrong binary rather than a wrong adjective. **`MainTest` is one class again.** Both sides grew one for the same class, in different packages, overlapping on arguments-refused, port-announced and two-sidecars. They are merged in `Main`'s own package keeping every assertion either side had, with the `UNIMPLEMENTED` test re-pointed at the no-engine service through the seam. `theProductionEntryPointHostsTheEngine` is new and is the substitution made assertable: without it the two suppliers could quietly swap back and every other test here would still pass, because the lifecycle is identical either way. ## The other three conflicts - **`.github/workflows/maven.yml`** - both sides added a matrix entry and neither displaced the other: this branch's `Demo (both entry points)` lane and master's `Chaos Pain Suite`. Union. - **`bin/check-copyright-headers.sh`** - master's refactor wins for the code (`check_derived`, `check_always_derived`, `fork_point_path`; this branch's copy still called `fp_path_of`, which no longer exists). The `Demo.java` recovery comment is this branch's, because master's version says the entry is ahead of its file and cites no ledger path *because none resolved on master* - both the file and `docs/inflight/branch-classic-comparison-demo.md` are here, so the path is cited. - **`ParentDeathWatchdog`** - a javadoc line only. This branch's wording wins because it is the one that is now true: the response to death is `DrainCoordinator`'s, and master's "once there is an engine it is the shutdown drain's" describes a future that has arrived. ## What was corrected because this merge falsified it `AGENTS.md`'s module table said the proxy module "hosts no engine, so a session is answered UNIMPLEMENTED", and the Python client's README called the engine-less sidecar "the production `Main`". Both are now wrong by one commit, so both were changed rather than left to rot. ## What was NOT corrected, and is a finding rather than an omission `docs/data/testing-evidence.d/parallel-consumer-proxy.yaml` says the module has no `src/test-integration` at all and that "there is no production entry point yet". Both were already false on this branch's own pre-merge tip - U10 landed `SidecarLifecycleIT` and `Main` without updating the record - so this is inherited rather than introduced, and rewriting an evidence record is the author's editorial call rather than a merge resolution. ## Two gates that the merge, not the code, turned red Both are 328's own files meeting analysers that arrived from `master` after this branch forked, so neither existed as a finding on either parent alone. - **SpotBugs now runs over this module** (`includeTests`, plus fb-contrib), and reported thirteen findings in the demo: three SLF4J format rules, class-envy and an over-concrete parameter on `ReferenceDemo`, a non-owned lock in its raw-gRPC arm, five exception-softening findings in `DemoBroker`, and two more in the test tree. **They arrive in two batches depending on whether `target/test-classes` exists when the gate runs** - the gate binds at `process-classes`, which is before `test-compile`, so a build that stops at `test-compile` sees the main-code findings only and a full `test` on a warm tree sees both. Do not read the first count as the whole set. Each was read at the site and dismissed in this module's own `spotbugs-exclude.xml`, class by class and pattern by pattern with the reason beside it - the shape the Java client aggregator's pom already prescribes and the api module already uses. None is a mute: a demo's output *is* pre-rendered text, a table renderer *does* read the rows it renders, and a `main()` whose only caller is a person reading a terminal is the one place `EXS_EXCEPTION_SOFTENING_NO_CONSTRAINTS` is the wrong rule. - **An unreadable exclude filter is a WARNING SpotBugs continues past.** The first version of that file quoted two command-line flags in an XML comment, and a double hyphen may not appear inside one - so the filter loaded as nothing, the gate kept failing with the same count, and the only evidence was one line reading `Unable to read filter` buried in the build log. Worth knowing before writing the next one. ## How it was verified JDK 17, macOS/arm64, toolchains through `mise` (`ruby 3.4.10`, `buf 1.72.0`). - **Full default reactor `./mvnw --fail-at-end test`.** Every module green, including the conformance suite over eleven of the thirteen bindings - so `java-grpc` runs against a live engine, and the Java, Kotlin and Scala handshake tests exercise the re-pointed no-engine entry point end to end. - **`bin/check-all.sh` on its real exit code: 1**, three gates failing and all three inherited. Established rather than assumed: the same gates on the incoming tip report the same `check-file-refs` and `check-branch-self-reference` findings, `check-inflight-tags` reports the same problems, and every *additional* finding names a file this merge did not touch - each verified with `git diff --quiet HEAD -- <file>`. `check-proto-lint` and `check-proto-breaking`, which report CANNOT RUN without `buf`, both pass under `mise`. - **`COPYRIGHT_CHECK_REQUIRE_FORK_POINT=1 bin/check-copyright-headers.sh`** - 0 violations against a real fork point. **Two things this box cannot answer, both established as inherited rather than assumed:** - **The `cpp` and `swift` conformance cells fail all ten of their scenarios with exit 126**, "found but cannot execute". Their runners are built in a container and extracted, so on macOS the artefact is a Linux binary - `file` reports `ELF 64-bit ... GNU/Linux` for both on a `Darwin arm64` host, which is the whole diagnosis and has nothing to do with anything in this merge. It is the condition #331 documents and fixes with `bin/run-conformance-in-container.sh`, one rung up. They are deselected here with `-Dpc.conformance.language`, named rather than silently skipped. - **Six of the eight foreign `SidecarHandshakeTest`s were not run** - Go, Python, TypeScript, Rust, Ruby and C# sit behind `-Dpc.foreignClients`, and C++ and Swift build inside containers on a box under a low-disk warning. Kotlin and Scala are JVM modules in the default reactor and did run, as did the Java client's in-JVM equivalent. What changed in all of them is one constant, and the process they now spawn is the same lifecycle with one supplier swapped.
…ArchUnit exemptions written twice Forward-merges `feats/go-vendored-pc` (#340) into this branch, the last link of the descendant-chain half of the plan's retarget move - `docs/plans/2026-08-31-001-process-god-branch-decomposition-plan.md`. Ordinary merge: no history rewritten, nothing force-pushed, this PR's base unchanged. The pre-merge tip is preserved as `origin/backup/pre-forward-merge-334`. The foreign Streams wrappers - the topology assembler, the handle protocol, the windowed aggregation and the Python side of all of it - are untouched. Seven files conflicted, and none of them is Streams code. **This worktree was 26 commits behind its own remote** and was fast-forwarded before the merge: the work continued on another machine, exactly as this branch's own handoff note (`docs/inflight/pr-334-handoff.md`) says it would. Those commits are the windowed-aggregation spike and its verdict, and they are what the merge is against. ## The same exemption, invented twice, in two files Both sides independently hit the same wall - the proxy client's api module is dependency-free by design, so it cannot wire `TestConventionRules` out of core's test-jar - and both wrote an exemption for it. The incoming pair wins in both files, and in one of them that is not a preference: - **`EveryModuleWiresUpArchUnitTest`** - the incoming constant is `EXEMPT_BECAUSE_THEY_CANNOT_DEPEND_ON_CORES_TEST_JAR`, and its companion guard `runsSomeOtherArchUnitWrapper` **had already auto-merged into this file** and reads that name. Keeping this branch's `DEPENDENCY_FREE_BY_DESIGN` would have left the guard pointing at a constant that no longer existed. The incoming javadoc is also the one that says why the row is an exemption rather than a mute - the module runs `ClientSurfaceArchTest` instead, and deleting that makes this test fail again - and its path comparison normalises the separator rather than relying on `Path` equality. - **The conformance module's `TestConventionsArchTest`** - the same wrapper written twice. The incoming javadoc answers the question a reader of that file actually has: why the wrapper cannot be inherited from the test-jar (surefire's `dependenciesToScan` would pull the whole test-jar into every module). Both sides also wanted the durable fix recorded, and it survives: `docs/refactoring.md` keeps **Relocate TestConventionRules out of core's test-jar**. ## Four registries, resolved as unions - and one that could not be `CONCEPTS.md`, `docs/refactoring.md`, `docs/inflight/bug-857-family.md` and `docs/inflight/pr-blockers-and-collisions.md` all conflicted for the same uninteresting reason: both sides appended different entries and they landed adjacently. Every one is a union; picking a side would have silently dropped somebody's entry. **`docs/quarantined-tests.md` is the exception, and it is a union that would have failed a gate.** This branch carried three entries, the incoming side one. The registry is enforced against the `@Quarantined` annotations actually in the tree, and after the merge exactly one survives - `ProducerManagerTest.producedRecordsCantBeInTransactionWithoutItsOffsetDirect`. The other two were not shelved, they were **fixed**: `de1620d71` un-quarantined `PCMetricsTest.metricsRegisterBinding` by freezing it on offsets, and `7028d9063` fixed the back-pressure test that had asserted an offset it had frozen. So this branch's copies are pre-fix, their diagnosis note is correctly gone with the fix, and an entry without an annotation is a hard failure of `bin/check-quarantine-registry.sh` - which now reports `Quarantine registry consistent (1 entry)`. ## Two dependency pins the merge broke, both of the class #293 named This module was written before the language-proxy stack merged down, and the stack brings the root pom's shared `grpc.version`. `RequireUpperBoundDeps` then refused two of this module's numbers, and the enforcer failed the module in 0.036s - before a single test ran, which is why the first reactor run showed 23 modules green and this one red. - **`protobuf-java` 3.25.5 to 3.25.8** - `grpc-protobuf:1.75.0` asks transitively for 3.25.8. The same bump, for the same reason, that the protocol module already took on the stack. The property stays module-local: this module owns its own `v1alpha1` schema, and the workflow that reads the protocol module's pin **by path** must not start reading this one. - **`error_prone_annotations` pinned at 2.47.0**, newly. gRPC 1.75.0 asks 2.30.0 through `grpc-stub`, `grpc-api` and `grpc-netty-shaded`; guava 33.6.0-jre - managed up from gRPC's own 33.3.1-android - asks 2.47.0. **The sibling modules are deliberately unequal here** - the proxy and the Java client aggregator sit at 2.48.0, the protocol module at 2.47.0, each at its own highest requester - so this is the number this module resolves to, not a copy of a neighbour's. Both are `dependencyManagement` entries local to this module, alongside the Jackson import that was already there for the same reason: nothing here is published, so no consumer inherits either pin. ## A sixth record had to become a class, for the reason the stack already recorded `ForeignCall` and `TopologyAssembler.Minted` were Java records. Error Prone 2.42.0 arrives with the stack, and it does not report on a Jabel-desugared record - it **crashes**, failing the whole compilation. #293 converted five value types in the proxy module for exactly this and recorded that 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 target. These are the sixth and seventh. Both keep their record-shaped accessors, so no call site changed, and `Minted` keeps its constructor validation - the kind-versus-node check that makes the erased casts in `resolve` and `sink` safe. `grep` now finds no `record` in the module. ## How it was verified JDK 17, macOS/arm64, toolchains through `mise`. - The first reactor run over this merge: **23 modules green**, this module the only failure, and it failed in 0.036s at the enforcer - before a test could run. With the four pins and the two record conversions in, the module builds and **its own 113 tests pass** (`TopologyAssemblerTest`'s 40 among them, which is what exercises `Minted`'s validation). - **`parallel-consumer-core` was not re-run until it went green, deliberately.** It flaked on `ParallelEoSStreamProcessorTest`'s sanity check on both full runs here. That is the known master-state flake, its core is byte-identical to the three rungs below, and today's tally across all four is appended to its ledger - `docs/inflight/test-processinkeyorder-sanity-check-races-the-first-poll.md`. Re-running for a green would have been the retry this repository removed on purpose. - `bin/check-all.sh` on its real exit code: 1, with the same three gates failing as at every other link in this chain - and this branch's own `docs/inflight/ci-standing-citation-and-tag-debt.md` already records two of them as red on standing debt rather than on anyone's change. - `bin/check-quarantine-registry.sh` green, which is the one that had to be, because this merge changed its subject. **Not run, and named rather than skipped quietly:** the `cpp` and `swift` conformance cells, whose runners are Linux ELF binaries on this host; they are deselected with `-Dpc.conformance.language` and belong to CI's Linux rows.
…hen it is looked up (#378) Manual mutation testing - sabotage the behaviour a new test guards and watch the test fail before committing it - was standard practice for a whole session only because it was written by hand into every dispatch prompt. The repo had the lessons as separate incident write-ups, each findable only by someone already suspecting the problem, and no binding rule anywhere: AGENTS.md mentions mutation once, about the lane's stale package regex, and docs/testing.md said nothing. So it goes in the layer that fires on its own. A short write-time slice, docs/testing-at-write-time.md, holds the rule and its three worked examples - the sabotage that never reached the system, the sabotage the engine made invisible by batching, and the test named for a property it could not detect - and a CLAUDE.md bridge in each module test tree imports that one file, so nothing is stated twice and it arrives when a test file is touched. docs/testing.md keeps the topic and points at the slice. The bridges are the one gitignore family scoped by pattern rather than enumerated, argued in the file: they are generated per test tree from one shape, so enumerating them would silently drop a new module's bridge - the exact everything-looked-wired-up-locally failure that section exists to prevent, arriving when someone is adding a module rather than thinking about harnesses. docs/agent-harness.md is updated on both sides of that: a bridge does not have to import an AGENTS.md, and the ".gitignore negation is the only place the question is asked" gap does not cover a pattern-negated family at all. Turning the bridges on exposed a defect in the citation gate. file-ref-gate.js reported a real, resolving path as dangling whenever a nested CLAUDE.md imported ACROSS directories: `@` is in TOKEN's character class, so `@../../../docs/x.md` reached the resolver with the import prefix attached, and normalising `<citing dir>/@../../..` let the `@..` segment absorb one of the `..` pops - the citation stopped one directory short. Every existing bridge was a same-directory `@AGENTS.md`, a one-segment token that is never a citation, so nothing had shown it. Fixed by stripping a LEADING `@` only; no tracked path contains an `@`, and the two-segment rule still runs afterwards. Self-tested with both arms - the import must yield the bare path AND resolve from the bridge that wrote it, and the adjacent `@AGENTS.md` must stay a non-citation - and proven able to fail by reverting the strip against the new case. VERIFICATION, and the trap in it. `git check-ignore` SKIPS a tracked path unless you pass --no-index, so the obvious check - stage the bridges, ask whether they are ignored - answers "not ignored" whatever .gitignore says. The control arm ran green against a .gitignore with the negation deleted before that was noticed. With --no-index it separates properly: negation removed, the bare CLAUDE.md rule claims the bridge; negation present, the `!**/src/test/CLAUDE.md` line does. Extracted from #293 (99a9ca1 on feats/proxy-requirements), which invented this while building the language proxy. Its proxy-module bridges stay behind with the modules they belong to; what is here is a bridge per module master already has, and the gate fix is the repository's, not the proxy's. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
Part of #242.
depends on #390
Description
This PR reviews the remote Parallel Consumer execution model. Everything around it - the polyglot
build machinery, the frozen wire contract, the sidecar's packaging and process boundary, the native
image, the Java reference client, the conformance framework and the ten language shells - is reviewed
in its own rung PR and is no longer part of this diff.
What is left is the thing that genuinely needs the whole in view: a Parallel Consumer engine whose
worker is in another process, and everything that follows from that being true.
ProxyProcessor- anExternalEnginewhose user function is a foreign worker over gRPC:dispatch waves and their coalescing window, the derived in-flight ceiling, the mailbox hand-back,
and the rule that every verdict-free return funnels through one method.
ConfigureHandlerandOptionsMapper: the session service thatbuilds an engine when a client's
Configurearrives, because nothing is read from a file, anenvironment variable or a flag.
reconciliation: what happens to records a worker was holding when its process went away.
cannot end a live one.
WorkManager#onAbandonedResultand the delivery-keyedabandon marker) - the one piece of this that is not in the proxy module, because a disconnect is
not a processing attempt and must not consume a retry.
Where the rest of it went
-Dpc.foreignClientsmachineryproxy.proto, its specification, the authoring guide, golden bytes, thebuf breakinggateThis PR is now based on
feats/proxy-dispatch-clientsrather thanmaster, so the diff below is theresidue rather than the whole continent. The rungs merge bottom-up; this one merges last. The plan is
docs/plans/2026-08-31-001-process-god-branch-decomposition-plan.md.What the merge that made this possible had to reconcile
The stack was cut out of this branch and then refined, and it carries current
masterwhile thisbranch did not - so merging it back was not a formality. The commit body of the merge carries the
full account; the three that a reviewer should know about:
ConformanceHarnessout of this branch'sProxyHarnesswith the engine lane removed, and recordedin three places that whoever landed the engine had to merge them rather than keep both - leaving
open which module it lands in. It lands here, in the sidecar module's test tree, keeping the
extracted rung's name: the engine lane needs
ConfigureHandlerandProxyProcessor, andTestModeMain- the process a spawned foreign runner actually starts - is in the same tree andwould otherwise close a dependency cycle.
ConformanceDriver.spawnAgainstand the ten foreign bindings are back, so the threeself-retiring guards test(conformance) astubbs#242: one definition of correct for every client, anchored by Parallel Consumer itself #387 and feat(clients) astubbs#242: ten foreign dispatch clients, and the runner registry that will drive them #390 left
behind now pass because the gap they held open is closed, not because anyone edited them.
masterafter this branch forked, and it crashes rather than reports on a Jabel-desugared record - so the
engine did not compile until they were converted.
What it does not do yet
Every client declares exactly one capability,
dispatch. Leases, heartbeats, reconnect,worker-death reporting, terminal outcomes and the shutdown drain exist in the engine and on the wire,
but no client negotiates them - so they are un-negotiated rather than half-built. Nothing is
published to any package registry. Batching, serialization/Schema Registry, demos and publishing are
recorded as follow-ups.
The production entry point still hosts no engine, and that is deliberate rather than an
oversight:
Main#sessionServiceFactoryreturnsNoEngineSessionService, because the drain that hasto land with the engine belongs to unit U10, which #384 named as still
having no PR. The engine is exercised end to end through
TestModeMaininstead.docs/inflight/proxy-the-production-entry-point-still-hosts-no-engine.mdcarries what U10 has to doand why eight cross-language handshake tests depend on the current shape.
Where the reasoning lives
docs/plans/2026-08-14-001-feat-language-proxy-plan.mdis the plan. Notes underdocs/inflight/carry what is open, what is parked and why. Commit bodies carry the diagnosis, the experiment and the
rejected alternative, because the release notes are generated from them.
roadmap-stage: the
language-proxy-sidecarentry'sstage_detailis updated in this PR to record that the artifact now arrives as a stack rather than as one PR, and that this rung carries the execution model and merges last. The stage itself is unchanged, and deliberately: the same artifact exists either way, so splitting how it is reviewed makes no more advanced artifact exist.Checklist
parallel-consumer-proxy-protocol/docs/), per-module READMEs, the module-maturity andtesting-evidence records, and the
docs/inflight/note for the deferred entry point.docs/features/- N/A - the sidecarpublishes to no registry and cannot yet be depended on; the release-documentation records it earns
are the maturity and evidence rows, which are in this PR and state in both directions what is and
is not evidenced.
and the three self-retiring guards that now pass honestly.
non-loopback opt-in that warns; credentials travel the protocol and appear in no log line at any
level (audited across all eleven clients); the unallowlisted Kafka property map remains a recorded
posture, mitigated by the loopback default.
ce-simplifyandce-code-reviewlocally - N/A - the code here was reviewed as part ofthis PR's earlier life and the rungs cut out of it; what is new is the reconciliation named above,
which is what a reviewer should spend attention on.
Known-red, and not a fault
Check PR Dependenciesis red until feat(clients) astubbs#242: ten foreign dispatch clients, and the runner registry that will drive them #390 and its parents merge.That is the dependency gate doing its job on the top of a stack.
claude-reviewis red until somebody asks for a review.