You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Parallelizes CompositeMessageDeliverer.deliverBatch so that transport groups — which are I/O-independent by construction — no longer block each other. Previously each group was delivered in turn via flatMap, so a batch spanning N transports cost T₁ + T₂ + ... + Tₙ; it now costs ~max(Tᵢ).
KOJAK-81 — Parallelize CompositeMessageDeliverer per-transport-group deliverBatch
Prerequisite KOJAK-74 (HTTP deliverBatch) and groundwork KOJAK-73 (Kafka fire-flush-await) are both on main, so there are finally 2+ overlappable transports for this to matter. Distinct from KOJAK-77, which parallelizes across scheduler workers rather than within one deliverBatch call.
Design decisions
One virtual thread per group, last group inline. N-1 groups are submitted to a per-call Executors.newVirtualThreadPerTaskExecutor(); the last runs on the calling thread rather than being handed to a thread the caller would only sit and wait for. One fewer thread per batch, and at least one transport keeps the caller's thread context.
Per-call executor, no injected Executor. Zero shared state, zero configuration burden, nothing to size or shut down. Parallelism is bounded by the number of registered transports (2–3), not by batch size or load, so there is no resource knob to get wrong.
Single-group fast path. A batch with one delivery type — the common case — starts no thread and allocates no executor, so homogeneous workloads are byte-for-byte the old path. Acceptance criterion "homogeneous batch has zero regression" is covered by a dedicated test, not just by inspection.
Future.get(), not join(). An interrupt on the caller (scheduler shutdown) is observed rather than swallowed: the flag is restored on the interrupted thread itself, which makes the remaining get() calls fail fast and makes the enclosing use { } escalate to shutdownNow() for still-running groups. Matches the reasoning already documented on HttpMessageDeliverer.awaitResult.
Order preservation is unchanged — final assembly still goes through the Map<OutboxEntry, DeliveryResult> lookup, and the existing order test passes untouched.
MessageDeliverer.deliverBatch is documented as MUST NOT throw, and CompositeMessageDeliverer is the one place that wraps third-party deliverers — but it violated that contract itself: the result-assembly step did error("missing result for entry …"), and any exception from a delegate propagated and aborted the whole batch.
Both are now contained. A deliverer that throws, or that returns no result for some of its entries, fails only its own entries as RetriableFailure (retriable rather than permanent: a broken transport must not dead-letter messages it never actually attempted), leaving the other transports' results intact.
This changes behavior for a transport that throws: previously the exception propagated → transaction rollback → the same rows were re-claimed indefinitely without ever incrementing retries. They now retry under the configured RetryPolicy and eventually reach FAILED. That is the documented contract and it stops one broken transport from blocking a mixed batch, but it is a real difference and is called out in CHANGELOG.
Escape hatch: TransportDispatch
The genuinely new constraint is not thread-safety — a single deliverer instance is already shared across workers whenever okapi.processor.concurrency > 1 — it is caller-thread affinity. Composite is the first place a third-party deliverer's own code runs on a thread the caller did not provide, so a deliverer reading through DataSourceUtils inside deliverBatch would silently get a fresh connection instead of joining the outbox transaction, and MDC correlation would drop. That requirement is now stated explicitly on the MessageDeliverer KDoc — i.e. this PR tightens a published contract in a minor release, post-1.0.
So the parallel path ships on by default (matching how HttpMessageDeliverer.deliverBatch and Kafka's fire-flush-await shipped, both of which are parallelism inside a deliverer rather than a resource knob like concurrency), with a one-line way back out:
CompositeMessageDeliverer(deliverers, TransportDispatch.SEQUENTIAL) — @jvmoverloads, so existing Kotlin and Java call sites of the one-arg constructor keep compiling and linking.
okapi.processor.transport-dispatch=sequential in okapi-spring-boot.
SEQUENTIAL restores pre-PR dispatch only; it shares deliverGroup and therefore keeps the contract fixes above, which would otherwise be a knob whose "off" position restores a contract violation.
Test coverage
New in CompositeMessageDelivererTest (core):
Transport groups overlap — proved by a 3-party CyclicBarrier rendezvous, not wall-clock timing, so it cannot flake on a loaded CI box; sequential dispatch simply can never trip the barrier.
Single-transport batch stays on the calling thread (the zero-regression criterion).
A throwing transport fails only its own entries, on both the parallel and the single-group fast path.
Entries a transport returns no result for fail retriably instead of throwing.
Caller interrupted mid-flight → retriable results, interrupt flag restored (and consumed, so it cannot leak into the next test).
SEQUENTIAL runs every group on the calling thread, and still isolates a throwing transport.
New in OutboxProcessorAutoConfigurationTest (spring-boot): two thread-recording deliverer beans and a real processBatch prove where the property actually lands — default yields 2 distinct threads (one of them the caller), sequential yields exactly 1 (the caller). No reflection into private fields, so it exercises the wiring rather than a field. Plus property binding, default value, and a bad-value startup-failure test.
Verified the concurrency tests are real guards, not tautologies: forcing dispatch back to sequential (awaiting each future at submit time) makes both the barrier test and the interrupt test fail, and they pass again on revert.
Please add the new okapi.processor.transport-dispatch entry to okapi-spring-boot/src/main/resources/META-INF/spring-configuration-metadata.json. That file is the module's published Spring configuration metadata, so binding works at runtime but IDE completion and generated configuration documentation will omit this newly supported public property.
This public property documentation says parallel dispatch creates one virtual thread per transport group, but CompositeMessageDeliverer intentionally runs the last group inline and creates only N-1 virtual threads (CompositeMessageDeliverer.kt:89-94). Please align this KDoc with the actual resource/thread behavior so users do not misestimate the dispatch cost.
Add transport-dispatch property to Spring configuration metadata
The new public okapi.processor.transport-dispatch property is not represented in the checked-in okapi-spring-boot/src/main/resources/META-INF/spring-configuration-metadata.json. Runtime binding works, but Spring Boot IDE metadata/autocomplete and generated configuration documentation will not expose this setting; add its type, default, description, and enum values to the metadata alongside this property.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Parallelizes CompositeMessageDeliverer.deliverBatch so that transport groups — which are I/O-independent by construction — no longer block each other. Previously each group was delivered in turn via flatMap, so a batch spanning N transports cost T₁ + T₂ + ... + Tₙ; it now costs ~
max(Tᵢ).KOJAK-81 — Parallelize CompositeMessageDeliverer per-transport-group deliverBatch
Prerequisite KOJAK-74 (HTTP deliverBatch) and groundwork KOJAK-73 (Kafka fire-flush-await) are both on main, so there are finally 2+ overlappable transports for this to matter. Distinct from KOJAK-77, which parallelizes across scheduler workers rather than within one deliverBatch call.
Design decisions
Contract hardening (behavioral change worth reviewing)
MessageDeliverer.deliverBatch is documented as MUST NOT throw, and CompositeMessageDeliverer is the one place that wraps third-party deliverers — but it violated that contract itself: the result-assembly step did error("missing result for entry …"), and any exception from a delegate propagated and aborted the whole batch.
Both are now contained. A deliverer that throws, or that returns no result for some of its entries, fails only its own entries as RetriableFailure (retriable rather than permanent: a broken transport must not dead-letter messages it never actually attempted), leaving the other transports' results intact.
This changes behavior for a transport that throws: previously the exception propagated → transaction rollback → the same rows were re-claimed indefinitely without ever incrementing retries. They now retry under the configured RetryPolicy and eventually reach FAILED. That is the documented contract and it stops one broken transport from blocking a mixed batch, but it is a real difference and is called out in CHANGELOG.
Escape hatch: TransportDispatch
The genuinely new constraint is not thread-safety — a single deliverer instance is already shared across workers whenever okapi.processor.concurrency > 1 — it is caller-thread affinity. Composite is the first place a third-party deliverer's own code runs on a thread the caller did not provide, so a deliverer reading through DataSourceUtils inside deliverBatch would silently get a fresh connection instead of joining the outbox transaction, and MDC correlation would drop. That requirement is now stated explicitly on the MessageDeliverer KDoc — i.e. this PR tightens a published contract in a minor release, post-1.0.
So the parallel path ships on by default (matching how HttpMessageDeliverer.deliverBatch and Kafka's fire-flush-await shipped, both of which are parallelism inside a deliverer rather than a resource knob like concurrency), with a one-line way back out:
SEQUENTIAL restores pre-PR dispatch only; it shares deliverGroup and therefore keeps the contract fixes above, which would otherwise be a knob whose "off" position restores a contract violation.
Test coverage
New in CompositeMessageDelivererTest (core):
New in OutboxProcessorAutoConfigurationTest (spring-boot): two thread-recording deliverer beans and a real processBatch prove where the property actually lands — default yields 2 distinct threads (one of them the caller), sequential yields exactly 1 (the caller). No reflection into private fields, so it exercises the wiring rather than a field. Plus property binding, default value, and a bad-value startup-failure test.
Verified the concurrency tests are real guards, not tautologies: forcing dispatch back to sequential (awaiting each future at submit time) makes both the barrier test and the interrupt test fail, and they pass again on revert.