diff --git a/CHANGELOG.md b/CHANGELOG.md index edcfce7..9e5c602 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/README.md b/README.md index ff0583e..25b7be0 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/CompositeMessageDeliverer.kt b/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/CompositeMessageDeliverer.kt index 1ba4ef6..5084b98 100644 --- a/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/CompositeMessageDeliverer.kt +++ b/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/CompositeMessageDeliverer.kt @@ -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 { +class CompositeMessageDeliverer @JvmOverloads constructor( + deliverers: List, + private val dispatch: TransportDispatch = TransportDispatch.PARALLEL, +) : MessageDeliverer { override val type: String = "composite" private val registry: Map = deliverers.associateBy { it.type } @@ -23,26 +32,131 @@ class CompositeMessageDeliverer(deliverers: List) : 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): List { if (entries.isEmpty()) return emptyList() - val resultByEntry: Map = 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 = 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` lookup. + */ + private fun deliverGroupsInParallel(groups: List): List { + 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() + } + } + + /** + * 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 = 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 { + 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) { + fun failAll(result: DeliveryResult): List = entries.map { DeliveryOutcome(it, result) } + } } diff --git a/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/MessageDeliverer.kt b/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/MessageDeliverer.kt index 43ea43e..6950506 100644 --- a/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/MessageDeliverer.kt +++ b/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/MessageDeliverer.kt @@ -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 diff --git a/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/TransportDispatch.kt b/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/TransportDispatch.kt new file mode 100644 index 0000000..e654081 --- /dev/null +++ b/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/TransportDispatch.kt @@ -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, +} diff --git a/okapi-core/src/test/kotlin/com/softwaremill/okapi/core/CompositeMessageDelivererTest.kt b/okapi-core/src/test/kotlin/com/softwaremill/okapi/core/CompositeMessageDelivererTest.kt index 81eb719..43ac96d 100644 --- a/okapi-core/src/test/kotlin/com/softwaremill/okapi/core/CompositeMessageDelivererTest.kt +++ b/okapi-core/src/test/kotlin/com/softwaremill/okapi/core/CompositeMessageDelivererTest.kt @@ -1,10 +1,19 @@ package com.softwaremill.okapi.core +import io.kotest.assertions.withClue import io.kotest.core.spec.style.FunSpec import io.kotest.matchers.shouldBe import io.kotest.matchers.string.shouldContain import io.kotest.matchers.types.shouldBeInstanceOf +import java.time.Duration import java.time.Instant +import java.util.concurrent.ConcurrentLinkedQueue +import java.util.concurrent.CountDownLatch +import java.util.concurrent.CyclicBarrier +import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicReference + +private const val AWAIT_TIMEOUT_SECONDS = 10L private fun deliveryInfo(t: String) = object : DeliveryInfo { override val type = t @@ -19,6 +28,22 @@ private fun fixedDeliverer(t: String, result: DeliveryResult) = object : Message override fun deliver(entry: OutboxEntry): DeliveryResult = result } +private fun batchDeliverer(t: String, batch: (List) -> List) = object : MessageDeliverer { + override val type = t + override fun deliver(entry: OutboxEntry): DeliveryResult = error("deliver() should not be used by deliverBatch") + override fun deliverBatch(entries: List): List = batch(entries) +} + +/** + * Rendezvous deliverer: blocks until every other transport group in the batch has also entered + * its `deliverBatch`. Sequential dispatch can never get past the barrier, so an overlap failure + * surfaces as a timeout (mapped to a retriable failure by the composite) rather than a hang. + */ +private fun rendezvousDeliverer(t: String, barrier: CyclicBarrier) = batchDeliverer(t) { entries -> + barrier.await(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS) + entries.map { DeliveryOutcome(it, DeliveryResult.Success) } +} + class CompositeMessageDelivererTest : FunSpec({ test("deliverBatch groups entries by type, delegates to each transport, preserves input order") { val composite = CompositeMessageDeliverer( @@ -98,4 +123,204 @@ class CompositeMessageDelivererTest : FunSpec({ batchCallsKafka shouldBe 1 batchCallsHttp shouldBe 1 } + + test("deliverBatch dispatches transport groups concurrently, not one after another") { + // Three parties: each of the three groups must be inside deliverBatch at the same moment + // for the barrier to trip, so the whole batch costs ~max(Tᵢ) rather than sum(Tᵢ). + val barrier = CyclicBarrier(3) + val composite = CompositeMessageDeliverer( + listOf( + rendezvousDeliverer("kafka", barrier), + rendezvousDeliverer("http", barrier), + rendezvousDeliverer("sqs", barrier), + ), + ) + val entries = listOf( + entryOfType("kafka", 1), + entryOfType("http", 2), + entryOfType("sqs", 3), + entryOfType("kafka", 4), + ) + + val results = composite.deliverBatch(entries) + + results.map { it.entry } shouldBe entries + results.map { it.result } shouldBe List(4) { DeliveryResult.Success } + } + + test("deliverBatch with a single transport group stays on the calling thread") { + var deliveryThread: Thread? = null + val composite = CompositeMessageDeliverer( + listOf( + batchDeliverer("kafka") { entries -> + deliveryThread = Thread.currentThread() + entries.map { DeliveryOutcome(it, DeliveryResult.Success) } + }, + ), + ) + + val results = composite.deliverBatch(listOf(entryOfType("kafka", 1), entryOfType("kafka", 2))) + + results.map { it.result } shouldBe listOf(DeliveryResult.Success, DeliveryResult.Success) + deliveryThread shouldBe Thread.currentThread() + } + + test("SEQUENTIAL dispatch runs every transport group on the calling thread") { + val deliveryThreads = ConcurrentLinkedQueue() + val recordingDeliverer = { t: String -> + batchDeliverer(t) { entries -> + deliveryThreads += Thread.currentThread() + entries.map { DeliveryOutcome(it, DeliveryResult.Success) } + } + } + val composite = CompositeMessageDeliverer( + listOf(recordingDeliverer("kafka"), recordingDeliverer("http"), recordingDeliverer("sqs")), + TransportDispatch.SEQUENTIAL, + ) + val entries = listOf(entryOfType("kafka", 1), entryOfType("http", 2), entryOfType("sqs", 3), entryOfType("http", 4)) + + val results = composite.deliverBatch(entries) + + results.map { it.entry } shouldBe entries + results.map { it.result } shouldBe List(4) { DeliveryResult.Success } + deliveryThreads.size shouldBe 3 + deliveryThreads.toSet() shouldBe setOf(Thread.currentThread()) + } + + test("SEQUENTIAL dispatch still isolates a throwing transport from the rest of the batch") { + val composite = CompositeMessageDeliverer( + listOf( + fixedDeliverer("kafka", DeliveryResult.Success), + batchDeliverer("http") { throw IllegalStateException("connection pool closed") }, + ), + TransportDispatch.SEQUENTIAL, + ) + val entries = listOf(entryOfType("kafka", 1), entryOfType("http", 2)) + + val results = composite.deliverBatch(entries) + + results[0].result shouldBe DeliveryResult.Success + results[1].result.shouldBeInstanceOf().error shouldContain "connection pool closed" + } + + test("a transport that throws fails only its own entries, retriably") { + val composite = CompositeMessageDeliverer( + listOf( + fixedDeliverer("kafka", DeliveryResult.Success), + batchDeliverer("http") { throw IllegalStateException("connection pool closed") }, + ), + ) + val entries = listOf(entryOfType("kafka", 1), entryOfType("http", 2), entryOfType("http", 3)) + + val results = composite.deliverBatch(entries) + + results.map { it.entry } shouldBe entries + results[0].result shouldBe DeliveryResult.Success + results[1].result.shouldBeInstanceOf().error shouldContain "connection pool closed" + results[2].result.shouldBeInstanceOf().error shouldContain "connection pool closed" + } + + test("a single throwing transport is caught on the fast path too") { + val composite = CompositeMessageDeliverer( + listOf(batchDeliverer("kafka") { throw IllegalStateException("producer closed") }), + ) + + val results = composite.deliverBatch(listOf(entryOfType("kafka", 1))) + + results.size shouldBe 1 + results[0].result.shouldBeInstanceOf().error shouldContain "producer closed" + } + + test("entries a transport returns no result for fail retriably instead of throwing") { + val composite = CompositeMessageDeliverer( + listOf( + batchDeliverer("kafka") { entries -> entries.take(1).map { DeliveryOutcome(it, DeliveryResult.Success) } }, + ), + ) + val entries = listOf(entryOfType("kafka", 1), entryOfType("kafka", 2)) + + val results = composite.deliverBatch(entries) + + results.map { it.entry } shouldBe entries + results[0].result shouldBe DeliveryResult.Success + results[1].result.shouldBeInstanceOf().error shouldContain "no result" + } + + test("a transport that ignores interruption cannot hold up the batch once the caller is interrupted") { + val delivererStarted = CountDownLatch(1) + val release = CountDownLatch(1) + // Models a transport that swallows interruption — exactly the task ExecutorService.close() + // would wait on (awaitTermination(1, DAYS) in a loop) before letting deliverBatch return. + val uninterruptible = batchDeliverer("kafka") { entries -> + delivererStarted.countDown() + var released = false + while (!released) { + released = try { + release.await(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS) + } catch (_: InterruptedException) { + false + } + } + entries.map { DeliveryOutcome(it, DeliveryResult.Success) } + } + val composite = CompositeMessageDeliverer(listOf(uninterruptible, fixedDeliverer("http", DeliveryResult.Success))) + val results = AtomicReference>() + val caller = Thread.ofVirtual().unstarted { + results.set(composite.deliverBatch(listOf(entryOfType("kafka", 1), entryOfType("http", 2)))) + } + + try { + caller.start() + delivererStarted.await(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS) shouldBe true + caller.interrupt() + + caller.join(Duration.ofSeconds(AWAIT_TIMEOUT_SECONDS)) + withClue("deliverBatch must not wait for a transport group it has already given up on") { + caller.isAlive shouldBe false + } + // Still latched: deliverBatch returned while that transport was demonstrably still running. + release.count shouldBe 1L + results.get()[0].result.shouldBeInstanceOf() + results.get()[1].result shouldBe DeliveryResult.Success + } finally { + release.countDown() + caller.join() + } + } + + test("interrupting the caller while a transport group is in flight yields retriable results and restores the flag") { + val delivererStarted = CountDownLatch(1) + val release = CountDownLatch(1) + val blocking = batchDeliverer("kafka") { entries -> + delivererStarted.countDown() + try { + release.await(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS) + entries.map { DeliveryOutcome(it, DeliveryResult.Success) } + } catch (e: InterruptedException) { + Thread.currentThread().interrupt() + entries.map { DeliveryOutcome(it, DeliveryResult.RetriableFailure(e.javaClass.simpleName)) } + } + } + val composite = CompositeMessageDeliverer(listOf(blocking, fixedDeliverer("http", DeliveryResult.Success))) + val caller = Thread.currentThread() + val interrupter = Thread.ofVirtual().start { + delivererStarted.await(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS) + caller.interrupt() + } + + try { + // "kafka" is the first-seen type, so it runs on a virtual thread while "http" — the last + // group — runs inline on this thread; the interrupt therefore lands while awaiting kafka. + val results = composite.deliverBatch(listOf(entryOfType("kafka", 1), entryOfType("http", 2))) + + results[0].result.shouldBeInstanceOf().error shouldContain "Interrupted" + results[1].result shouldBe DeliveryResult.Success + // Consumes the flag the composite restored, so it cannot leak into the next test. + Thread.interrupted() shouldBe true + } finally { + release.countDown() + interrupter.join() + Thread.interrupted() + } + } }) diff --git a/okapi-spring-boot/build.gradle.kts b/okapi-spring-boot/build.gradle.kts index 5e705e5..7226f67 100644 --- a/okapi-spring-boot/build.gradle.kts +++ b/okapi-spring-boot/build.gradle.kts @@ -98,6 +98,7 @@ dependencies { // LiquibaseDisabledNotice breadcrumb + our PTM↔DS validation cannot-verify WARN) — slf4j-simple // does not provide an introspectable appender. testImplementation(libs.logbackClassic) + testImplementation(kotlin("test")) } // CI version override: ./gradlew :okapi-spring-boot:test -PspringBootVersion=4.0.4 -PspringVersion=7.0.6 @@ -116,3 +117,6 @@ if (springBootVersion != null || springVersion != null) { } } } +repositories { + mavenCentral() +} diff --git a/okapi-spring-boot/src/main/kotlin/com/softwaremill/okapi/springboot/OutboxAutoConfiguration.kt b/okapi-spring-boot/src/main/kotlin/com/softwaremill/okapi/springboot/OutboxAutoConfiguration.kt index 4c49f1a..8d0556b 100644 --- a/okapi-spring-boot/src/main/kotlin/com/softwaremill/okapi/springboot/OutboxAutoConfiguration.kt +++ b/okapi-spring-boot/src/main/kotlin/com/softwaremill/okapi/springboot/OutboxAutoConfiguration.kt @@ -171,7 +171,7 @@ class OutboxAutoConfiguration( clock: ObjectProvider, ): OutboxEntryProcessor { return OutboxEntryProcessor( - deliverer = if (deliverers.size == 1) deliverers.single() else CompositeMessageDeliverer(deliverers), + deliverer = if (deliverers.size == 1) deliverers.single() else CompositeMessageDeliverer(deliverers, props.transportDispatch), retryPolicy = retryPolicy.getIfAvailable { RetryPolicy(maxRetries = props.maxRetries) }, clock = clock.getIfAvailable { Clock.systemUTC() }, ) diff --git a/okapi-spring-boot/src/main/kotlin/com/softwaremill/okapi/springboot/OutboxProcessorProperties.kt b/okapi-spring-boot/src/main/kotlin/com/softwaremill/okapi/springboot/OutboxProcessorProperties.kt index f7a8fb3..a0eb16f 100644 --- a/okapi-spring-boot/src/main/kotlin/com/softwaremill/okapi/springboot/OutboxProcessorProperties.kt +++ b/okapi-spring-boot/src/main/kotlin/com/softwaremill/okapi/springboot/OutboxProcessorProperties.kt @@ -1,10 +1,11 @@ package com.softwaremill.okapi.springboot +import com.softwaremill.okapi.core.TransportDispatch import org.springframework.boot.context.properties.ConfigurationProperties import java.time.Duration @ConfigurationProperties(prefix = "okapi.processor") -data class OutboxProcessorProperties( +data class OutboxProcessorProperties @JvmOverloads constructor( val interval: Duration = Duration.ofSeconds(1), val batchSize: Int = 10, val maxRetries: Int = 5, @@ -13,6 +14,14 @@ data class OutboxProcessorProperties( * via `FOR UPDATE SKIP LOCKED`. See [com.softwaremill.okapi.core.OutboxSchedulerConfig]. */ val concurrency: Int = 1, + /** + * How a batch spanning several `MessageDeliverer` beans is dispatched: `parallel` (default, + * one virtual thread per transport group) or `sequential` (one group after another on the + * calling thread). Only applies when 2+ deliverer beans are present — a single deliverer is + * never wrapped in a `CompositeMessageDeliverer`. See + * [com.softwaremill.okapi.core.TransportDispatch]. + */ + val transportDispatch: TransportDispatch = TransportDispatch.PARALLEL, ) { init { require(!interval.isNegative && interval.toMillis() > 0) { "interval must be at least 1ms" } diff --git a/okapi-spring-boot/src/test/kotlin/com/softwaremill/okapi/springboot/OutboxProcessorAutoConfigurationTest.kt b/okapi-spring-boot/src/test/kotlin/com/softwaremill/okapi/springboot/OutboxProcessorAutoConfigurationTest.kt index fcd5239..25d9cc6 100644 --- a/okapi-spring-boot/src/test/kotlin/com/softwaremill/okapi/springboot/OutboxProcessorAutoConfigurationTest.kt +++ b/okapi-spring-boot/src/test/kotlin/com/softwaremill/okapi/springboot/OutboxProcessorAutoConfigurationTest.kt @@ -1,24 +1,34 @@ package com.softwaremill.okapi.springboot +import com.softwaremill.okapi.core.DeliveryInfo +import com.softwaremill.okapi.core.DeliveryResult import com.softwaremill.okapi.core.MessageDeliverer +import com.softwaremill.okapi.core.OutboxEntry import com.softwaremill.okapi.core.OutboxEntryProcessor +import com.softwaremill.okapi.core.OutboxMessage import com.softwaremill.okapi.core.OutboxStore import com.softwaremill.okapi.core.TransactionRunner +import com.softwaremill.okapi.core.TransportDispatch import com.softwaremill.okapi.micrometer.MicrometerOutboxListener import com.softwaremill.okapi.micrometer.MicrometerOutboxMetrics import com.softwaremill.okapi.micrometer.OutboxMetricsRefresher import io.kotest.assertions.withClue import io.kotest.core.spec.style.FunSpec +import io.kotest.matchers.collections.shouldContain import io.kotest.matchers.collections.shouldNotBeEmpty import io.kotest.matchers.nulls.shouldNotBeNull import io.kotest.matchers.shouldBe +import org.springframework.beans.factory.getBean import org.springframework.boot.autoconfigure.AutoConfigurations import org.springframework.boot.autoconfigure.AutoConfigureAfter import org.springframework.boot.test.context.runner.ApplicationContextRunner import org.springframework.jdbc.datasource.SimpleDriverDataSource +import java.time.Duration import java.time.Duration.ofMillis import java.time.Duration.ofMinutes import java.time.Duration.ofSeconds +import java.time.Instant +import java.util.concurrent.ConcurrentLinkedQueue import javax.sql.DataSource class OutboxProcessorAutoConfigurationTest : FunSpec({ @@ -51,29 +61,80 @@ class OutboxProcessorAutoConfigurationTest : FunSpec({ "okapi.processor.batch-size=20", "okapi.processor.max-retries=3", "okapi.processor.concurrency=4", + "okapi.processor.transport-dispatch=sequential", ) .run { ctx -> - val props = ctx.getBean(OutboxProcessorProperties::class.java) + val props = ctx.getBean() props.interval shouldBe ofMillis(500) props.batchSize shouldBe 20 props.maxRetries shouldBe 3 props.concurrency shouldBe 4 + props.transportDispatch shouldBe TransportDispatch.SEQUENTIAL } } test("default properties when nothing is configured") { contextRunner.run { ctx -> - val props = ctx.getBean(OutboxProcessorProperties::class.java) + val props = ctx.getBean() props.interval shouldBe ofSeconds(1) props.batchSize shouldBe 10 props.maxRetries shouldBe 5 props.concurrency shouldBe 1 + props.transportDispatch shouldBe TransportDispatch.PARALLEL } } + test("multi-transport batches are dispatched in parallel by default") { + val deliveryThreads = ConcurrentLinkedQueue() + dispatchContextRunner(deliveryThreads).run { ctx -> + ctx.getBean() + .processBatch(listOf(entryOfType("kafka-like"), entryOfType("http-like"))) + + withClue("one group should run on a virtual thread while the other runs inline: $deliveryThreads") { + deliveryThreads.toSet().size shouldBe 2 + } + deliveryThreads shouldContain Thread.currentThread() + } + } + + test("okapi.processor.transport-dispatch=sequential keeps every transport group on the calling thread") { + val deliveryThreads = ConcurrentLinkedQueue() + dispatchContextRunner(deliveryThreads) + .withPropertyValues("okapi.processor.transport-dispatch=sequential") + .run { ctx -> + ctx.getBean() + .processBatch(listOf(entryOfType("kafka-like"), entryOfType("http-like"))) + + deliveryThreads.size shouldBe 2 + deliveryThreads.toSet() shouldBe setOf(Thread.currentThread()) + } + } + + // transportDispatch was appended to the primary constructor; without @JvmOverloads that would + // delete the 4-arg JVM constructor Java callers had before, breaking them at compile and link + // time. @JvmOverloads regenerates it — and this asserts Spring's Kotlin-aware constructor + // binding still picks the primary constructor now that the class has several. + test("OutboxProcessorProperties keeps the pre-transportDispatch JVM constructor for Java callers") { + val signatures = OutboxProcessorProperties::class.java.constructors.map { it.parameterTypes.toList() } + + withClue("available constructors: $signatures") { + signatures shouldContain listOf(Duration::class.java, Int::class.java, Int::class.java, Int::class.java) + signatures shouldContain + listOf(Duration::class.java, Int::class.java, Int::class.java, Int::class.java, TransportDispatch::class.java) + } + } + + test("invalid transport-dispatch triggers startup failure") { + contextRunner + .withPropertyValues("okapi.processor.transport-dispatch=concurrent") + .run { ctx -> + ctx.startupFailure.shouldNotBeNull() + } + } + test("SmartLifecycle is running after context start, and stop() actually halts it") { contextRunner.run { ctx -> - val scheduler = ctx.getBean(OutboxProcessorScheduler::class.java) + val scheduler = ctx.getBean() scheduler.isRunning shouldBe true scheduler.stop() scheduler.isRunning shouldBe false @@ -82,7 +143,7 @@ class OutboxProcessorAutoConfigurationTest : FunSpec({ test("getPhase returns PROCESSOR_PHASE constant (orders before purger)") { contextRunner.run { ctx -> - val scheduler = ctx.getBean(OutboxProcessorScheduler::class.java) + val scheduler = ctx.getBean() scheduler.phase shouldBe OutboxProcessorScheduler.PROCESSOR_PHASE } } @@ -105,7 +166,7 @@ class OutboxProcessorAutoConfigurationTest : FunSpec({ test("stop(callback) invokes callback AND actually halts the scheduler") { contextRunner.run { ctx -> - val scheduler = ctx.getBean(OutboxProcessorScheduler::class.java) + val scheduler = ctx.getBean() var callbackInvoked = false scheduler.stop { callbackInvoked = true } callbackInvoked shouldBe true @@ -117,7 +178,7 @@ class OutboxProcessorAutoConfigurationTest : FunSpec({ contextRunner .withBean("secondDeliverer", MessageDeliverer::class.java, { stubDelivererWithType("second") }) .run { ctx -> - val processor = ctx.getBean(OutboxEntryProcessor::class.java) + val processor = ctx.getBean() processor.shouldNotBeNull() ctx.getBeansOfType(MessageDeliverer::class.java).size shouldBe 2 } @@ -142,7 +203,7 @@ class OutboxProcessorAutoConfigurationTest : FunSpec({ }) .withPropertyValues("okapi.metrics.refresh-interval=1m") .run { ctx -> - val props = ctx.getBean(OkapiMetricsProperties::class.java) + val props = ctx.getBean() props.refreshInterval shouldBe ofMinutes(1) } } @@ -153,7 +214,7 @@ class OutboxProcessorAutoConfigurationTest : FunSpec({ io.micrometer.core.instrument.simple.SimpleMeterRegistry() }) .run { ctx -> - val props = ctx.getBean(OkapiMetricsProperties::class.java) + val props = ctx.getBean() props.refreshInterval shouldBe ofSeconds(15) } } @@ -191,10 +252,10 @@ class OutboxProcessorAutoConfigurationTest : FunSpec({ .withBean(DataSource::class.java, { SimpleDriverDataSource() }) .withBean(TransactionRunner::class.java, { noOpTransactionRunner() }) .run { ctx -> - ctx.getBean(io.micrometer.core.instrument.MeterRegistry::class.java).shouldNotBeNull() - ctx.getBean(MicrometerOutboxListener::class.java).shouldNotBeNull() - ctx.getBean(MicrometerOutboxMetrics::class.java).shouldNotBeNull() - ctx.getBean(OutboxMetricsRefresher::class.java).shouldNotBeNull() + ctx.getBean().shouldNotBeNull() + ctx.getBean().shouldNotBeNull() + ctx.getBean().shouldNotBeNull() + ctx.getBean().shouldNotBeNull() } } @@ -227,6 +288,37 @@ class OutboxProcessorAutoConfigurationTest : FunSpec({ }) // Loads a Spring Boot auto-config class by trying version-specific FQCNs in order. +/** + * Context with two thread-recording deliverers, so a batch spanning both transports proves where + * `okapi.processor.transport-dispatch` actually lands: the recorded threads are the observable + * difference between parallel and sequential dispatch. + */ +private fun dispatchContextRunner(deliveryThreads: ConcurrentLinkedQueue) = ApplicationContextRunner() + .withConfiguration(AutoConfigurations.of(OutboxAutoConfiguration::class.java)) + .withBean(OutboxStore::class.java, { stubStore() }) + .withBean("kafkaLikeDeliverer", MessageDeliverer::class.java, { threadRecordingDeliverer("kafka-like", deliveryThreads) }) + .withBean("httpLikeDeliverer", MessageDeliverer::class.java, { threadRecordingDeliverer("http-like", deliveryThreads) }) + .withBean(DataSource::class.java, { SimpleDriverDataSource() }) + .withBean(TransactionRunner::class.java, { noOpTransactionRunner() }) + +private fun threadRecordingDeliverer(t: String, into: ConcurrentLinkedQueue) = object : MessageDeliverer { + override val type = t + + override fun deliver(entry: OutboxEntry): DeliveryResult { + into += Thread.currentThread() + return DeliveryResult.Success + } +} + +private fun entryOfType(t: String): OutboxEntry { + val deliveryInfo = object : DeliveryInfo { + override val type = t + + override fun serialize(): String = "{}" + } + return OutboxEntry.createPending(OutboxMessage("evt", "{}"), deliveryInfo, Instant.EPOCH) +} + // Lets a single test exercise both the 3.5.x (`...actuate.autoconfigure.metrics...`) and 4.0.x (`...micrometer.metrics.autoconfigure...`) layouts. private fun resolveSpringBootClass(vararg candidateFqcns: String): Class<*> { val classLoader = OkapiMicrometerAutoConfiguration::class.java.classLoader