Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 34 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,42 @@ Until `1.0.0`, breaking changes may appear in any release and are flagged with *

## [Unreleased]

### Changed

- **`CompositeMessageDeliverer.deliverBatch`** now dispatches transport groups concurrently — one
virtual thread per group, with the last group run on the calling thread — instead of blocking on
each group in turn. A batch spanning N transports costs ~`max(Tᵢ)` rather than `sum(Tᵢ)`; a batch
with a single delivery type takes a fast path that starts no thread and allocates no executor, so
homogeneous workloads are unchanged. `MessageDeliverer` implementations must therefore be
thread-safe and must not depend on caller thread-locals (MDC, `TransactionSynchronizationManager`,
security context). The transaction scope is unchanged — the scheduler still wraps the whole
`OutboxProcessor.processNext` cycle (claim, deliver, update) in one transaction, so it remains
open across delivery — but a transport group dispatched to a virtual thread no longer inherits
the caller's transaction-bound resources, so a deliverer that implicitly joined the outbox
transaction (e.g. via Spring's `DataSourceUtils`) now gets a fresh connection. Deliverers that
cannot satisfy that can restore the previous behaviour with the new
`TransportDispatch.SEQUENTIAL` constructor argument, exposed in `okapi-spring-boot` as
`okapi.processor.transport-dispatch=sequential`. For a heterogeneous batch, parallel dispatch
also shortens how long the transaction stays open, from `sum(Tᵢ)` to ~`max(Tᵢ)`. (KOJAK-81)
- **`CompositeMessageDeliverer` now upholds `deliverBatch`'s "must not throw" contract even when a
transport does not.** A deliverer that throws, or that returns no result for some of its entries,
previously propagated the exception (or an `IllegalStateException` from the result-assembly step)
and aborted the whole batch; it now fails only its own entries, as `RetriableFailure`, leaving the
other transports' results intact. Such entries are retried under the configured `RetryPolicy`
rather than rolled back and re-claimed indefinitely. (KOJAK-81)

### Changed (BREAKING)

- **`OutboxProcessorProperties` gained a `transportDispatch` constructor parameter** (see the
`CompositeMessageDeliverer` entry above). `@JvmOverloads` preserves the previous four-argument
JVM constructor, so Java callers are unaffected, but adding a property to a Kotlin `data class`
necessarily changes the generated `copy`/`copy$default` and default-argument constructor
signatures. Kotlin code compiled against an earlier release and not recompiled will fail with
`NoSuchMethodError` if it calls `copy(...)` on this class or constructs it using default
arguments — recompile consumers against this release. Binding from `application.yml` /
`application.properties` is unaffected, as is `CompositeMessageDeliverer`, whose new parameter
is covered by `@JvmOverloads` for previously compiled Kotlin and Java callers alike.
([#113](https://github.com/softwaremill/okapi/pull/113))
- **`ExposedConnectionProvider`** now requires a `database: Database` constructor argument and
reads the active transaction from `database.transactionManager.currentOrNull()` instead of the
global `TransactionManager.currentOrNull()`. Previously, in a multi-database Exposed app, the
Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,7 @@ In a typical single-DataSource application, all properties are optional. Multi-D
| `okapi.processor.batch-size` | `10` | Maximum entries claimed per worker per tick. |
| `okapi.processor.max-retries` | `5` | Retries after the initial attempt before an entry becomes `FAILED`. |
| `okapi.processor.concurrency` | `1` | Parallel workers per tick, each claiming its own batch. Tune based on database capacity and delivery latency; see [Performance](#performance). |
| `okapi.processor.transport-dispatch` | `parallel` | How a batch spanning several `MessageDeliverer` beans is dispatched: `parallel` (all groups but one go to a virtual thread each — N transports use N-1 extra threads — while the last runs inline on the calling thread, so the batch costs ~`max(Tᵢ)`) or `sequential` (one group after another, all on the calling thread). Use `sequential` only if a custom deliverer depends on caller thread context — MDC, security context, or a transaction-bound resource such as a connection obtained via `DataSourceUtils`, since a group on a virtual thread does not inherit the caller's transaction. Batches with a single delivery type always run inline, regardless of this setting. |

### Purger

Expand Down
Original file line number Diff line number Diff line change
@@ -1,13 +1,22 @@
package com.softwaremill.okapi.core

import java.util.concurrent.Callable
import java.util.concurrent.Executors
import java.util.concurrent.Future

/**
* Routes delivery to the correct [MessageDeliverer] based on [OutboxEntry.deliveryType].
*
* Fails permanently for any entry whose type has no registered deliverer,
* rather than throwing — so the outbox processor can move it to FAILED
* and continue with remaining entries.
*
* [dispatch] selects how a batch spanning several transports is executed; see [TransportDispatch].
*/
class CompositeMessageDeliverer(deliverers: List<MessageDeliverer>) : MessageDeliverer {
class CompositeMessageDeliverer @JvmOverloads constructor(
deliverers: List<MessageDeliverer>,
private val dispatch: TransportDispatch = TransportDispatch.PARALLEL,
) : MessageDeliverer {
override val type: String = "composite"

private val registry: Map<String, MessageDeliverer> = deliverers.associateBy { it.type }
Expand All @@ -23,26 +32,131 @@ class CompositeMessageDeliverer(deliverers: List<MessageDeliverer>) : MessageDel
* to the matching deliverer's [MessageDeliverer.deliverBatch]. Results are
* re-assembled in original input order.
*
* Entries whose type has no registered deliverer are mapped to
* [DeliveryResult.PermanentFailure] (consistent with [deliver]).
* Transport groups are I/O-independent, so under [TransportDispatch.PARALLEL] (the default)
* they are dispatched **concurrently**: one virtual thread per group, with the last group run
* on the calling thread rather than handed to a thread that the caller would only wait for. A
* heterogeneous batch therefore costs ~`max(Tᵢ)` instead of `T₁ + T₂ + ... + Tₙ`. A batch with
* a single delivery type — the common case — takes a fast path that starts no thread and
* allocates no executor, whatever the configured [TransportDispatch].
*
* Because a group may then run on a thread other than the caller's, [MessageDeliverer]
* implementations must be thread-safe and must not depend on caller thread-locals (SLF4J MDC,
* Spring's `TransactionSynchronizationManager`, security context); deliverers that cannot
* satisfy that opt out via [TransportDispatch.SEQUENTIAL].
*
* The outbox transaction is open across delivery, and stays that way: the scheduler wraps the
* whole [OutboxProcessor.processNext] cycle — claim, deliver, update — in one transaction, so
* that `FOR UPDATE SKIP LOCKED` holds its row locks until the results are written. What this
* dispatch changes is which thread the transport runs on, not that scope. The store round-trips
* still run on the calling thread inside the transaction, but a group handed to a virtual
* thread does **not** inherit the caller's transaction-bound resources: a deliverer that
* implicitly joins the outbox transaction (Spring's `DataSourceUtils`, Exposed's thread-bound
* transaction) gets a fresh connection instead — precisely what [TransportDispatch.SEQUENTIAL]
* is for. For a heterogeneous batch, parallel dispatch also shortens how long that transaction
* stays open, from `sum(Tᵢ)` to ~`max(Tᵢ)`.
*
* Upholds [MessageDeliverer.deliverBatch]'s "must not throw" contract even when a transport
* does not: a deliverer that throws, or that returns no result for some of its entries, fails
* only its own entries — as [DeliveryResult.RetriableFailure], so nothing is silently dropped —
* while the other transports' results stand. Entries whose type has no registered deliverer are
* mapped to [DeliveryResult.PermanentFailure], consistent with [deliver].
*/
override fun deliverBatch(entries: List<OutboxEntry>): List<DeliveryOutcome> {
if (entries.isEmpty()) return emptyList()

val resultByEntry: Map<OutboxEntry, DeliveryResult> = entries
.groupBy { it.deliveryType }
.flatMap { (type, group) ->
val deliverer = registry[type]
if (deliverer != null) {
deliverer.deliverBatch(group)
} else {
group.map { DeliveryOutcome(it, DeliveryResult.PermanentFailure("No deliverer registered for type '$type'")) }
}
}
.associate { it.entry to it.result }
val groups = entries.groupBy { it.deliveryType }.map { (type, group) -> TransportGroup(type, group) }
val outcomes = when {
groups.size == 1 -> deliverGroup(groups.single())
dispatch == TransportDispatch.SEQUENTIAL -> groups.flatMap { group -> deliverGroup(group) }
else -> deliverGroupsInParallel(groups)
}
val resultByEntry: Map<OutboxEntry, DeliveryResult> = outcomes.associate { it.entry to it.result }

return entries.map { entry ->
DeliveryOutcome(entry, resultByEntry[entry] ?: error("missing result for entry ${entry.outboxId}"))
val result = resultByEntry[entry]
?: DeliveryResult.RetriableFailure(
"Deliverer for type '${entry.deliveryType}' returned no result for entry ${entry.outboxId}",
)
DeliveryOutcome(entry, result)
}
}

/**
* Fires every group but the last on its own virtual thread, then runs the last one on the
* calling thread before awaiting the rest. Returned outcomes are unordered — [deliverBatch]
* re-assembles input order from its `Map<OutboxEntry, DeliveryResult>` lookup.
*/
private fun deliverGroupsInParallel(groups: List<TransportGroup>): List<DeliveryOutcome> {
val executor = Executors.newVirtualThreadPerTaskExecutor()
return try {
val inFlight = groups.dropLast(1).map { group -> group to executor.submit(Callable { deliverGroup(group) }) }
val onCallerThread = deliverGroup(groups.last())
inFlight.flatMap { (group, future) -> await(group, future) } + onCallerThread
} finally {
// shutdownNow() without awaiting termination, deliberately not `use`/close(): close()
// loops on awaitTermination(1, DAYS) until every task ends, so one transport that
// ignores interruption would pin the caller here — and with it the scheduler's
// shutdown — long after we have results for the whole batch.
//
// On the normal path this is a no-op: every group was awaited above, so nothing is
// running. It only bites on the interrupt path, and there it is a deliberate trade:
// interruption is the only lever available (an HTTP request already on the wire or a
// `producer.send()` in flight cannot be cancelled), so a deliverer that ignores it may
// complete its send after this method has reported that group as retriable — the entry
// is then delivered again on a later claim. That is the at-least-once contract okapi
// already documents, and the same window exists in HttpMessageDeliverer.deliverBatch,
// which likewise abandons in-flight `sendAsync` futures when the caller is interrupted.
// Waiting would narrow it — the real result could then be persisted instead of a
// retriable one — but only by reintroducing the unbounded block this `finally` exists
// to avoid. Duplicate delivery is the documented contract; an indefinite hang is not.
// The tasks hold no shared state, and virtual threads are daemon threads, so an
// abandoned one cannot hold up JVM exit.
executor.shutdownNow()
Comment thread
rawo marked this conversation as resolved.
}
}

/**
* Awaits via [Future.get] rather than [java.util.concurrent.CompletableFuture.join] so an
* interrupt on the caller (scheduler shutdown) is observed instead of ignored. The flag is
* restored on the interrupted thread itself, so the remaining `get()` calls fail fast and the
* `finally` in [deliverGroupsInParallel] interrupts the still-running groups — without waiting
* for them, so a transport that ignores interruption cannot hold up the caller.
*
* Reporting an interrupted group as [DeliveryResult.RetriableFailure] while its send may still
* be in flight is what makes okapi at-least-once rather than exactly-once: that entry can be
* delivered twice. See the trade-off spelled out in [deliverGroupsInParallel], and the
* "Duplicate delivery is possible" guidance in the README — consumers must be idempotent.
*/
private fun await(group: TransportGroup, future: Future<List<DeliveryOutcome>>): List<DeliveryOutcome> = try {
future.get()
} catch (e: InterruptedException) {
Thread.currentThread().interrupt()
group.failAll(DeliveryResult.RetriableFailure("Interrupted while delivering type '${group.type}': ${describe(e)}"))
} catch (e: Exception) {
// deliverGroup() already converts every Exception into outcomes, so this only sees an
// ExecutionException wrapping an Error, or a CancellationException from shutdownNow().
// Retriable rather than permanent: a broken transport must not dead-letter messages that
// were never actually attempted.
group.failAll(DeliveryResult.RetriableFailure("Delivery of type '${group.type}' failed: ${describe(e.cause ?: e)}"))
}

private fun deliverGroup(group: TransportGroup): List<DeliveryOutcome> {
val deliverer = registry[group.type]
?: return group.failAll(DeliveryResult.PermanentFailure("No deliverer registered for type '${group.type}'"))

return try {
deliverer.deliverBatch(group.entries)
} catch (e: InterruptedException) {
Thread.currentThread().interrupt()
group.failAll(DeliveryResult.RetriableFailure("Interrupted while delivering type '${group.type}': ${describe(e)}"))
} catch (e: Exception) {
group.failAll(DeliveryResult.RetriableFailure("Deliverer for type '${group.type}' threw: ${describe(e)}"))
}
}

private fun describe(e: Throwable): String = "${e.javaClass.simpleName}: ${e.message ?: "no message"}"

private data class TransportGroup(val type: String, val entries: List<OutboxEntry>) {
fun failAll(result: DeliveryResult): List<DeliveryOutcome> = entries.map { DeliveryOutcome(it, result) }
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,11 @@ package com.softwaremill.okapi.core
*
* [type] must match the [DeliveryInfo.type] of the entries this deliverer handles.
* [CompositeMessageDeliverer] uses this to route entries to the correct implementation.
*
* Implementations must be thread-safe and must not depend on caller thread-locals: a single
* instance is shared by every scheduler worker (see `OutboxSchedulerConfig.concurrency`), and
* [CompositeMessageDeliverer] may invoke [deliverBatch] on a thread other than the calling one
* when a batch spans several transports.
*/
interface MessageDeliverer {
val type: String
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
package com.softwaremill.okapi.core

/**
* How [CompositeMessageDeliverer] dispatches a batch that spans more than one transport.
*
* A batch with a single delivery type — the common case — is unaffected by this setting: it always
* runs inline on the calling thread, with no executor and no thread hop.
*/
enum class TransportDispatch {
/**
* Transport groups run concurrently, one virtual thread per group (the last group runs on the
* calling thread), so a heterogeneous batch costs ~`max(Tᵢ)` instead of `T₁ + T₂ + ... + Tₙ`.
*
* The default. Requires [MessageDeliverer] implementations to tolerate running on a thread
* other than the caller's — see [SEQUENTIAL] for when that does not hold.
*/
PARALLEL,

/**
* Transport groups run one after another, entirely on the calling thread — the behaviour from
* before [PARALLEL] became the default.
*
* An escape hatch for custom deliverers that depend on caller thread context: SLF4J MDC,
* Spring's `TransactionSynchronizationManager` (e.g. a deliverer reading through
* `DataSourceUtils` expecting to join the outbox transaction), or a thread-bound security
* context. Costs `sum(Tᵢ)` for a heterogeneous batch.
*/
SEQUENTIAL,
}
Loading
Loading