From 5d5431d01d929c437fd8e8c613fdbcad100ef5a1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rafa=C5=82=20Wokacz?= Date: Fri, 25 Sep 2026 15:46:48 +0200 Subject: [PATCH] feat(message-delivery): Add x-outbox-id header to kafka and http transport for handling deduplications on client side #KOJAK-78 --- CHANGELOG.md | 13 ++++++ README.md | 26 +++++++++++- .../softwaremill/okapi/core/OutboxHeaders.kt | 24 +++++++++++ .../okapi/http/HttpMessageDeliverer.kt | 11 +++++ .../okapi/http/HttpMessageDelivererTest.kt | 31 ++++++++++++++ .../KafkaTransportIntegrationTest.kt | 36 ++++++++++++++++ .../okapi/kafka/KafkaMessageDeliverer.kt | 12 ++++++ .../okapi/kafka/KafkaMessageDelivererTest.kt | 41 +++++++++++++++++++ 8 files changed, 193 insertions(+), 1 deletion(-) create mode 100644 okapi-core/src/main/kotlin/com/softwaremill/okapi/core/OutboxHeaders.kt diff --git a/CHANGELOG.md b/CHANGELOG.md index 9e5c602..1fdd6ed 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,19 @@ Until `1.0.0`, breaking changes may appear in any release and are flagged with * ## [Unreleased] +### Added + +- **`x-outbox-id` header on every delivery.** Both `KafkaMessageDeliverer` and `HttpMessageDeliverer` + now attach the entry's UUID under `OutboxHeaders.OUTBOX_ID` (`okapi-core`, so consumers can + reference the name without depending on a transport module). The value is stable across retries of + an entry, giving consumers something to deduplicate redeliveries on. okapi sets it **last**, so it + overrides a header of the same name supplied via `DeliveryInfo` — replaced outright over HTTP; + appended after yours in Kafka, whose headers are multi-valued, which is why consumers must read it + with `lastHeader(...)`. It does not deduplicate at the `publish()` level: two `publish()` calls for + the same business event are two entries with two ids. Previously the README stated that okapi did + not send the `OutboxId`; that is no longer the case, and the "Deduplicating on `x-outbox-id`" + section documents the recipe and both caveats. (KOJAK-78) + ### Changed - **`CompositeMessageDeliverer.deliverBatch`** now dispatches transport groups concurrently — one diff --git a/README.md b/README.md index 25b7be0..837c41c 100644 --- a/README.md +++ b/README.md @@ -91,11 +91,35 @@ Runnable, self-contained applications live in [okapi-examples](https://github.co ### Guarantees and limits -- **Duplicate delivery is possible.** A crash between a successful delivery and the status update means the message may be sent again after restart. If processing a message more than once would cause unwanted effects, make the consumer idempotent — for example, deduplicate on a business key in the payload or a header set in the `DeliveryInfo`. okapi sends your payload and configured headers, but not the `OutboxId` returned by `publish()`; that identifier stays on the publisher side for correlation and logging. +- **Duplicate delivery is possible.** A crash between a successful delivery and the status update means the message may be sent again after restart. If processing a message more than once would cause unwanted effects, make the consumer idempotent. okapi sends your payload, your configured headers, and an `x-outbox-id` header to deduplicate on — see [Deduplicating on `x-outbox-id`](#deduplicating-on-x-outbox-id). - **Best-effort ordering.** Rows are claimed by `created_at`, oldest first. However, parallel delivery and retries mean messages may reach consumers in a different order. Strict delivery ordering is not guaranteed. - **Failure classification is the transport's job.** Each deliverer decides what is retriable. HTTP: 5xx, 429, 408 and connection errors are retriable; other responses and TLS errors are permanent. Kafka: broker-side retriable exceptions are retried; authorization and configuration errors are not. - **Retry budget.** `okapi.processor.max-retries` (default 5) counts retries *after* the first attempt — six attempts in total before a row becomes `FAILED`. `FAILED` is terminal. Retriable messages become eligible again on the next processor poll; there is no per-message backoff. +### Deduplicating on `x-outbox-id` + +Every delivery carries an `x-outbox-id` header holding the outbox entry's UUID — the same value `publish()` returns. It is set by all transports (the name is the constant `OutboxHeaders.OUTBOX_ID` in `okapi-core`, so consumers can reference it without depending on a transport module), and the value is stable across retries: if okapi delivers the same entry twice, both copies carry the same id. A consumer that records ids it has already processed can therefore drop the repeat. + +```kotlin +// Kafka consumer +val outboxId = record.headers().lastHeader(OutboxHeaders.OUTBOX_ID)?.let { String(it.value()) } +if (outboxId != null && !seenIds.add(outboxId)) return // already processed, skip +``` + +```kotlin +// HTTP receiver (Spring MVC) +@PostMapping("/webhook") +fun receive(@RequestHeader("x-outbox-id") outboxId: String, @RequestBody payload: String) { + if (!seenIds.add(outboxId)) return // already processed, skip + // ... +} +``` + +Notice that: + +- **okapi sets the header last, so it overrides any `x-outbox-id` you set yourself** in `DeliveryInfo`. Over HTTP your value is replaced outright. Kafka headers are multi-valued, so your value is still present in the record, but okapi's is the one appended last — which is why consumers must read it with `lastHeader(...)` (or take the last of `headers(...)`) rather than the first match. Pick a different header name if you need to pass an identifier of your own. +- **It does not deduplicate at the `publish()` level.** The id identifies an *outbox entry*, not a business event. Calling `publish()` twice for the same event creates two entries with two different ids, and a consumer deduplicating on `x-outbox-id` will process both. Guarding against that is the publisher's job — deduplicate on a business key in the payload, or make the publish itself idempotent. + ## Configuration In a typical single-DataSource application, all properties are optional. Multi-DataSource setups may require explicit qualifiers, as described below. diff --git a/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/OutboxHeaders.kt b/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/OutboxHeaders.kt new file mode 100644 index 0000000..e804ede --- /dev/null +++ b/okapi-core/src/main/kotlin/com/softwaremill/okapi/core/OutboxHeaders.kt @@ -0,0 +1,24 @@ +package com.softwaremill.okapi.core + +/** + * Header names okapi attaches to every delivered message, independent of transport. + * + * Defined here rather than in a transport module so that each [MessageDeliverer] emits the same + * name and the same value format, and so consumers can reference the constant without depending + * on `okapi-kafka` or `okapi-http`. + */ +object OutboxHeaders { + /** + * Carries [OutboxEntry.outboxId] — the canonical UUID string, i.e. `outboxId.raw.toString()` — + * so consumers can deduplicate redeliveries of the same entry. + * + * Set by okapi on every delivery, and set **last**, so it always wins over a header of the same + * name supplied through `DeliveryInfo`: HTTP replaces the value outright, while Kafka appends + * (headers there are multi-valued), making okapi's the one `lastHeader` returns. + * + * The value is stable across retries of an entry, which is what makes it usable for + * deduplication. It does not identify a *business* event: two `OutboxPublisher.publish` calls + * for the same event produce two entries with two different ids. + */ + const val OUTBOX_ID: String = "x-outbox-id" +} diff --git a/okapi-http/src/main/kotlin/com/softwaremill/okapi/http/HttpMessageDeliverer.kt b/okapi-http/src/main/kotlin/com/softwaremill/okapi/http/HttpMessageDeliverer.kt index c1444bf..96ace8f 100644 --- a/okapi-http/src/main/kotlin/com/softwaremill/okapi/http/HttpMessageDeliverer.kt +++ b/okapi-http/src/main/kotlin/com/softwaremill/okapi/http/HttpMessageDeliverer.kt @@ -5,6 +5,7 @@ import com.softwaremill.okapi.core.DeliveryOutcome import com.softwaremill.okapi.core.DeliveryResult import com.softwaremill.okapi.core.MessageDeliverer import com.softwaremill.okapi.core.OutboxEntry +import com.softwaremill.okapi.core.OutboxHeaders import java.io.IOException import java.net.URI import java.net.http.HttpClient @@ -17,6 +18,9 @@ import javax.net.ssl.SSLException /** * [MessageDeliverer] that sends outbox entries as HTTP requests via JDK [HttpClient]. * + * Every request carries [OutboxHeaders.OUTBOX_ID] for consumer-side deduplication; see + * [buildRequest] for how it interacts with headers supplied through [HttpDeliveryInfo]. + * * Status code classification: * - 2xx → [DeliveryResult.Success] * - 5xx, 429, 408 → [DeliveryResult.RetriableFailure] (configurable via [retriableStatusCodes]) @@ -101,6 +105,12 @@ class HttpMessageDeliverer @JvmOverloads constructor( SendAttempt.ImmediateFailure(classifyThrowable(e)) } + /** + * [OutboxHeaders.OUTBOX_ID] is set after the caller's own headers, so it wins over an + * [HttpDeliveryInfo] header of the same name — `setHeader` replaces, so the caller's value is + * discarded rather than merely shadowed (Kafka, whose headers are multi-valued, keeps both and + * resolves by `lastHeader`). + */ private fun buildRequest(entry: OutboxEntry): HttpRequest { val info = HttpDeliveryInfo.deserialize(entry.deliveryMetadata) val url = urlResolver.resolve(info.serviceName) + info.endpointPath @@ -115,6 +125,7 @@ class HttpMessageDeliverer @JvmOverloads constructor( HttpRequest.BodyPublishers.ofString(entry.payload), ) .apply { info.headers.forEach { (k, v) -> setHeader(k, v) } } + .setHeader(OutboxHeaders.OUTBOX_ID, entry.outboxId.raw.toString()) .build() } diff --git a/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererTest.kt b/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererTest.kt index e6e2e8c..9d79ccd 100644 --- a/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererTest.kt +++ b/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererTest.kt @@ -9,6 +9,7 @@ import com.github.tomakehurst.wiremock.client.WireMock.urlEqualTo import com.github.tomakehurst.wiremock.core.WireMockConfiguration.wireMockConfig import com.softwaremill.okapi.core.DeliveryResult import com.softwaremill.okapi.core.OutboxEntry +import com.softwaremill.okapi.core.OutboxHeaders import com.softwaremill.okapi.core.OutboxMessage import io.kotest.core.spec.style.FunSpec import io.kotest.datatest.WithDataTestName @@ -83,6 +84,36 @@ class HttpMessageDelivererTest : FunSpec({ wiremock.verify(postRequestedFor(urlEqualTo("/test")).withHeader("Content-Type", equalTo("text/plain"))) } + test("sends x-outbox-id carrying the entry's UUID") { + wiremock.stubFor(post(urlEqualTo("/test")).willReturn(aResponse().withStatus(200))) + val e = entry() + + deliverer.deliver(e) + + wiremock.verify( + postRequestedFor(urlEqualTo("/test")) + .withHeader(OutboxHeaders.OUTBOX_ID, equalTo(e.outboxId.raw.toString())), + ) + } + + test("okapi's x-outbox-id replaces a caller-supplied header of the same name") { + wiremock.stubFor(post(urlEqualTo("/test")).willReturn(aResponse().withStatus(200))) + val info = httpDeliveryInfo { + serviceName = "svc" + endpointPath = "/test" + header(OutboxHeaders.OUTBOX_ID, "caller-supplied") + } + val e = OutboxEntry.createPending(OutboxMessage("test", """{"k":"v"}"""), info, Instant.now()) + + deliverer.deliver(e) + + // setHeader replaces, so unlike Kafka the caller's value is gone, not merely shadowed. + wiremock.verify( + postRequestedFor(urlEqualTo("/test")) + .withHeader(OutboxHeaders.OUTBOX_ID, equalTo(e.outboxId.raw.toString())), + ) + } + test("connection error -> RetriableFailure") { wiremock.stubFor( post(urlEqualTo("/test")) diff --git a/okapi-integration-tests/src/test/kotlin/com/softwaremill/okapi/test/transport/KafkaTransportIntegrationTest.kt b/okapi-integration-tests/src/test/kotlin/com/softwaremill/okapi/test/transport/KafkaTransportIntegrationTest.kt index 7fd2dc5..61eee7c 100644 --- a/okapi-integration-tests/src/test/kotlin/com/softwaremill/okapi/test/transport/KafkaTransportIntegrationTest.kt +++ b/okapi-integration-tests/src/test/kotlin/com/softwaremill/okapi/test/transport/KafkaTransportIntegrationTest.kt @@ -2,6 +2,7 @@ package com.softwaremill.okapi.test.transport import com.softwaremill.okapi.core.DeliveryResult import com.softwaremill.okapi.core.OutboxEntry +import com.softwaremill.okapi.core.OutboxHeaders import com.softwaremill.okapi.core.OutboxMessage import com.softwaremill.okapi.kafka.KafkaDeliveryInfo import com.softwaremill.okapi.kafka.KafkaMessageDeliverer @@ -83,6 +84,41 @@ class KafkaTransportIntegrationTest : FunSpec({ headerMap["source"] shouldBe "okapi" } + test("consumer reads x-outbox-id from the broker and it matches the entry UUID") { + val entry = entryWithInfo( + topic = "outbox-id-topic-${UUID.randomUUID()}", + headers = mapOf("traceId" to "trace-abc"), + ) + deliverer.deliver(entry) shouldBe DeliveryResult.Success + + val consumer = kafka.createConsumer(groupId = "test-outbox-id-${UUID.randomUUID()}") + consumer.subscribe(listOf(KafkaDeliveryInfo.deserialize(entry.deliveryMetadata).topic)) + val records = consumer.poll(Duration.ofSeconds(10)) + consumer.close() + + records.count() shouldBe 1 + val record = records.first() + // lastHeader is the documented way for consumers to read it. + String(record.headers().lastHeader(OutboxHeaders.OUTBOX_ID).value()) shouldBe entry.outboxId.raw.toString() + // Caller headers still survive alongside it. + record.headers().associate { it.key() to String(it.value()) }["traceId"] shouldBe "trace-abc" + } + + test("redelivering the same entry carries the same x-outbox-id, which is what makes dedup work") { + val entry = entryWithInfo(topic = "outbox-id-retry-topic-${UUID.randomUUID()}") + deliverer.deliver(entry) shouldBe DeliveryResult.Success + deliverer.deliver(entry) shouldBe DeliveryResult.Success + + val consumer = kafka.createConsumer(groupId = "test-outbox-id-retry-${UUID.randomUUID()}") + consumer.subscribe(listOf(KafkaDeliveryInfo.deserialize(entry.deliveryMetadata).topic)) + val records = consumer.poll(Duration.ofSeconds(10)) + consumer.close() + + records.count() shouldBe 2 + records.map { String(it.headers().lastHeader(OutboxHeaders.OUTBOX_ID).value()) }.toSet() shouldBe + setOf(entry.outboxId.raw.toString()) + } + test("deliver uses partition key") { val entry = entryWithInfo( topic = "key-topic-${UUID.randomUUID()}", diff --git a/okapi-kafka/src/main/kotlin/com/softwaremill/okapi/kafka/KafkaMessageDeliverer.kt b/okapi-kafka/src/main/kotlin/com/softwaremill/okapi/kafka/KafkaMessageDeliverer.kt index b1d8796..c0cf3f8 100644 --- a/okapi-kafka/src/main/kotlin/com/softwaremill/okapi/kafka/KafkaMessageDeliverer.kt +++ b/okapi-kafka/src/main/kotlin/com/softwaremill/okapi/kafka/KafkaMessageDeliverer.kt @@ -4,6 +4,7 @@ import com.softwaremill.okapi.core.DeliveryOutcome import com.softwaremill.okapi.core.DeliveryResult import com.softwaremill.okapi.core.MessageDeliverer import com.softwaremill.okapi.core.OutboxEntry +import com.softwaremill.okapi.core.OutboxHeaders import org.apache.kafka.clients.producer.Producer import org.apache.kafka.clients.producer.ProducerRecord import org.apache.kafka.clients.producer.RecordMetadata @@ -18,6 +19,9 @@ import java.util.concurrent.TimeoutException /** * [MessageDeliverer] that publishes outbox entries to Kafka topics. * + * Every record carries [OutboxHeaders.OUTBOX_ID] for consumer-side deduplication; see [buildRecord] + * for how it interacts with headers supplied through [KafkaDeliveryInfo]. + * * Exception classification: * - Kafka [RetriableException] → [DeliveryResult.RetriableFailure] * - Thread interrupts ([InterruptException], [InterruptedException]) → @@ -93,10 +97,18 @@ class KafkaMessageDeliverer( } } + /** + * [OutboxHeaders.OUTBOX_ID] is added after the caller's own headers. Kafka headers are + * multi-valued — `add()` appends rather than replaces — so ordering is what makes + * `headers().lastHeader(OutboxHeaders.OUTBOX_ID)` resolve to the outbox id even when a + * [KafkaDeliveryInfo] header of the same name is present. Such a header is not dropped; it + * stays visible via `headers().headers(...)`, just no longer last. + */ private fun buildRecord(entry: OutboxEntry): ProducerRecord { val info = KafkaDeliveryInfo.deserialize(entry.deliveryMetadata) return ProducerRecord(info.topic, info.partitionKey, entry.payload).apply { info.headers.forEach { (k, v) -> headers().add(k, v.toByteArray()) } + headers().add(OutboxHeaders.OUTBOX_ID, entry.outboxId.raw.toString().toByteArray()) } } diff --git a/okapi-kafka/src/test/kotlin/com/softwaremill/okapi/kafka/KafkaMessageDelivererTest.kt b/okapi-kafka/src/test/kotlin/com/softwaremill/okapi/kafka/KafkaMessageDelivererTest.kt index 26a3d6a..a8eec23 100644 --- a/okapi-kafka/src/test/kotlin/com/softwaremill/okapi/kafka/KafkaMessageDelivererTest.kt +++ b/okapi-kafka/src/test/kotlin/com/softwaremill/okapi/kafka/KafkaMessageDelivererTest.kt @@ -2,6 +2,7 @@ package com.softwaremill.okapi.kafka import com.softwaremill.okapi.core.DeliveryResult import com.softwaremill.okapi.core.OutboxEntry +import com.softwaremill.okapi.core.OutboxHeaders import com.softwaremill.okapi.core.OutboxMessage import io.kotest.core.spec.style.FunSpec import io.kotest.matchers.shouldBe @@ -75,4 +76,44 @@ class KafkaMessageDelivererTest : FunSpec({ Thread.interrupted() } } + + test("deliver attaches x-outbox-id carrying the entry's UUID") { + val producer = MockProducer(true, null, StringSerializer(), StringSerializer()) + val deliverer = KafkaMessageDeliverer(producer) + val entry = entry() + + deliverer.deliver(entry) + + val record = producer.history().single() + String(record.headers().lastHeader(OutboxHeaders.OUTBOX_ID).value()) shouldBe entry.outboxId.raw.toString() + } + + test("deliverBatch attaches x-outbox-id to every record, each carrying its own entry's UUID") { + val producer = MockProducer(true, null, StringSerializer(), StringSerializer()) + val deliverer = KafkaMessageDeliverer(producer) + val entries = listOf(entry(), entry(), entry()) + + deliverer.deliverBatch(entries) + + val sentIds = producer.history().map { String(it.headers().lastHeader(OutboxHeaders.OUTBOX_ID).value()) } + sentIds shouldBe entries.map { it.outboxId.raw.toString() } + } + + test("okapi's x-outbox-id is last, so it wins over a caller-supplied header of the same name") { + val producer = MockProducer(true, null, StringSerializer(), StringSerializer()) + val deliverer = KafkaMessageDeliverer(producer) + val info = kafkaDeliveryInfo { + topic = "test-topic" + header(OutboxHeaders.OUTBOX_ID, "caller-supplied") + } + val entry = OutboxEntry.createPending(OutboxMessage("test", """{"k":"v"}"""), info, Instant.now()) + + deliverer.deliver(entry) + + val headers = producer.history().single().headers() + String(headers.lastHeader(OutboxHeaders.OUTBOX_ID).value()) shouldBe entry.outboxId.raw.toString() + // Kafka headers are multi-valued: the caller's value is shadowed, not dropped. + headers.headers(OutboxHeaders.OUTBOX_ID).map { String(it.value()) } shouldBe + listOf("caller-supplied", entry.outboxId.raw.toString()) + } })