Skip to content

feat(composite-delivery): parallelize CompositeMessageDeliverer per-transport-group deliverBatch - #113

Merged
rawo merged 5 commits into
mainfrom
feat/kojak-81
Sep 24, 2026
Merged

rawo merged 5 commits into
mainfrom
feat/kojak-81

Conversation

@rawo

@rawo rawo commented Sep 23, 2026

Copy link
Copy Markdown
Contributor

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.

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:

  • 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.

…threads by default

Fallback to queued batch deliveries on single thread with configuration option.
#KOJAK-81
@rawo
rawo requested review from endrju19 and ramafasa and a lite review from Copilot September 23, 2026 13:15

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Interrupt cleanup can clear the caller’s interrupt flag, and documentation and configuration metadata corrections remain unresolved.

Get a fresh assessment by requesting another Copilot review.

Review effort: Lite
Findings: 3 Low severity

Open (3)
What changed in this PR

Parallelizes composite transport delivery using virtual threads, adds sequential dispatch configuration, and hardens delegate failure handling.

Changes:

  • Adds parallel and sequential transport-group dispatch.
  • Converts delegate errors and missing results into retriable failures.
  • Adds Spring Boot wiring, tests, documentation, and changelog updates.
File Summary
README.md Documents dispatch configuration; clarify that the final group runs inline.
okapi-spring-boot/​src/​test/​kotlin/​com/​softwaremill/​okapi/​springboot/​OutboxProcessorAutoConfigurationTest.kt Tests property binding and dispatch wiring.
okapi-spring-boot/​src/​main/​kotlin/​com/​softwaremill/​okapi/​springboot/​OutboxProcessorProperties.kt Adds dispatch configuration; metadata and executor-description updates remain needed.
okapi-spring-boot/​src/​main/​kotlin/​com/​softwaremill/​okapi/​springboot/​OutboxAutoConfiguration.kt Wires dispatch configuration into the composite deliverer.
okapi-spring-boot/​build.gradle.kts Updates test dependencies and repositories.
okapi-core/​src/​test/​kotlin/​com/​softwaremill/​okapi/​core/​CompositeMessageDelivererTest.kt Tests parallelism, interruption, and failure isolation.
okapi-core/​src/​main/​kotlin/​com/​softwaremill/​okapi/​core/​TransportDispatch.kt Defines dispatch modes.
okapi-core/​src/​main/​kotlin/​com/​softwaremill/​okapi/​core/​MessageDeliverer.kt Documents thread-safety and caller-thread requirements.
okapi-core/​src/​main/​kotlin/​com/​softwaremill/​okapi/​core/​CompositeMessageDeliverer.kt Implements dispatch and failure containment; interrupt cleanup and transaction-scope documentation need correction.
CHANGELOG.md Records behavioral changes; transaction-boundary wording needs correction.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread CHANGELOG.md Outdated
Comment thread README.md Outdated

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Unbounded executor cleanup can block indefinitely, and the public properties constructor needs binary compatibility preservation.

Get a fresh assessment by requesting another Copilot review.

Review effort: Lite
Findings: 2 High severity

Open (2)
Resolved since last review (3)
Previously missed (1)

In code that hasn't changed since last review

Low severity Add transport-dispatch to Spring configuration metadata

okapi-spring-boot/​src/​main/​kotlin/​com/​softwaremill/​okapi/​springboot/​OutboxProcessorProperties.kt:24

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.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Critical findings remain around the termination assertion and Kotlin binary compatibility.

Get a fresh assessment by requesting another Copilot review.

Review effort: Lite
Findings: 2 High severity

Open (2)
Resolved since last review (2)
Previously missed (2)

In code that hasn't changed since last review

Low severity Correct KDoc for parallel transport dispatch thread behavior

okapi-spring-boot/​src/​main/​kotlin/​com/​softwaremill/​okapi/​springboot/​OutboxProcessorProperties.kt:20

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.

Low severity Add transport-dispatch property to Spring configuration metadata

okapi-spring-boot/​src/​main/​kotlin/​com/​softwaremill/​okapi/​springboot/​OutboxProcessorProperties.kt:24

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.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Cancellation can return before interrupted delivery tasks stop, risking duplicate deliveries and delayed shutdown.

Get a fresh assessment by requesting another Copilot review.

Review effort: Lite
Findings: 2 High severity

Open (2)
Resolved since last review (1)

…tdownNow()

Explicitly explaining that the duplication is possible. It is part of the contract.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🔵 Needs a closer look

Address missing-deliverer task fan-out and add Spring configuration metadata for the new property.

Review effort: Lite
Findings: None

Resolved since last review (2)

@rawo rawo self-assigned this Sep 24, 2026
@rawo
rawo merged commit d28b882 into main Sep 24, 2026
9 checks passed
@rawo
rawo deleted the feat/kojak-81 branch September 24, 2026 10:34
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants