diff --git a/benchmarks/http-deliverbatch.json b/benchmarks/http-deliverbatch.json new file mode 100644 index 0000000..1c31194 --- /dev/null +++ b/benchmarks/http-deliverbatch.json @@ -0,0 +1,598 @@ +[ + { + "jmhVersion" : "1.37", + "benchmark" : "com.softwaremill.okapi.benchmarks.HttpThroughputBenchmark.drainAll", + "mode" : "avgt", + "threads" : 1, + "forks" : 2, + "jvm" : "/usr/lib/jvm/25.0.2-open/bin/java", + "jvmArgs" : [ + "-Xms8g", + "-Xmx8g", + "-XX:+UseG1GC", + "-Dliquibase.duplicateFileMode=WARN" + ], + "jdkVersion" : "25.0.2", + "vmName" : "OpenJDK 64-Bit Server VM", + "vmVersion" : "25.0.2+10-69", + "warmupIterations" : 3, + "warmupTime" : "10 s", + "warmupBatchSize" : 1, + "measurementIterations" : 5, + "measurementTime" : "30 s", + "measurementBatchSize" : 1, + "params" : { + "batchSize" : "10", + "httpLatencyMs" : "0" + }, + "primaryMetric" : { + "score" : 1.3326614469266165, + "scoreError" : 0.06602008494480224, + "scoreConfidence" : [ + 1.2666413619818142, + 1.3986815318714187 + ], + "scorePercentiles" : { + "0.0" : 1.2568559525714287, + "50.0" : 1.3349724140421053, + "90.0" : 1.3823694920473684, + "95.0" : 1.3830754144210526, + "99.0" : 1.3830754144210526, + "99.9" : 1.3830754144210526, + "99.99" : 1.3830754144210526, + "99.999" : 1.3830754144210526, + "99.9999" : 1.3830754144210526, + "100.0" : 1.3830754144210526 + }, + "scoreUnit" : "ms/op", + "rawData" : [ + [ + 1.3046833416, + 1.2568559525714287, + 1.30207778535, + 1.3830754144210526, + 1.3760161906842105 + ], + [ + 1.3520587696842106, + 1.29142953545, + 1.3178860584, + 1.372032945105263, + 1.370498476 + ] + ] + }, + "secondaryMetrics" : { + } + }, + { + "jmhVersion" : "1.37", + "benchmark" : "com.softwaremill.okapi.benchmarks.HttpThroughputBenchmark.drainAll", + "mode" : "avgt", + "threads" : 1, + "forks" : 2, + "jvm" : "/usr/lib/jvm/25.0.2-open/bin/java", + "jvmArgs" : [ + "-Xms8g", + "-Xmx8g", + "-XX:+UseG1GC", + "-Dliquibase.duplicateFileMode=WARN" + ], + "jdkVersion" : "25.0.2", + "vmName" : "OpenJDK 64-Bit Server VM", + "vmVersion" : "25.0.2+10-69", + "warmupIterations" : 3, + "warmupTime" : "10 s", + "warmupBatchSize" : 1, + "measurementIterations" : 5, + "measurementTime" : "30 s", + "measurementBatchSize" : 1, + "params" : { + "batchSize" : "10", + "httpLatencyMs" : "20" + }, + "primaryMetric" : { + "score" : 5.040692756166666, + "scoreError" : 0.37267453246678134, + "scoreConfidence" : [ + 4.668018223699884, + 5.413367288633448 + ], + "scorePercentiles" : { + "0.0" : 4.743718756833333, + "50.0" : 4.985558770833333, + "90.0" : 5.5493497707833335, + "95.0" : 5.591932076333333, + "99.0" : 5.591932076333333, + "99.9" : 5.591932076333333, + "99.99" : 5.591932076333333, + "99.999" : 5.591932076333333, + "99.9999" : 5.591932076333333, + "100.0" : 5.591932076333333 + }, + "scoreUnit" : "ms/op", + "rawData" : [ + [ + 5.143818479, + 4.962072875, + 5.000442485833333, + 4.9706750558333335, + 4.922914437333334 + ], + [ + 5.166109020833333, + 4.744830791333333, + 4.743718756833333, + 5.160413583333334, + 5.591932076333333 + ] + ] + }, + "secondaryMetrics" : { + } + }, + { + "jmhVersion" : "1.37", + "benchmark" : "com.softwaremill.okapi.benchmarks.HttpThroughputBenchmark.drainAll", + "mode" : "avgt", + "threads" : 1, + "forks" : 2, + "jvm" : "/usr/lib/jvm/25.0.2-open/bin/java", + "jvmArgs" : [ + "-Xms8g", + "-Xmx8g", + "-XX:+UseG1GC", + "-Dliquibase.duplicateFileMode=WARN" + ], + "jdkVersion" : "25.0.2", + "vmName" : "OpenJDK 64-Bit Server VM", + "vmVersion" : "25.0.2+10-69", + "warmupIterations" : 3, + "warmupTime" : "10 s", + "warmupBatchSize" : 1, + "measurementIterations" : 5, + "measurementTime" : "30 s", + "measurementBatchSize" : 1, + "params" : { + "batchSize" : "10", + "httpLatencyMs" : "100" + }, + "primaryMetric" : { + "score" : 14.118910350816668, + "scoreError" : 2.2777656483925592, + "scoreConfidence" : [ + 11.841144702424108, + 16.396675999209226 + ], + "scorePercentiles" : { + "0.0" : 13.179229250333334, + "50.0" : 13.416770770833335, + "90.0" : 17.673442560750004, + "95.0" : 17.9172068545, + "99.0" : 17.9172068545, + "99.9" : 17.9172068545, + "99.99" : 17.9172068545, + "99.999" : 17.9172068545, + "99.9999" : 17.9172068545, + "100.0" : 17.9172068545 + }, + "scoreUnit" : "ms/op", + "rawData" : [ + [ + 14.111793861333334, + 13.892910791666667, + 13.468068041666667, + 13.3654735, + 13.266195222333334 + ], + [ + 17.9172068545, + 15.479563917, + 13.290906611, + 13.179229250333334, + 13.217755458333333 + ] + ] + }, + "secondaryMetrics" : { + } + }, + { + "jmhVersion" : "1.37", + "benchmark" : "com.softwaremill.okapi.benchmarks.HttpThroughputBenchmark.drainAll", + "mode" : "avgt", + "threads" : 1, + "forks" : 2, + "jvm" : "/usr/lib/jvm/25.0.2-open/bin/java", + "jvmArgs" : [ + "-Xms8g", + "-Xmx8g", + "-XX:+UseG1GC", + "-Dliquibase.duplicateFileMode=WARN" + ], + "jdkVersion" : "25.0.2", + "vmName" : "OpenJDK 64-Bit Server VM", + "vmVersion" : "25.0.2+10-69", + "warmupIterations" : 3, + "warmupTime" : "10 s", + "warmupBatchSize" : 1, + "measurementIterations" : 5, + "measurementTime" : "30 s", + "measurementBatchSize" : 1, + "params" : { + "batchSize" : "50", + "httpLatencyMs" : "0" + }, + "primaryMetric" : { + "score" : 0.40358162615988774, + "scoreError" : 0.03151884569476366, + "scoreConfidence" : [ + 0.3720627804651241, + 0.4351004718546514 + ], + "scorePercentiles" : { + "0.0" : 0.3743532155576923, + "50.0" : 0.40496053347836736, + "90.0" : 0.42950256860333025, + "95.0" : 0.43018804717391307, + "99.0" : 0.43018804717391307, + "99.9" : 0.43018804717391307, + "99.99" : 0.43018804717391307, + "99.999" : 0.43018804717391307, + "99.9999" : 0.43018804717391307, + "100.0" : 0.43018804717391307 + }, + "scoreUnit" : "ms/op", + "rawData" : [ + [ + 0.3743532155576923, + 0.38934304821568627, + 0.37919615871153844, + 0.3857983651960784, + 0.4233332614680851 + ], + [ + 0.4204744467234043, + 0.42320865159574467, + 0.4142155918367347, + 0.43018804717391307, + 0.39570547512 + ] + ] + }, + "secondaryMetrics" : { + } + }, + { + "jmhVersion" : "1.37", + "benchmark" : "com.softwaremill.okapi.benchmarks.HttpThroughputBenchmark.drainAll", + "mode" : "avgt", + "threads" : 1, + "forks" : 2, + "jvm" : "/usr/lib/jvm/25.0.2-open/bin/java", + "jvmArgs" : [ + "-Xms8g", + "-Xmx8g", + "-XX:+UseG1GC", + "-Dliquibase.duplicateFileMode=WARN" + ], + "jdkVersion" : "25.0.2", + "vmName" : "OpenJDK 64-Bit Server VM", + "vmVersion" : "25.0.2+10-69", + "warmupIterations" : 3, + "warmupTime" : "10 s", + "warmupBatchSize" : 1, + "measurementIterations" : 5, + "measurementTime" : "30 s", + "measurementBatchSize" : 1, + "params" : { + "batchSize" : "50", + "httpLatencyMs" : "20" + }, + "primaryMetric" : { + "score" : 2.4637369849165154, + "scoreError" : 0.22172371793583026, + "scoreConfidence" : [ + 2.242013266980685, + 2.6854607028523456 + ], + "scorePercentiles" : { + "0.0" : 2.3148564686666666, + "50.0" : 2.3679583734166667, + "90.0" : 2.679749229028182, + "95.0" : 2.6842740499, + "99.0" : 2.6842740499, + "99.9" : 2.6842740499, + "99.99" : 2.6842740499, + "99.999" : 2.6842740499, + "99.9999" : 2.6842740499, + "100.0" : 2.6842740499 + }, + "scoreUnit" : "ms/op", + "rawData" : [ + [ + 2.6842740499, + 2.639025841181818, + 2.6176993445454544, + 2.5823947764545454, + 2.3148564686666666 + ], + [ + 2.3707573195833334, + 2.3559355000833335, + 2.36515942725, + 2.3576010035, + 2.349666118 + ] + ] + }, + "secondaryMetrics" : { + } + }, + { + "jmhVersion" : "1.37", + "benchmark" : "com.softwaremill.okapi.benchmarks.HttpThroughputBenchmark.drainAll", + "mode" : "avgt", + "threads" : 1, + "forks" : 2, + "jvm" : "/usr/lib/jvm/25.0.2-open/bin/java", + "jvmArgs" : [ + "-Xms8g", + "-Xmx8g", + "-XX:+UseG1GC", + "-Dliquibase.duplicateFileMode=WARN" + ], + "jdkVersion" : "25.0.2", + "vmName" : "OpenJDK 64-Bit Server VM", + "vmVersion" : "25.0.2+10-69", + "warmupIterations" : 3, + "warmupTime" : "10 s", + "warmupBatchSize" : 1, + "measurementIterations" : 5, + "measurementTime" : "30 s", + "measurementBatchSize" : 1, + "params" : { + "batchSize" : "50", + "httpLatencyMs" : "100" + }, + "primaryMetric" : { + "score" : 7.4837112907000005, + "scoreError" : 0.12976034179141005, + "scoreConfidence" : [ + 7.35395094890859, + 7.613471632491411 + ], + "scorePercentiles" : { + "0.0" : 7.34852254175, + "50.0" : 7.494728619875, + "90.0" : 7.58549876165, + "95.0" : 7.5875377085, + "99.0" : 7.5875377085, + "99.9" : 7.5875377085, + "99.99" : 7.5875377085, + "99.999" : 7.5875377085, + "99.9999" : 7.5875377085, + "100.0" : 7.5875377085 + }, + "scoreUnit" : "ms/op", + "rawData" : [ + [ + 7.40089597925, + 7.56714824, + 7.5669819895, + 7.54069975, + 7.535344094 + ], + [ + 7.5875377085, + 7.440802302, + 7.39506715625, + 7.45411314575, + 7.34852254175 + ] + ] + }, + "secondaryMetrics" : { + } + }, + { + "jmhVersion" : "1.37", + "benchmark" : "com.softwaremill.okapi.benchmarks.HttpThroughputBenchmark.drainAll", + "mode" : "avgt", + "threads" : 1, + "forks" : 2, + "jvm" : "/usr/lib/jvm/25.0.2-open/bin/java", + "jvmArgs" : [ + "-Xms8g", + "-Xmx8g", + "-XX:+UseG1GC", + "-Dliquibase.duplicateFileMode=WARN" + ], + "jdkVersion" : "25.0.2", + "vmName" : "OpenJDK 64-Bit Server VM", + "vmVersion" : "25.0.2+10-69", + "warmupIterations" : 3, + "warmupTime" : "10 s", + "warmupBatchSize" : 1, + "measurementIterations" : 5, + "measurementTime" : "30 s", + "measurementBatchSize" : 1, + "params" : { + "batchSize" : "100", + "httpLatencyMs" : "0" + }, + "primaryMetric" : { + "score" : 0.3195939434455658, + "scoreError" : 0.03573571177856272, + "scoreConfidence" : [ + 0.2838582316670031, + 0.3553296552241285 + ], + "scorePercentiles" : { + "0.0" : 0.29674805806557375, + "50.0" : 0.3125038394628872, + "90.0" : 0.36809225234857945, + "95.0" : 0.3699338622653061, + "99.0" : 0.3699338622653061, + "99.9" : 0.3699338622653061, + "99.99" : 0.3699338622653061, + "99.999" : 0.3699338622653061, + "99.9999" : 0.3699338622653061, + "100.0" : 0.3699338622653061 + }, + "scoreUnit" : "ms/op", + "rawData" : [ + [ + 0.31907520177192983, + 0.31043724430508474, + 0.3699338622653061, + 0.3515177630980392, + 0.3074357938474576 + ], + [ + 0.29983399381967213, + 0.30261940623333333, + 0.29674805806557375, + 0.3145704346206897, + 0.3237676764285714 + ] + ] + }, + "secondaryMetrics" : { + } + }, + { + "jmhVersion" : "1.37", + "benchmark" : "com.softwaremill.okapi.benchmarks.HttpThroughputBenchmark.drainAll", + "mode" : "avgt", + "threads" : 1, + "forks" : 2, + "jvm" : "/usr/lib/jvm/25.0.2-open/bin/java", + "jvmArgs" : [ + "-Xms8g", + "-Xmx8g", + "-XX:+UseG1GC", + "-Dliquibase.duplicateFileMode=WARN" + ], + "jdkVersion" : "25.0.2", + "vmName" : "OpenJDK 64-Bit Server VM", + "vmVersion" : "25.0.2+10-69", + "warmupIterations" : 3, + "warmupTime" : "10 s", + "warmupBatchSize" : 1, + "measurementIterations" : 5, + "measurementTime" : "30 s", + "measurementBatchSize" : 1, + "params" : { + "batchSize" : "100", + "httpLatencyMs" : "20" + }, + "primaryMetric" : { + "score" : 2.174603333653846, + "scoreError" : 0.032096370050147184, + "scoreConfidence" : [ + 2.142506963603699, + 2.2066997037039933 + ], + "scorePercentiles" : { + "0.0" : 2.150771185846154, + "50.0" : 2.1670441426153846, + "90.0" : 2.2067764270153845, + "95.0" : 2.2074136475384614, + "99.0" : 2.2074136475384614, + "99.9" : 2.2074136475384614, + "99.99" : 2.2074136475384614, + "99.999" : 2.2074136475384614, + "99.9999" : 2.2074136475384614, + "100.0" : 2.2074136475384614 + }, + "scoreUnit" : "ms/op", + "rawData" : [ + [ + 2.1631431763076923, + 2.1862946826153844, + 2.150771185846154, + 2.157385826923077, + 2.152124092846154 + ], + [ + 2.1957775416153846, + 2.170945108923077, + 2.2074136475384614, + 2.2010414423076923, + 2.161136631615385 + ] + ] + }, + "secondaryMetrics" : { + } + }, + { + "jmhVersion" : "1.37", + "benchmark" : "com.softwaremill.okapi.benchmarks.HttpThroughputBenchmark.drainAll", + "mode" : "avgt", + "threads" : 1, + "forks" : 2, + "jvm" : "/usr/lib/jvm/25.0.2-open/bin/java", + "jvmArgs" : [ + "-Xms8g", + "-Xmx8g", + "-XX:+UseG1GC", + "-Dliquibase.duplicateFileMode=WARN" + ], + "jdkVersion" : "25.0.2", + "vmName" : "OpenJDK 64-Bit Server VM", + "vmVersion" : "25.0.2+10-69", + "warmupIterations" : 3, + "warmupTime" : "10 s", + "warmupBatchSize" : 1, + "measurementIterations" : 5, + "measurementTime" : "30 s", + "measurementBatchSize" : 1, + "params" : { + "batchSize" : "100", + "httpLatencyMs" : "100" + }, + "primaryMetric" : { + "score" : 7.0586413417000005, + "scoreError" : 0.04465353994315141, + "scoreConfidence" : [ + 7.013987801756849, + 7.103294881643152 + ], + "scorePercentiles" : { + "0.0" : 7.0226183418, + "50.0" : 7.0565465044, + "90.0" : 7.10314267146, + "95.0" : 7.1043565664, + "99.0" : 7.1043565664, + "99.9" : 7.1043565664, + "99.99" : 7.1043565664, + "99.999" : 7.1043565664, + "99.9999" : 7.1043565664, + "100.0" : 7.1043565664 + }, + "scoreUnit" : "ms/op", + "rawData" : [ + [ + 7.1043565664, + 7.092217617, + 7.0778627, + 7.0226183418, + 7.0272288414 + ], + [ + 7.0312766, + 7.0537570168, + 7.0343162082, + 7.0834435334, + 7.059335992 + ] + ] + }, + "secondaryMetrics" : { + } + } +] + + diff --git a/benchmarks/results-postopt-kojak-74.md b/benchmarks/results-postopt-kojak-74.md new file mode 100644 index 0000000..08f9438 --- /dev/null +++ b/benchmarks/results-postopt-kojak-74.md @@ -0,0 +1,116 @@ +# HTTP deliverBatch parallel sendAsync — Results (KOJAK-74) + +Measured on MacBook M3 Max, JDK 25 (25.0.2), Postgres 16 + WireMock (in-JVM) via Testcontainers, +full JMH config: `fork=2, warmup=3 × 10s, iter=5 × 30s` — n=10 samples per benchmark. + +## Headline numbers — HTTP throughput + +Baseline is sequential blocking `httpClient.send()` (from +[`results-kafka-deliverbatch.md`](results-kafka-deliverbatch.md#http-throughput-companion-benchmark)). +Optimized is parallel `httpClient.sendAsync()` fire-all + `thenApply`/`exceptionally` + `get`. + +| batchSize | Baseline (ms/op) | Optimized (ms/op) | **Improvement** | +|-----------|-------------------|---------------------|------------------| +| **latency 0 ms** | | | | +| 10 | 0.638 | 1.333 ± 0.066 | **0.48×** (slower) | +| 50 | 0.321 | 0.404 ± 0.032 | **0.79×** (slower) | +| 100 | 0.290 | 0.320 ± 0.036 | **0.91×** (slower) | +| **latency 20 ms** | | | | +| 10 | 26.429 | 5.041 ± 0.373 | **5.24×** | +| 50 | 24.892 | 2.464 ± 0.222 | **10.10×** | +| 100 | 26.545 | 2.175 ± 0.032 | **12.21×** | +| **latency 100 ms**| | | | +| 10 | 108.515 | 14.119 ± 2.278 | **7.69×** | +| 50 | 105.313 | 7.484 ± 0.130 | **14.07×** | +| 100 | 107.714 | 7.059 ± 0.045 | **15.26×** | + +Translated to msg/s (`@OperationsPerInvocation(1000)`): + +| batchSize | latency 0 ms | latency 20 ms | latency 100 ms | +|-----------|--------------|----------------|-----------------| +| 10 | ~750 msg/s | ~198 msg/s | ~71 msg/s | +| 50 | ~2,475 msg/s | ~406 msg/s | ~134 msg/s | +| 100 | ~3,125 msg/s | ~460 msg/s | ~142 msg/s | + +Baseline msg/s for comparison: ~1,567–3,448 (latency 0), ~38–40 (latency 20), ~9.2–9.5 (latency 100), +flat across `batchSize` — confirming baseline was fully sequential. + +Raw JSON: [`http-deliverbatch.json`](http-deliverbatch.json). + +## What changed + +`HttpMessageDeliverer.deliverBatch` now fires all requests before awaiting any of them: +1. **Fire** — call `httpClient.sendAsync()` for every entry (non-blocking; a synchronous failure + building one entry's request, e.g. corrupt `deliveryMetadata`, is isolated to that entry via a + `SendAttempt` sealed type and does not prevent the rest of the batch from firing — mirrors the + `SendOutcome` pattern from `KafkaMessageDeliverer`). +2. **Classify inline** — each future is chained with `.thenApply { classifyResponse(...) }.exceptionally { classifyThrowable(...) }` + so every future always completes successfully with a `DeliveryResult`, never exceptionally. +3. **Await** — `.get()` per entry, in input order; since step 2 guarantees no exceptional + completion, this only returns normally or throws `InterruptedException` if the caller is + interrupted while waiting (`get()`, unlike `join()`, observes interrupts — handled by + restoring the flag and classifying the entry as retriable). The JDK `HttpClient`'s own + per-request 30s timeout (`HttpRequest.timeout()`) bounds the wait — no separate await-timeout + needed (unlike Kafka's `flush()`, which has no equivalent per-record deadline). + +`deliver()` was refactored to share `buildRequest`/`classifyResponse`/`classifyThrowable` helpers +with `deliverBatch`, with no single-entry behavior change (existing `HttpMessageDelivererTest` +passes unmodified). `classifyThrowable` also gained an explicit `SSLException` → `PermanentFailure` +branch (checked before `IOException`, since `SSLException` is a subtype) — a TLS handshake/cert +failure is a configuration problem that won't fix itself on retry, unlike a transient connection +reset. + +## Reading the table + +- **Real gains at realistic latency, but short of the ticket's optimistic 1000+ msg/s target.** + At `batchSize=50, latencyMs=20` (the ticket's headline scenario) we measured **~406 msg/s** + (10.1×), not 1000+. The parallelization itself is proven working — see the concurrency section + below — but per-batch fixed costs (DB claim query, `batchSize` individual `UPDATE` statements, + transaction commit, JVM/connection-pool bookkeeping) that were previously hidden behind + sequential network wait now make up a larger share of wall time once the network wait is + parallelized away. This is the same effect the Kafka results doc flagged as "next bottleneck": + KOJAK-75 (batch `UPDATE` via JDBC `executeBatch`) directly attacks this, and should lift these + numbers further without any change to `HttpMessageDeliverer` itself. +- **`latencyMs=0` regressed slightly (0.79×–0.91×), `batchSize=10` more noticeably (0.48×).** + Expected: with no network wait to hide, `sendAsync`'s extra `CompletableFuture` allocation and + chaining (`thenApply`/`exceptionally`) is pure overhead versus a single blocking `send()` call. + At `batchSize=10` this fixed per-request overhead is a larger fraction of an already-tiny + 10-message batch. This confirms the AC's own framing — `latencyMs=0` is "already CPU-bound; + library + DB + WireMock overhead" — the optimization targets I/O-bound webhook delivery, not + the zero-latency floor. +- **Improvement grows with both `batchSize` and `latencyMs`** (5.2× → 15.3× from + `(10, 20ms)` to `(100, 100ms)`), consistent with the mechanism: more concurrent in-flight + requests and a larger per-request wait to overlap both increase the ratio of hidden-vs-exposed + latency. +- **Baseline is flat across `batchSize`** (as previously documented) — confirms the "before" + picture really was fully sequential, so the comparison is apples-to-apples. + +## Concurrency proof + +`HttpMessageDelivererBatchConcurrencyTest` fires a 10-entry batch against a WireMock stub with a +fixed 300 ms per-request delay and inspects WireMock's request journal (`allServeEvents`): all 10 +requests' `loggedDate` timestamps fall within a single 300 ms window of each other (spread `< +delayMs`), and the overall `deliverBatch` call completes in well under half of the 3000 ms a fully +sequential implementation would need. This is direct evidence — not just an aggregate timing +inference — that the requests were in flight concurrently. + +## Verification context + +- Unit tests: `HttpMessageDelivererBatchTest` covers empty input, all-success ordering, mixed + success/failure classification, all-fail batch, connection reset, poison-pill metadata isolation, + a custom injected `HttpClient`, and a TLS handshake failure classifying as `PermanentFailure`. +- `HttpMessageDelivererBatchConcurrencyTest` proves concurrent in-flight requests via WireMock's + request log (see above). +- Existing single-entry `HttpMessageDelivererTest` passes unmodified after the `deliver()` refactor. +- Integration tests (`HttpEndToEndTest`, `MysqlHttpEndToEndTest` in `okapi-integration-tests`) + continue to pass with real Postgres + WireMock, now exercising `deliverBatch` end-to-end via + `OutboxEntryProcessor`. +- ktlint clean project-wide (`ktlintCheck`); full `./gradlew build` green. + +## What's next + +1. **Batch UPDATE via JDBC `executeBatch`** (KOJAK-75) — now the more clearly load-bearing + bottleneck for HTTP too, for the same reason it already was for Kafka: with network wait + parallelized away, the N individual per-entry `UPDATE` statements are next in line. +2. **Multi-threaded scheduler with pluggable executor** (KOJAK-77) — multiplies both the Kafka + and HTTP `deliverBatch` gains across N concurrent workers. 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 8a93abb..c1444bf 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 @@ -1,6 +1,7 @@ package com.softwaremill.okapi.http import com.fasterxml.jackson.core.JsonProcessingException +import com.softwaremill.okapi.core.DeliveryOutcome import com.softwaremill.okapi.core.DeliveryResult import com.softwaremill.okapi.core.MessageDeliverer import com.softwaremill.okapi.core.OutboxEntry @@ -10,6 +11,8 @@ import java.net.http.HttpClient import java.net.http.HttpRequest import java.net.http.HttpResponse import java.time.Duration +import java.util.concurrent.CompletableFuture +import javax.net.ssl.SSLException /** * [MessageDeliverer] that sends outbox entries as HTTP requests via JDK [HttpClient]. @@ -20,9 +23,16 @@ import java.time.Duration * - other → [DeliveryResult.PermanentFailure] * * Exception classification: + * - [SSLException] (bad cert, protocol mismatch — won't fix itself) → [DeliveryResult.PermanentFailure] * - [IOException] (connection/timeout) → [DeliveryResult.RetriableFailure] - * - [InterruptedException] → [DeliveryResult.RetriableFailure] (interrupt flag restored) - * - other (corrupt metadata, unknown service, malformed URI) → [DeliveryResult.PermanentFailure] + * - [InterruptedException] → [DeliveryResult.RetriableFailure] (interrupt flag restored at the + * synchronous call site that observed it; classification itself has no side effects) + * - other (corrupt metadata, unknown service, malformed URI, illegal argument) → [DeliveryResult.PermanentFailure] + * + * [deliverBatch] fires all requests via [HttpClient.sendAsync] in parallel instead of blocking + * sequentially on [HttpClient.send]. The JDK [HttpClient] will use HTTP/2 multiplexing when + * supported by the server; otherwise it will fall back to HTTP/1.1 (potentially with multiple + * connections) — either way, requests can be overlapped without explicit per-host grouping. */ class HttpMessageDeliverer @JvmOverloads constructor( private val urlResolver: ServiceUrlResolver, @@ -32,41 +42,107 @@ class HttpMessageDeliverer @JvmOverloads constructor( override val type: String = HttpDeliveryInfo.TYPE override fun deliver(entry: OutboxEntry): DeliveryResult = try { - val info = HttpDeliveryInfo.deserialize(entry.deliveryMetadata) - val url = urlResolver.resolve(info.serviceName) + info.endpointPath + val request = buildRequest(entry) + val response = httpClient.send(request, HttpResponse.BodyHandlers.ofString()) + classifyResponse(response.statusCode(), response.body()) + } catch (e: InterruptedException) { + Thread.currentThread().interrupt() + classifyThrowable(e) + } catch (e: Exception) { + classifyThrowable(e) + } - val request = - HttpRequest - .newBuilder() - .uri(URI.create(url)) - .timeout(REQUEST_TIMEOUT) - .header("Content-Type", "application/json") - .method( - info.httpMethod.name, - HttpRequest.BodyPublishers.ofString(entry.payload), - ) - .apply { info.headers.forEach { (k, v) -> header(k, v) } } - .build() + /** + * Fires all requests via [HttpClient.sendAsync] before awaiting any of them, so N requests + * incur ~max(latency) instead of ~sum(latency). A synchronous failure building one entry's + * request (corrupt metadata, unknown service) is isolated to that entry and does not prevent + * the rest of the batch from firing. + */ + override fun deliverBatch(entries: List): List { + if (entries.isEmpty()) return emptyList() - val response = httpClient.send(request, HttpResponse.BodyHandlers.ofString()) - val status = response.statusCode() - val body = response.body() + val attempts: List> = entries.map { entry -> entry to fireOne(entry) } - when { - status in 200..299 -> DeliveryResult.Success - status in retriableStatusCodes -> DeliveryResult.RetriableFailure("HTTP $status: $body") - else -> DeliveryResult.PermanentFailure("HTTP $status: $body") + return attempts.map { (entry, attempt) -> + val result = when (attempt) { + is SendAttempt.ImmediateFailure -> attempt.result + is SendAttempt.InFlight -> awaitResult(attempt.future) + } + DeliveryOutcome(entry, result) } - } catch (e: JsonProcessingException) { - // Subtype of IOException, but corrupt metadata won't fix itself — classify before IOException. - DeliveryResult.PermanentFailure(e.message ?: e.javaClass.simpleName) - } catch (e: IOException) { - DeliveryResult.RetriableFailure(e.message ?: "Connection failed") + } + + /** + * Awaits via [CompletableFuture.get] rather than [CompletableFuture.join] so that an + * interrupt on the caller (processor thread shutdown/backpressure) is observed instead of + * being ignored — the flag is restored here, on the interrupted thread itself, and the + * remaining not-yet-awaited futures in the batch will then also fail fast as retriable. + */ + private fun awaitResult(future: CompletableFuture): DeliveryResult = try { + future.get() } catch (e: InterruptedException) { Thread.currentThread().interrupt() - DeliveryResult.RetriableFailure(e.message ?: "Interrupted") + classifyThrowable(e) + } catch (e: Exception) { + // fireOne()'s .exceptionally already converts every failure into a normal completion, so + // get() should never throw ExecutionException/CancellationException in practice — but + // deliverBatch's contract is MUST NOT throw, so this stays defensive. + classifyThrowable(e.cause ?: e) + } + + private fun fireOne(entry: OutboxEntry): SendAttempt = try { + val request = buildRequest(entry) + val future: CompletableFuture = httpClient + .sendAsync(request, HttpResponse.BodyHandlers.ofString()) + .thenApply { response -> classifyResponse(response.statusCode(), response.body()) } + .exceptionally { e -> classifyThrowable(e.cause ?: e) } + SendAttempt.InFlight(future) } catch (e: Exception) { - DeliveryResult.PermanentFailure(e.message ?: e.javaClass.simpleName) + SendAttempt.ImmediateFailure(classifyThrowable(e)) + } + + private fun buildRequest(entry: OutboxEntry): HttpRequest { + val info = HttpDeliveryInfo.deserialize(entry.deliveryMetadata) + val url = urlResolver.resolve(info.serviceName) + info.endpointPath + + return HttpRequest + .newBuilder() + .uri(URI.create(url)) + .timeout(REQUEST_TIMEOUT) + .setHeader("Content-Type", "application/json") + .method( + info.httpMethod.name, + HttpRequest.BodyPublishers.ofString(entry.payload), + ) + .apply { info.headers.forEach { (k, v) -> setHeader(k, v) } } + .build() + } + + private fun classifyResponse(status: Int, body: String): DeliveryResult = when { + status in 200..299 -> DeliveryResult.Success + status in retriableStatusCodes -> DeliveryResult.RetriableFailure("HTTP $status: $body") + else -> DeliveryResult.PermanentFailure("HTTP $status: $body") + } + + private fun classifyThrowable(e: Throwable): DeliveryResult { + val message = e.message ?: e.javaClass.simpleName + return when (e) { + // Subtypes of IOException, but neither fixes itself on retry — classify before IOException. + is JsonProcessingException -> DeliveryResult.PermanentFailure(message) + is SSLException -> DeliveryResult.PermanentFailure(message) + is IOException -> DeliveryResult.RetriableFailure(message) + // Deliberately side-effect free: this can run on an HttpClient completion thread + // (via fireOne()'s .exceptionally), not the caller's thread, so interrupting + // "current thread" here would interrupt an unrelated shared pool thread. Callers on + // a genuinely interrupted thread (deliver(), awaitResult()) restore the flag themselves. + is InterruptedException -> DeliveryResult.RetriableFailure(message) + else -> DeliveryResult.PermanentFailure(message) + } + } + + private sealed interface SendAttempt { + data class InFlight(val future: CompletableFuture) : SendAttempt + data class ImmediateFailure(val result: DeliveryResult) : SendAttempt } companion object { diff --git a/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererBatchConcurrencyTest.kt b/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererBatchConcurrencyTest.kt new file mode 100644 index 0000000..b766a35 --- /dev/null +++ b/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererBatchConcurrencyTest.kt @@ -0,0 +1,68 @@ +package com.softwaremill.okapi.http + +import com.github.tomakehurst.wiremock.WireMockServer +import com.github.tomakehurst.wiremock.client.WireMock.aResponse +import com.github.tomakehurst.wiremock.client.WireMock.post +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.OutboxMessage +import io.kotest.core.spec.style.FunSpec +import io.kotest.matchers.longs.shouldBeLessThan +import io.kotest.matchers.shouldBe +import java.time.Instant + +private const val SCHEDULING_JITTER_TOLERANCE_MS = 200L + +private fun entry(suffix: String): OutboxEntry { + val info = httpDeliveryInfo { + serviceName = "svc" + endpointPath = "/test" + } + return OutboxEntry.createPending(OutboxMessage("evt-$suffix", """{"k":"v-$suffix"}"""), info, Instant.now()) +} + +/** + * Proves `deliverBatch` fires requests concurrently rather than one-at-a-time, using WireMock's + * request journal (`allServeEvents`) rather than just overall wall-clock time — the timestamp + * spread across requests is direct evidence they overlapped in flight. + */ +class HttpMessageDelivererBatchConcurrencyTest : FunSpec({ + val wiremock = WireMockServer(wireMockConfig().dynamicPort()) + val deliverer by lazy { + HttpMessageDeliverer({ "http://localhost:${wiremock.port()}" }) + } + + beforeSpec { wiremock.start() } + afterSpec { wiremock.stop() } + beforeEach { wiremock.resetAll() } + + test("deliverBatch fires N requests concurrently: request-log timestamps overlap within one delay window") { + val batchSize = 10 + val delayMs = 300L + wiremock.stubFor(post(urlEqualTo("/test")).willReturn(aResponse().withStatus(200).withFixedDelay(delayMs.toInt()))) + val entries = (1..batchSize).map { entry("e$it") } + + val start = System.nanoTime() + val results = deliverer.deliverBatch(entries) + val elapsedMs = (System.nanoTime() - start) / 1_000_000 + + results.forEach { (_, r) -> r shouldBe DeliveryResult.Success } + + // Sequential delivery would take ~batchSize * delayMs; parallel delivery should land near + // one delay window plus scheduling overhead — well under half the sequential bound. + elapsedMs.shouldBeLessThan(batchSize * delayMs / 2) + + // Direct evidence of overlap: every request's logged arrival time falls within a single + // delay window of each other, i.e. WireMock received all of them before the first one + // could possibly have completed and freed up a sequential caller to send the next. + val loggedTimestamps = wiremock.allServeEvents.map { it.request.loggedDate.time } + loggedTimestamps.size shouldBe batchSize + val spreadMs = loggedTimestamps.maxOrNull()!! - loggedTimestamps.minOrNull()!! + // A fully sequential implementation would spread these across ~(batchSize - 1) * delayMs + // (2700 ms here); a small tolerance above one delay window still clearly distinguishes + // "overlapped" from "sequential" while absorbing CI scheduling/connection-setup jitter. + spreadMs.shouldBeLessThan(delayMs + SCHEDULING_JITTER_TOLERANCE_MS) + } +}) diff --git a/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererBatchTest.kt b/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererBatchTest.kt new file mode 100644 index 0000000..2de2992 --- /dev/null +++ b/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererBatchTest.kt @@ -0,0 +1,136 @@ +package com.softwaremill.okapi.http + +import com.github.tomakehurst.wiremock.WireMockServer +import com.github.tomakehurst.wiremock.client.WireMock.aResponse +import com.github.tomakehurst.wiremock.client.WireMock.post +import com.github.tomakehurst.wiremock.client.WireMock.urlEqualTo +import com.github.tomakehurst.wiremock.core.WireMockConfiguration.wireMockConfig +import com.github.tomakehurst.wiremock.http.Fault +import com.softwaremill.okapi.core.DeliveryOutcome +import com.softwaremill.okapi.core.DeliveryResult +import com.softwaremill.okapi.core.OutboxEntry +import com.softwaremill.okapi.core.OutboxMessage +import io.kotest.core.spec.style.FunSpec +import io.kotest.matchers.shouldBe +import io.kotest.matchers.types.shouldBeInstanceOf +import java.net.http.HttpClient +import java.time.Duration +import java.time.Instant + +private fun entry(suffix: String, path: String = "/test", metadataOverride: String? = null): OutboxEntry { + val info = httpDeliveryInfo { + serviceName = "svc" + endpointPath = path + } + val base = OutboxEntry.createPending(OutboxMessage("evt-$suffix", """{"k":"v-$suffix"}"""), info, Instant.now()) + return if (metadataOverride != null) base.copy(deliveryMetadata = metadataOverride) else base +} + +class HttpMessageDelivererBatchTest : FunSpec({ + val wiremock = WireMockServer(wireMockConfig().dynamicPort().dynamicHttpsPort()) + val deliverer by lazy { + HttpMessageDeliverer(ServiceUrlResolver { "http://localhost:${wiremock.port()}" }) + } + + beforeSpec { wiremock.start() } + afterSpec { wiremock.stop() } + beforeEach { wiremock.resetAll() } + + test("deliverBatch on empty input returns empty list and fires no requests") { + deliverer.deliverBatch(emptyList()) shouldBe emptyList() + wiremock.allServeEvents.size shouldBe 0 + } + + test("deliverBatch with all-success preserves input order") { + wiremock.stubFor(post(urlEqualTo("/test")).willReturn(aResponse().withStatus(200))) + val entries = listOf(entry("a"), entry("b"), entry("c")) + + val results = deliverer.deliverBatch(entries) + + results.size shouldBe 3 + results.map { it.entry } shouldBe entries + results.forEach { (_, r) -> r shouldBe DeliveryResult.Success } + } + + test("deliverBatch with mixed status codes classifies each entry independently, in order") { + wiremock.stubFor(post(urlEqualTo("/ok")).willReturn(aResponse().withStatus(200))) + wiremock.stubFor(post(urlEqualTo("/retriable")).willReturn(aResponse().withStatus(500))) + wiremock.stubFor(post(urlEqualTo("/permanent")).willReturn(aResponse().withStatus(400))) + val entries = listOf( + entry("a", path = "/ok"), + entry("b", path = "/retriable"), + entry("c", path = "/permanent"), + ) + + val results = deliverer.deliverBatch(entries) + + results.map { it.entry } shouldBe entries + results[0].result shouldBe DeliveryResult.Success + results[1].result.shouldBeInstanceOf() + results[2].result.shouldBeInstanceOf() + } + + test("deliverBatch with all-fail batch classifies every entry as RetriableFailure") { + wiremock.stubFor(post(urlEqualTo("/test")).willReturn(aResponse().withStatus(503))) + val entries = listOf(entry("a"), entry("b"), entry("c")) + + val results = deliverer.deliverBatch(entries) + + results.size shouldBe 3 + results.forEach { (_, r) -> r.shouldBeInstanceOf() } + } + + test("deliverBatch connection reset -> RetriableFailure for the affected entry") { + wiremock.stubFor( + post(urlEqualTo("/test")).willReturn(aResponse().withFault(Fault.CONNECTION_RESET_BY_PEER)), + ) + + val results = deliverer.deliverBatch(listOf(entry("a"))) + + results[0].result.shouldBeInstanceOf() + } + + test("deliverBatch poison-pill metadata yields PermanentFailure for bad entry, others unaffected") { + wiremock.stubFor(post(urlEqualTo("/test")).willReturn(aResponse().withStatus(200))) + val good1 = entry("good1") + val poisoned = entry("bad", metadataOverride = "{not valid http info json}") + val good2 = entry("good2") + + val results = deliverer.deliverBatch(listOf(good1, poisoned, good2)) + + results.size shouldBe 3 + results.map { it.entry } shouldBe listOf(good1, poisoned, good2) + results[0].result shouldBe DeliveryResult.Success + results[1].result.shouldBeInstanceOf() + results[2].result shouldBe DeliveryResult.Success + // Only the two well-formed entries actually reached the server. + wiremock.allServeEvents.size shouldBe 2 + } + + test("deliverBatch with a custom HttpClient uses it for every entry in the batch") { + wiremock.stubFor(post(urlEqualTo("/test")).willReturn(aResponse().withStatus(200))) + val customClient = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(2)).build() + val custom = HttpMessageDeliverer(ServiceUrlResolver { "http://localhost:${wiremock.port()}" }, customClient) + val entries = listOf(entry("a"), entry("b")) + + val results = custom.deliverBatch(entries) + + results.forEach { (_, r) -> r shouldBe DeliveryResult.Success } + wiremock.allServeEvents.size shouldBe 2 + } + + test("deliverBatch TLS handshake failure -> PermanentFailure (does not throw)") { + // WireMock's HTTPS port serves its own self-signed cert, which the JDK's default trust + // store rejects -- a deterministic SSLHandshakeException (cert/config problem, won't fix + // itself on retry). Deliberately not testing this via https:// against the *plaintext* + // port: that failure mode depends on exactly how/when the peer aborts the connection, and + // the exception type it produces is a genuine JDK-level race -- usually SSLException, but + // occasionally java.net.http.HttpConnectTimeoutException (an IOException, misclassified + // as RetriableFailure), causing sporadic CI failures. + val tlsDeliverer = HttpMessageDeliverer(ServiceUrlResolver { "https://localhost:${wiremock.httpsPort()}" }) + + val outcome: DeliveryOutcome = tlsDeliverer.deliverBatch(listOf(entry("a"))).single() + + outcome.result.shouldBeInstanceOf() + } +}) diff --git a/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererInterruptionTest.kt b/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererInterruptionTest.kt new file mode 100644 index 0000000..78e560b --- /dev/null +++ b/okapi-http/src/test/kotlin/com/softwaremill/okapi/http/HttpMessageDelivererInterruptionTest.kt @@ -0,0 +1,110 @@ +package com.softwaremill.okapi.http + +import com.github.tomakehurst.wiremock.WireMockServer +import com.github.tomakehurst.wiremock.client.WireMock.aResponse +import com.github.tomakehurst.wiremock.client.WireMock.post +import com.github.tomakehurst.wiremock.client.WireMock.urlEqualTo +import com.github.tomakehurst.wiremock.core.WireMockConfiguration.wireMockConfig +import com.softwaremill.okapi.core.DeliveryOutcome +import com.softwaremill.okapi.core.DeliveryResult +import com.softwaremill.okapi.core.OutboxEntry +import com.softwaremill.okapi.core.OutboxMessage +import io.kotest.core.spec.style.FunSpec +import io.kotest.matchers.shouldBe +import io.kotest.matchers.types.shouldBeInstanceOf +import java.time.Instant + +private fun entry(suffix: String): OutboxEntry { + val info = httpDeliveryInfo { + serviceName = "svc" + endpointPath = "/test" + } + return OutboxEntry.createPending(OutboxMessage("evt-$suffix", """{"k":"v-$suffix"}"""), info, Instant.now()) +} + +/** + * Polls WireMock's request journal until [count] requests have been received, rather than relying + * on [Thread.State] or a fixed sleep — a thread blocked inside the JDK HTTP client isn't guaranteed + * to report `WAITING`/`TIMED_WAITING` on every JVM/OS combination, but the request actually + * reaching the (deliberately delayed) stub is an unambiguous, environment-independent signal that + * the caller is now blocked awaiting the response. + */ +private fun awaitRequestsReceived(wiremock: WireMockServer, count: Int, timeoutMs: Long = 5_000) { + val deadline = System.currentTimeMillis() + timeoutMs + while (wiremock.allServeEvents.size < count) { + check(System.currentTimeMillis() < deadline) { + "Expected $count request(s), only ${wiremock.allServeEvents.size} received" + } + Thread.sleep(5) + } +} + +/** + * A timed [Thread.join] only establishes a happens-before edge for the joined thread's writes if + * it returns because the thread actually terminated, not because the timeout elapsed — if it + * times out and the thread happens to die a moment later, `isAlive` can read `false` with no + * visibility guarantee for what the thread wrote. Confirming termination and then joining again + * unconditionally (returns immediately on an already-dead thread) closes that gap. + */ +private fun awaitTermination(thread: Thread, timeoutMs: Long = 5_000) { + thread.join(timeoutMs) + check(!thread.isAlive) { "Thread did not terminate within ${timeoutMs}ms" } + thread.join() +} + +/** + * Proves that interrupting the calling thread while it is blocked inside [HttpMessageDeliverer] + * is observed promptly and the interrupt flag ends up restored on that same thread — not on some + * unrelated `HttpClient` completion thread — for both the synchronous and batched delivery paths. + */ +class HttpMessageDelivererInterruptionTest : FunSpec({ + val wiremock = WireMockServer(wireMockConfig().dynamicPort()) + val deliverer by lazy { + HttpMessageDeliverer({ "http://localhost:${wiremock.port()}" }) + } + + beforeSpec { wiremock.start() } + afterSpec { wiremock.stop() } + beforeEach { wiremock.resetAll() } + + test("interrupting the caller during deliver() yields RetriableFailure and restores the flag on that thread") { + wiremock.stubFor(post(urlEqualTo("/test")).willReturn(aResponse().withStatus(200).withFixedDelay(2_000))) + + var result: DeliveryResult? = null + var interruptedAfterReturn = false + val thread = Thread { + result = deliverer.deliver(entry("a")) + interruptedAfterReturn = Thread.currentThread().isInterrupted + } + thread.start() + awaitRequestsReceived(wiremock, count = 1) + thread.interrupt() + awaitTermination(thread) + + result.shouldBeInstanceOf() + interruptedAfterReturn shouldBe true + } + + test( + "interrupting the caller while awaiting deliverBatch yields RetriableFailure for every entry " + + "and restores the flag on that thread", + ) { + wiremock.stubFor(post(urlEqualTo("/test")).willReturn(aResponse().withStatus(200).withFixedDelay(2_000))) + val entries = listOf(entry("a"), entry("b"), entry("c")) + + var results: List? = null + var interruptedAfterReturn = false + val thread = Thread { + results = deliverer.deliverBatch(entries) + interruptedAfterReturn = Thread.currentThread().isInterrupted + } + thread.start() + awaitRequestsReceived(wiremock, count = entries.size) + thread.interrupt() + awaitTermination(thread) + + results?.size shouldBe 3 + results?.forEach { it.result.shouldBeInstanceOf() } + interruptedAfterReturn shouldBe true + } +})