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
13 changes: 13 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
26 changes: 25 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
@@ -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"
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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])
Expand Down Expand Up @@ -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
Expand All @@ -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()
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()}",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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]) →
Expand Down Expand Up @@ -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<String?, String> {
val info = KafkaDeliveryInfo.deserialize(entry.deliveryMetadata)
return ProducerRecord<String?, String>(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())
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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())
}
})
Loading