diff --git a/docs/guides/client-pool.md b/docs/guides/client-pool.md index f9b36a431..0edd1535e 100644 --- a/docs/guides/client-pool.md +++ b/docs/guides/client-pool.md @@ -136,6 +136,50 @@ canonical reference; refer to the per-language builder or constructor for exact old-template sandbox IDs into the shared buffer during a rolling deploy. - `resize(max_idle)` and `release_all_idle()` can be called from any node. +### Connection reuse at high warmup concurrency + +**Problem.** Pool-created sandboxes go through the SDK transport per sandbox, and by +default each sandbox's HTTP client opens its own TCP connections — a warmup create +typically uses 2-4 fresh connections (create + endpoint lookups + readiness probe + +renew). At high `warmup_concurrency` the resulting connection burst can exceed what the +server listener can absorb (for example macOS accept backlog 128), producing intermittent +TCP-level `Connection reset` / `Broken pipe` failures. Failed warmups are retried +(backoff-gated), amplifying attempts several-fold and making fills slower and burstier +than at lower concurrency — measured in the Kotlin SDK benchmark harness at +`warmup_concurrency=1000`: ~80% of warmup attempts failed with 5x attempt amplification. + +**Fix: share one transport connection pool.** Warmup creates then reuse connections +instead of opening new ones. Measured fill of 2000 idles at `warmup_concurrency=1000` +(Kotlin SDK, mock server): + +| shared pool size | fill time | create failures | attempt amplification | +|---|---|---|---| +| 0 (per-sandbox connections) | ~20 s | ~8000 | 5x | +| 100 | ~9 s | ~2700 | 2.3x | +| 200 | ~7 s | ~1300 | 1.7x | +| 500 | ~4 s | 0 | 1x | + +**Kotlin/Java.** Inject a shared `okhttp3.ConnectionPool` through the standard +`ConnectionConfig` — no pool-specific option is needed: + +```kotlin +ConnectionConfig.builder() + .connectionPool(ConnectionPool(500, 5, TimeUnit.MINUTES)) // ~= warmup_concurrency + .build() +``` + +The pool uses this config for every sandbox it creates (warmup, direct create, idle +connect). A user-provided pool is treated as user-managed and is never evicted by the +SDK. A pool-created shared pool sized by `warmup_concurrency` will become the SDK +default once the companion change lands; until then configure it explicitly. + +**Rule of thumb.** `warmup_concurrency` beyond ~200-300 only pays off together with a +shared connection pool — without reuse the extra threads mostly produce +connection-reset retries. If connections cannot be shared, keep +`warmup_concurrency` in the 200-300 range. The same principle applies to any SDK whose +transport opens connections per sandbox; see the language-local transport docs for how +to share a connection pool. + ## Minimal usage ### Python (sync) diff --git a/sdks/sandbox/kotlin/sandbox/src/main/kotlin/com/alibaba/opensandbox/sandbox/pool/PoolWarmupDiagnostics.kt b/sdks/sandbox/kotlin/sandbox/src/main/kotlin/com/alibaba/opensandbox/sandbox/pool/PoolWarmupDiagnostics.kt new file mode 100644 index 000000000..d2330c589 --- /dev/null +++ b/sdks/sandbox/kotlin/sandbox/src/main/kotlin/com/alibaba/opensandbox/sandbox/pool/PoolWarmupDiagnostics.kt @@ -0,0 +1,191 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.alibaba.opensandbox.sandbox.pool + +/** + * Diagnostic counters for the warmup pipeline. + * + * Benchmark/diagnostic aid only — not part of the public API contract. + * Overhead is a single `System.nanoTime()` read plus a locked list append per + * warmup event, and it is a no-op when never read. State is global (all + * pools in the JVM share it), so reset before a single-pool experiment and + * snapshot after it settles. + * + * Measured phases of one warmup task: + * - queue wait: submission -> task picked up by a warmup worker + * - create: `createOneSandbox` (create + readiness + renew) + * - commit: lock + store putIdle when the sandbox enters the idle buffer + * Plus the reconcile-tick cadence (submission driver) and the in-flight + * warmup count trajectory. + */ +object PoolWarmupDiagnostics { + data class PhaseStats( + val count: Long, + val meanMs: Double, + val p50Ms: Long, + val p95Ms: Long, + val maxMs: Long, + ) + + data class Snapshot( + val queueWaitMs: PhaseStats, + val createDurationMs: PhaseStats, + val commitDurationMs: PhaseStats, + val tickIntervalMs: PhaseStats, + val tickDurationMs: PhaseStats, + val submitBurst: PhaseStats, + val submitCalls: Long, + val inFlightPeak: Int, + val inFlightMean: Double, + val createFailures: Map, + ) + + private val lock = Any() + private val queueWaitNanos = ArrayList() + private val createNanos = ArrayList() + private val commitNanos = ArrayList() + private val tickIntervalNanos = ArrayList() + private val tickDurationNanos = ArrayList() + private val submitBursts = ArrayList() + private val failureReasons = LinkedHashMap() + private var submitCalls = 0L + private var inFlightSum = 0L + private var inFlightSamples = 0L + private var inFlightPeak = 0 + private var lastTickNanos = 0L + + fun reset() { + synchronized(lock) { + queueWaitNanos.clear() + createNanos.clear() + commitNanos.clear() + tickIntervalNanos.clear() + tickDurationNanos.clear() + submitBursts.clear() + failureReasons.clear() + submitCalls = 0L + inFlightSum = 0L + inFlightSamples = 0L + inFlightPeak = 0 + lastTickNanos = 0L + } + } + + fun recordQueueWait(nanos: Long) = synchronized(lock) { queueWaitNanos.add(nanos) } + + fun recordCreate(nanos: Long) = synchronized(lock) { createNanos.add(nanos) } + + fun recordCommit(nanos: Long) = synchronized(lock) { commitNanos.add(nanos) } + + /** Records the wall-clock spacing between reconcile ticks plus each tick's execution time. */ + fun recordTick(nowNanos: Long, durationNanos: Long) { + synchronized(lock) { + tickDurationNanos.add(durationNanos) + if (lastTickNanos != 0L) { + tickIntervalNanos.add(nowNanos - lastTickNanos) + } + lastTickNanos = nowNanos + } + } + + fun recordSubmitBurst(size: Int) { + synchronized(lock) { + submitCalls++ + submitBursts.add(size) + } + } + + fun recordInFlight(current: Int) { + synchronized(lock) { + if (current > inFlightPeak) inFlightPeak = current + inFlightSum += current + inFlightSamples++ + } + } + + /** Records a failed createOneSandbox attempt: exception class + first words of the message. */ + fun recordCreateFailure(failure: Throwable) { + val message = failure.message?.trim()?.take(80) ?: "" + val key = "${failure.javaClass.simpleName}: $message" + synchronized(lock) { + failureReasons[key] = (failureReasons[key] ?: 0L) + 1 + } + } + + fun snapshot(): Snapshot { + val queueWait: LongArray + val create: LongArray + val commit: LongArray + val tickInterval: LongArray + val tickDuration: LongArray + val bursts: IntArray + val calls: Long + val peak: Int + val meanInFlight: Double + val failures: Map + synchronized(lock) { + queueWait = queueWaitNanos.toLongArray() + create = createNanos.toLongArray() + commit = commitNanos.toLongArray() + tickInterval = tickIntervalNanos.toLongArray() + tickDuration = tickDurationNanos.toLongArray() + bursts = submitBursts.toIntArray() + calls = submitCalls + peak = inFlightPeak + meanInFlight = if (inFlightSamples == 0L) 0.0 else inFlightSum.toDouble() / inFlightSamples + failures = failureReasons.toMap() + } + return Snapshot( + queueWaitMs = phase(queueWait), + createDurationMs = phase(create), + commitDurationMs = phase(commit), + tickIntervalMs = phase(tickInterval), + tickDurationMs = phase(tickDuration), + submitBurst = burstPhase(bursts), + submitCalls = calls, + inFlightPeak = peak, + inFlightMean = meanInFlight, + createFailures = failures, + ) + } + + private fun phase(nanos: LongArray): PhaseStats { + if (nanos.isEmpty()) return PhaseStats(0, 0.0, 0L, 0L, 0L) + val sorted = nanos.copyOf() + sorted.sort() + return PhaseStats( + count = sorted.size.toLong(), + meanMs = sorted.average() / 1_000_000.0, + p50Ms = sorted[(sorted.size - 1) / 2] / 1_000_000, + p95Ms = sorted[((sorted.size - 1) * 95) / 100] / 1_000_000, + maxMs = sorted.last() / 1_000_000, + ) + } + + private fun burstPhase(bursts: IntArray): PhaseStats { + if (bursts.isEmpty()) return PhaseStats(0, 0.0, 0L, 0L, 0L) + val sorted = bursts.copyOf() + sorted.sort() + return PhaseStats( + count = sorted.size.toLong(), + meanMs = sorted.average(), + p50Ms = sorted[(sorted.size - 1) / 2].toLong(), + p95Ms = sorted[((sorted.size - 1) * 95) / 100].toLong(), + maxMs = sorted.last().toLong(), + ) + } +} diff --git a/sdks/sandbox/kotlin/sandbox/src/main/kotlin/com/alibaba/opensandbox/sandbox/pool/SandboxPool.kt b/sdks/sandbox/kotlin/sandbox/src/main/kotlin/com/alibaba/opensandbox/sandbox/pool/SandboxPool.kt index dcbd7a38c..beac87498 100644 --- a/sdks/sandbox/kotlin/sandbox/src/main/kotlin/com/alibaba/opensandbox/sandbox/pool/SandboxPool.kt +++ b/sdks/sandbox/kotlin/sandbox/src/main/kotlin/com/alibaba/opensandbox/sandbox/pool/SandboxPool.kt @@ -906,6 +906,7 @@ class SandboxPool internal constructor( if (!isCurrentRun(run) || lifecycleState.get() != LifecycleState.RUNNING) return if (!isPoolNamespaceActive()) return val reconcileConfig = config.withMaxIdle(resolveMaxIdle()) + val tickStart = System.nanoTime() try { run.primaryOwned.set( PoolReconciler.runReconcileTick( @@ -920,6 +921,8 @@ class SandboxPool internal constructor( } catch (e: Exception) { run.primaryOwned.set(false) throw e + } finally { + PoolWarmupDiagnostics.recordTick(System.nanoTime(), System.nanoTime() - tickStart) } } finally { endOperation(run) @@ -999,6 +1002,7 @@ class SandboxPool internal constructor( run: RunContext, count: Int, ) { + PoolWarmupDiagnostics.recordSubmitBurst(count) repeat(count) { if (!isCurrentRun(run) || lifecycleState.get() != LifecycleState.RUNNING || @@ -1024,20 +1028,26 @@ class SandboxPool internal constructor( private inner class TrackedWarmupTask( private val run: RunContext, ) : Runnable { + private val submittedAtNanos = System.nanoTime() private val completed = AtomicBoolean(false) init { run.warmingCount.incrementAndGet() + PoolWarmupDiagnostics.recordInFlight(run.warmingCount.get()) beginOperation(run) } override fun run() { + PoolWarmupDiagnostics.recordQueueWait(System.nanoTime() - submittedAtNanos) + val createStart = System.nanoTime() val outcome = try { WarmupOutcome.Success(createOneSandbox()) } catch (failure: Throwable) { + PoolWarmupDiagnostics.recordCreateFailure(failure) WarmupOutcome.Failure(failure) } + PoolWarmupDiagnostics.recordCreate(System.nanoTime() - createStart) dispatchCompletion(outcome) } @@ -1066,6 +1076,7 @@ class SandboxPool internal constructor( handleWarmupOutcome(run, outcome) } finally { run.warmingCount.decrementAndGet() + PoolWarmupDiagnostics.recordInFlight(run.warmingCount.get()) endOperation(run) // Only successful completions trigger an immediate reconcile. A failed warmup // frees its slot but must not cause an immediate retry: fast-failing creates @@ -1142,6 +1153,7 @@ class SandboxPool internal constructor( ) { var cleanupSource: String? = null run.commitLock.lock() + val commitStart = System.nanoTime() try { val state = lifecycleState.get() if (!isCurrentRun(run) || (state != LifecycleState.RUNNING && state != LifecycleState.DRAINING)) { @@ -1189,6 +1201,7 @@ class SandboxPool internal constructor( } } } finally { + PoolWarmupDiagnostics.recordCommit(System.nanoTime() - commitStart) run.commitLock.unlock() } cleanupSource?.let { source -> diff --git a/tests/benchmark/.gitignore b/tests/benchmark/.gitignore new file mode 100644 index 000000000..129417619 --- /dev/null +++ b/tests/benchmark/.gitignore @@ -0,0 +1,2 @@ +bin/ +results/ diff --git a/tests/benchmark/README.md b/tests/benchmark/README.md new file mode 100644 index 000000000..2b286aee8 --- /dev/null +++ b/tests/benchmark/README.md @@ -0,0 +1,334 @@ +# OpenSandbox Pool Benchmark + +Reproducible, cross-SDK benchmark harness for the sandbox SDK **client pool** +(`SandboxPool` in the Kotlin/JVM SDK; Go/Python/JS pools can reuse the same +mock server). It runs a standalone mock of the lifecycle + execd API with +configurable provisioning latency and fault injection, drives the pool through +scenario workloads, and writes a JSON + Markdown report. + +``` +tests/benchmark/ +├── mockserver/ # standalone Go mock server (lifecycle + execd + control endpoints) +├── kotlin/ # benchmark driver (JVM, uses the published Kotlin SDK) +├── configs/ # mock server scenario configs (latency / fault profiles) +├── run.sh # one-shot orchestration: publish SDK -> start mock -> run driver +└── results/ # reports and mock logs (gitignored) +``` + +## Prerequisites + +- Go (any recent version; mock server uses stdlib only) +- JDK 17+ (driver + Kotlin SDK; `run.sh` picks a suitable JDK automatically + on macOS) +- A Gradle wrapper is included under `kotlin/`. The SDK sources are referenced + in place (composite build), so no separate install is needed. + +## Quick start + +```bash +# default config: create/delete 300-800ms, execd ping 100ms, others 50-100ms +./run.sh + +# smoke run with fast provisioning +./run.sh --mock-config configs/fast.json -- --scenarios cold-start,warm-latency + +# full control over the driver +./run.sh -- --max-idle 50 --warmup-concurrency 10 --steady-duration-s 120 +``` + +`run.sh` performs three steps: + +1. Builds the mock server (`go build` + exec). +2. Starts the mock server. +3. Runs the driver: `./gradlew run` with forwarded `--key value` args. + +The Kotlin SDK is **built from source**: the driver uses a Gradle composite +build (`includeBuild` in `kotlin/settings.gradle.kts`), so the benchmark always +runs the checked-out SDK code and picks up SDK changes without any +publish/install step. `com.alibaba.opensandbox:sandbox:1.0.18` in +`kotlin/build.gradle.kts` is a module coordinate that the composite build +substitutes with the local `:sandbox` project (the version is informational). + +Reports and every artifact of a run land in one directory, +`results/run-/`: + +``` +results/run-/ +├── report.json # full results: per-scenario metrics + per-scenario QPS + end stats +├── report.md # human-readable view of report.json +├── server-timeseries-.csv # per-second: alive + each API request count +├── client-.csv # per-500ms: threads, heap MB, idle, inFlight, degraded, backoff +├── mock-stats-end.json # raw /__stats snapshot at end of run +├── mock-config.json # mock server config used for this run +├── driver-args.txt # driver CLI arguments +└── mockserver.log # mock server log +``` + +The `server-timeseries-*.csv` files merge the mock's per-second per-API QPS +series with the alive-sandbox gauge into one table +(`second,alive,create,delete,get,renew,endpoint,execd.ping,execd.other`) for +direct plotting; the full per-second series is also embedded in `report.json`. + +## Mock server + +The mock implements the API surface the SDK pool actually drives, per +`specs/sandbox-lifecycle.yml` and `specs/execd-api.yaml`: + +| Endpoint | Behavior | +|---|---| +| `POST /v1/sandboxes` | Simulated provisioning: sleeps `createLatencyMs`, returns `Pending`, flips to `Running` after `bootDelayMs` | +| `GET /v1/sandboxes/{id}` | Sandbox info; 404 after kill/expiry | +| `POST /v1/sandboxes/{id}/renew-expiration` | Applies server-side TTL | +| `DELETE /v1/sandboxes/{id}` | Marks terminated; execd stops responding | +| `GET /v1/sandboxes/{id}/endpoints/{port}` | Returns the execd URL + a per-sandbox access token | +| execd (`/ping`, anything) | 200 when the sandbox is booted and healthy; **404** while `Pending`/expired/poisoned | + +**Why execd answers 404 while a sandbox is not ready**: the SDK's readiness +check polls execd `/ping` (`Sandbox.checkReady`), and its retry interceptor +retries 5xx and transport errors with backoff. A 404 is non-retryable, so the +client polls at its configured `healthCheckPollingInterval` instead of paying +policy backoff — keeping the benchmark's latency numbers clean. The +`execdFailureRate` fault injects 500s when you *want* the retry policy +engaged. + +**Why boot state lives in execd**: readiness is decided by execd `/ping`, not +the lifecycle GET, so the mock attributes each execd request to a sandbox via +the endpoint token and only answers once that sandbox is `Running`. + +Server-side state and counters are exposed for the driver to validate pool +behavior (no over-creation, stale cleanup, hit ratio): + +- `GET /__stats` — counters + per-route QPS/latency + live config +- `POST /__config` — runtime fault injection: `createFailureRate`, + `execdFailureRate`, `bootDelayMs`, `poisonExisting` (flips all alive + sandboxes to a failing state, simulating stale idles); latency knobs are + runtime-mutable too (`createLatencyMs`, `latencyOverrides` — the latter + replaces the whole per-route map) +- `POST /__reset` — zero counters and QPS history + +### Per-API QPS tracking + +The mock records an exact per-second count for every request on each route +(`lifecycle.create|get|delete|renew|endpoint`, `execd.ping|other`) in a ring +buffer (`-stats-window-sec`, default 1800s). `/__stats` returns per route: + +```json +"lifecycle.create": { + "total": 244, "qps1s": 2.0, "qps5s": 3.2, "qps60s": 1.8, + "avgMs": 100.2, "maxMs": 101, + "seriesStartUnixSec": 1786678587, "series": [2, 4, 2, ...] +} +``` + +- `total` counts since the last `__reset`; `series` is the per-second request + count covering `seriesStartUnixSec .. now` (older than the window is + dropped; the driver resets before each scenario so each section is + self-contained). +- The driver attaches a QPS snapshot to **every scenario section** + (`results..mockQps` in `report.json`), and the end-of-run totals + live under `results.mockServerStats`. This is the server-side ground truth + for offline analysis: warmup bursts, reconcile-tick load, replenish spikes, + and per-API request mixes (e.g. how many `renew-expiration` calls each + acquire generates). + +For runs longer than the window, poll `GET /__stats` from your analysis tool +and accumulate the series yourself. + +### Mock config (JSON) + +```json +{ + "createLatencyMs": { "distribution": "lognormal", "meanMs": 800, "stddevMs": 400, "minMs": 50 }, + "createFailureRate": 0.0, + "bootDelayMs": 300, + "execdFailureRate": 0.0, + "defaultTtlSeconds": 3600, + "latencyOverrides": { + "lifecycle.renew": { "distribution": "fixed", "meanMs": 20 }, + "lifecycle.endpoint": { "distribution": "lognormal", "meanMs": 100, "stddevMs": 30, "minMs": 10 }, + "execd.ping": { "distribution": "fixed", "meanMs": 40 } + } +} +``` + +**Response time model** — three independent knobs: + +| Knob | What it controls | +|---|---| +| `createLatencyMs` | `POST /v1/sandboxes` response time (the server sleeps this long) | +| `bootDelayMs` | How long a created sandbox stays `Pending`; during this window execd pings fail **immediately** with 404 (no latency is paid) and the SDK readiness poll keeps retrying at its polling interval | +| `latencyOverrides` | Response time per route: `lifecycle.get`, `lifecycle.delete`, `lifecycle.renew`, `lifecycle.endpoint`, `execd.ping`, `execd.other`. Routes without an override respond immediately; an override for `lifecycle.create` replaces `createLatencyMs`. Execd route latency only applies once the sandbox is booted — not-ready probes fail fast | + +Default profile (no `-config`): create/delete uniform **300-800ms**, execd +`/ping` fixed **100ms**, all other APIs uniform **50-100ms**. + +The readiness sequence a client observes is therefore: create latency, then a +few fast `404` polls while the sandbox boots, then one successful ping — +typically one ping for the default profile (fixed 100ms ping vs. max 300ms +boot window). + +So the full create-to-ready time a client observes is +`createLatencyMs + bootDelayMs` plus one successful ping (once booted, the +ready execd pays its route latency). All latency knobs take a `LatencySpec`: +`distribution` is `uniform` (random between `minMs` and `maxMs`), `fixed` +(always `meanMs`), or `lognormal` (`meanMs`/`stddevMs`, floored at `minMs`). +They are also runtime-mutable via `POST /__config`; sending `latencyOverrides` +replaces the whole per-route map, and the resulting response times are visible +in `/__stats` per-route `avgMs`/`maxMs` and the QPS series. + +Presets: `default.json`, `fast.json` (smoke tests), `slow.json`. + +## Driver + +Run the driver standalone (mock already up): + +```bash +cd kotlin +./gradlew --console=plain run --args="--mock-base-url http://127.0.0.1:18080 --scenarios all" +``` + +### Scenarios + +| Scenario | What it measures | +|---|---| +| `cold-start` | Time from `pool.start()` until idle buffer is full; over-creation check (server `created` vs `maxIdle`) | +| `warm-latency` | acquire p50/p90/p95/p99/p999 + hit ratio from a warm pool (`N` workers × `M` rounds) | +| `steady-state` | Sustained acquires/sec under concurrent loaders with hold time; idle trajectory (min/mean/empty ratio) | +| `replenish-lag` | Time for a released idle slot to be refilled (completion-driven reconcile) | +| `failure-injection` | Pool behavior at `createFailureRate` 60%: success rate, backoff, DEGRADED transition, recovery after fault removal | +| `stale-idle` | Poisoned idle candidates (`--stale-poison-rate`, default 1.0 = all): retry cost, stale cleanup, refill with fresh sandboxes | +| `idle-expiry` | Self-healing under short server-side TTL: reap + recreate keeps the buffer near `maxIdle` | +| `resize` | Shrink accuracy and speed (excess idles killed by reconcile), regrow speed, server-side alive check | +| `shutdown-race` | Workers hammering acquire while the pool drains: success rate in the running phase vs rejections during DRAINING (`poolNotRunning`), drain duration | +| `store-outage` | State-store outage (OSEP-0005): DIRECT_CREATE must fall through to direct create and stay available; FAIL_FAST must fail closed with `storeUnavailable`; refill after recovery | + +### Driver options + +| Option | Default | Meaning | +|---|---|---| +| `--scenarios` | `all` | Comma-separated list: `cold-start`, `warm-latency`, `steady-state`, `replenish-lag`, `failure-injection`, `stale-idle`, `idle-expiry` | +| `--mock-base-url` | `http://127.0.0.1:18080` | Mock lifecycle base URL | +| `--report-dir` | `results/run-` | Report output directory (`run.sh` passes an absolute path) | +| `--max-idle` | `20` | Pool idle-buffer target | +| `--warmup-concurrency` | `4` | Concurrent warmup creation workers | +| `--shared-connection-pool-size` | `500` | Inject one shared OkHttp `ConnectionPool` with this many idle slots across all sandbox clients (matches the guidance in the [client pool guide](/guides/client-pool#connection-reuse-at-high-warmup-concurrency)); `0` = per-sandbox fresh connections (reproduces the connection-reset pathology) | +| `--reconcile-interval-ms` | `1000` | Pool reconcile tick interval | +| `--idle-timeout-s` | `1800` | Server-side TTL applied to pool-created sandboxes | +| `--acquire-min-remaining-ttl-s` | `0` | Idle entries with less remaining TTL than this are discarded on acquire; `0` = SDK auto default (`min(60s, idleTimeout/2)`) | +| `--primary-lock-ttl-s` | `0` | Distributed primary-lock TTL; `0` = SDK default (60s). No effect with the in-memory state store (single node always holds the lock) | +| `--degraded-threshold` | `0` | Consecutive create failures before the pool enters DEGRADED; `0` = SDK default (3) | +| `--acquire-ready-timeout-ms` | `15000` | `checkReady` timeout when acquiring (idle connect + direct create) | +| `--warmup-ready-timeout-ms` | `15000` | `checkReady` timeout for warmup creations | +| `--health-check-polling-interval-ms` | `200` | `checkReady` probe interval (execd ping cadence) | +| `--cold-start-timeout-ms` | `120000` | Max time to wait for the pool to fill | +| `--warm-workers` | `16` | Loader threads in `warm-latency` | +| `--warm-rounds-per-worker` | `150` | Acquire rounds per `warm-latency` worker | +| `--steady-workers` | `16` | Loader threads in `steady-state` | +| `--steady-duration-s` | `60` | `steady-state` run duration | +| `--hold-min-ms` | `1000` | Lower bound of the random hold time per acquired sandbox in `steady-state` | +| `--hold-max-ms` | `5000` | Upper bound of the random hold time per acquired sandbox in `steady-state` | +| `--replenish-rounds` | `20` | Kill-and-wait repetitions in `replenish-lag` | +| `--replenish-wait-timeout-ms` | `15000` | Max time to wait for one replenished slot | +| `--failure-create-rate` | `0.6` | Create failure rate injected in `failure-injection` | +| `--failure-acquires` | `60` | Acquire attempts in `failure-injection` | +| `--stale-acquires` | `100` | Acquire attempts in `stale-idle` | +| `--stale-retries` | `3` | Pool `maxAcquireRetries` in `stale-idle` (idle candidates tried per acquire) | +| `--stale-poison-rate` | `1.0` | Fraction (0..1] of alive sandboxes poisoned in `stale-idle`; `1.0` = all (partial rates simulate real-world partial failure) | +| `--stale-acquire-ready-timeout-ms` | `3000` | `acquireReadyTimeout` in `stale-idle`; short because the SDK polls a failing execd for the full timeout before discarding a candidate | +| `--idle-expiry-idle-timeout-s` | `20` | `idleTimeout` in `idle-expiry` (short TTL so server-side expiry is exercised) | +| `--idle-expiry-duration-s` | `40` | `idle-expiry` run duration | + +Each scenario resets the mock's counters and QPS history first, so +`report.json`'s per-scenario QPS sections cover exactly that scenario. + +### Reproducing a production pool profile + +Any `PoolConfig`-level profile maps 1:1 onto driver knobs. Example — a +large-pool / high-frequency-acquire / high-frequency-replenish production +profile (`maxIdle=13815, warmupConcurrency=1000, idleTtl=4h, +acquireMinRemainingTtl=15min, reconcile=30s, acquireReady=60s, +warmupReady=180s, primaryLockTtl=360s, degradedThreshold=5`): + +```bash +./run.sh -- --max-idle 13815 --warmup-concurrency 1000 \ + --reconcile-interval-ms 30000 --idle-timeout-s 14400 \ + --acquire-min-remaining-ttl-s 900 --primary-lock-ttl-s 360 \ + --degraded-threshold 5 --acquire-ready-timeout-ms 60000 \ + --warmup-ready-timeout-ms 180000 --cold-start-timeout-ms 300000 \ + --scenarios cold-start,steady-state --steady-workers 300 \ + --steady-duration-s 120 --hold-min-ms 200 --hold-max-ms 2000 +``` + +Notes for this kind of run: + +- `--cold-start-timeout-ms` must cover a full fill: with + `warmupConcurrency=1000` and the default profile (~2-6s per sandbox) a + 13815-sandbox pool fills in roughly 30-90s. +- At this scale each scenario still starts a fresh pool, but the mock's + sandbox registry accumulates across scenarios (previous pools are not + killed on non-graceful shutdown) — budget mock memory accordingly or run + one scenario per invocation. +- For higher acquire frequency, shrink the `execd.ping` latency via + `latencyOverrides` in the mock config: each acquire pays one readiness ping + (100ms by default), which caps the sustainable acquire rate. +- `primaryLockTtl` and `drainTimeout` do not change single-node benchmark + behavior (in-memory state store always grants the lock; the driver never + shuts down gracefully); they are reproduced for config fidelity only. + +For `warm-latency`, keep `maxIdle` comfortably above `--warm-workers` if you +want to measure pure idle-hit latency: when workers outnumber the idle buffer, +acquires drain it and fall through to direct create (which the `hitRatio` +column will show). + +### Warmup throughput and connection reuse + +At high `warmupConcurrency`, per-sandbox HTTP clients opening fresh TCP connections +cause intermittent `Connection reset` failures and retry amplification (see +`warmupPipeline.createFailures` in `cold-start` reports). The problem, evidence, and +configuration guidance (including production `ConnectionConfig` setup) are documented +in the [client pool guide](/guides/client-pool#connection-reuse-at-high-warmup-concurrency). + +In the benchmark, the shared pool is **on by default** (fixed 500 idle +slots), so high-concurrency runs already follow the guidance. Pass +`--shared-connection-pool-size 0` to reproduce the per-sandbox-connection +pathology, or an explicit `N` to sweep pool sizes. + +### What is measured + +| Concern | Where it shows up | +|---|---| +| Acquire latency (p50/90/95/99/999) | `results..latency` | +| Acquire success rate | `successRate` (successful acquires / attempts) | +| Failure breakdown | `latency.failuresByType`: `readyTimeout` / `createFailed`-style `other` / `poolNotRunning` / `poolEmpty` / `acquireFailed` / `storeUnavailable` | +| Idle-hit vs direct-create | `hitRatio` (warm-latency), `directCreateRatio` (steady-state) | +| Pool health | `client.poolIdleCount` (min/mean/max samples), `poolIdleZeroRatio`, `poolDegradedSamples`, `poolBackoffSamples`, `poolInFlightMax`; failure scenarios also report `poolStateAfterBurst`/`backoffActive`/`failureCount` | +| Replenish throughput | `replenishRatePerSec` / `killRatePerSec` (server-observed), plus the per-second `lifecycle.create`/`lifecycle.delete` QPS series under `mockQps` | +| Pool-size trajectory / over-creation | mock `aliveStats.max` + per-second `alive` series (server view); pool idle should never exceed `maxIdle` + in-flight warmups | +| Client threads | `client.threads` (min/mean/max sampled every 500ms) + `client.threadPeakSinceProbeStart` | +| Client memory/GC | `client.heapUsedMb` (min/mean/max), `client.gcCollections`, `client.gcTimeMs` | +| Server QPS (all APIs) | `perScenarioQps..` — per-route totals, 1s/5s/60s rates, and the full per-second series | + +The `client` block is produced by a probe thread sampling +`SandboxPool.snapshot()` plus JVM thread/heap/GC beans every 500ms during the +scenario. + +## Reusing the mock from other SDKs + +Point any SDK's `ConnectionConfig` at the mock: + +```kotlin +ConnectionConfig.builder() + .domain("127.0.0.1:18080") // lifecycle API (no scheme; driver adds /v1) + .protocol("http") + .build() +``` + +The Go SDK's `pool_test.go` shows the same wiring for Go. Health checks and +endpoint lookups behave like a real server, so the mock doubles as a +deterministic test fixture for SDK e2e-style tests. + +## Adding a scenario + +Implement it in `kotlin/.../benchmark/Scenarios.kt`, returning a +`Map` (nested maps render as sections in Markdown), register it in +`ALL_SCENARIOS`, and add a `--` default in `Cli.kt` when it needs knobs. diff --git a/tests/benchmark/configs/default.json b/tests/benchmark/configs/default.json new file mode 100644 index 000000000..9e5baa450 --- /dev/null +++ b/tests/benchmark/configs/default.json @@ -0,0 +1,42 @@ +{ + "createLatencyMs": { + "distribution": "uniform", + "minMs": 300, + "maxMs": 800 + }, + "createFailureRate": 0.0, + "bootDelayMs": 300, + "execdFailureRate": 0.0, + "defaultTtlSeconds": 3600, + "latencyOverrides": { + "lifecycle.delete": { + "distribution": "uniform", + "minMs": 300, + "maxMs": 800 + }, + "lifecycle.get": { + "distribution": "uniform", + "minMs": 50, + "maxMs": 100 + }, + "lifecycle.renew": { + "distribution": "uniform", + "minMs": 50, + "maxMs": 100 + }, + "lifecycle.endpoint": { + "distribution": "uniform", + "minMs": 50, + "maxMs": 100 + }, + "execd.ping": { + "distribution": "fixed", + "meanMs": 100 + }, + "execd.other": { + "distribution": "uniform", + "minMs": 50, + "maxMs": 100 + } + } +} diff --git a/tests/benchmark/configs/fast.json b/tests/benchmark/configs/fast.json new file mode 100644 index 000000000..2aa4e2328 --- /dev/null +++ b/tests/benchmark/configs/fast.json @@ -0,0 +1,11 @@ +{ + "createLatencyMs": { + "distribution": "fixed", + "meanMs": 100 + }, + "createFailureRate": 0.0, + "bootDelayMs": 50, + "execdFailureRate": 0.0, + "defaultTtlSeconds": 3600, + "latencyOverrides": {} +} diff --git a/tests/benchmark/configs/slow.json b/tests/benchmark/configs/slow.json new file mode 100644 index 000000000..64af9e101 --- /dev/null +++ b/tests/benchmark/configs/slow.json @@ -0,0 +1,13 @@ +{ + "createLatencyMs": { + "distribution": "lognormal", + "meanMs": 2000, + "stddevMs": 1000, + "minMs": 100 + }, + "createFailureRate": 0.0, + "bootDelayMs": 1000, + "execdFailureRate": 0.0, + "defaultTtlSeconds": 3600, + "latencyOverrides": {} +} diff --git a/tests/benchmark/kotlin/build.gradle.kts b/tests/benchmark/kotlin/build.gradle.kts new file mode 100644 index 000000000..f2a0f14d5 --- /dev/null +++ b/tests/benchmark/kotlin/build.gradle.kts @@ -0,0 +1,58 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +plugins { + kotlin("jvm") version "2.2.21" + application +} + +group = "com.alibaba.opensandbox" +version = "1.0.0" + +java { + sourceCompatibility = JavaVersion.VERSION_11 + targetCompatibility = JavaVersion.VERSION_11 +} + +repositories { + mavenCentral() +} + +dependencies { + // OpenSandbox Kotlin SDK, built from source via composite build + // (see settings.gradle.kts). The module coordinate is substituted by the + // included build's :sandbox project; the version is informational only. + implementation("com.alibaba.opensandbox:sandbox:1.0.18") + + implementation("com.squareup.okhttp3:okhttp:4.12.0") + implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.9.0") + implementation("org.slf4j:slf4j-simple:2.0.9") +} + +application { + mainClass.set("com.alibaba.opensandbox.benchmark.MainKt") +} + +tasks.withType { + compilerOptions { + jvmTarget.set(org.jetbrains.kotlin.gradle.dsl.JvmTarget.JVM_11) + } +} + +tasks.withType { + sourceCompatibility = "11" + targetCompatibility = "11" +} diff --git a/tests/benchmark/kotlin/gradle.properties b/tests/benchmark/kotlin/gradle.properties new file mode 100644 index 000000000..e50be0b79 --- /dev/null +++ b/tests/benchmark/kotlin/gradle.properties @@ -0,0 +1,4 @@ +org.gradle.jvmargs=-Xmx2g -XX:MaxMetaspaceSize=512m +org.gradle.parallel=true +org.gradle.caching=true +org.gradle.configuration-cache=true diff --git a/tests/benchmark/kotlin/gradle/wrapper/gradle-wrapper.jar b/tests/benchmark/kotlin/gradle/wrapper/gradle-wrapper.jar new file mode 100644 index 000000000..f8e1ee312 Binary files /dev/null and b/tests/benchmark/kotlin/gradle/wrapper/gradle-wrapper.jar differ diff --git a/tests/benchmark/kotlin/gradle/wrapper/gradle-wrapper.properties b/tests/benchmark/kotlin/gradle/wrapper/gradle-wrapper.properties new file mode 100644 index 000000000..4eac4a84c --- /dev/null +++ b/tests/benchmark/kotlin/gradle/wrapper/gradle-wrapper.properties @@ -0,0 +1,7 @@ +distributionBase=GRADLE_USER_HOME +distributionPath=wrapper/dists +distributionUrl=https\://services.gradle.org/distributions/gradle-9.2.1-all.zip +networkTimeout=10000 +validateDistributionUrl=true +zipStoreBase=GRADLE_USER_HOME +zipStorePath=wrapper/dists diff --git a/tests/benchmark/kotlin/gradlew b/tests/benchmark/kotlin/gradlew new file mode 100755 index 000000000..adff685a0 --- /dev/null +++ b/tests/benchmark/kotlin/gradlew @@ -0,0 +1,248 @@ +#!/bin/sh + +# +# Copyright © 2015 the original authors. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# +# SPDX-License-Identifier: Apache-2.0 +# + +############################################################################## +# +# Gradle start up script for POSIX generated by Gradle. +# +# Important for running: +# +# (1) You need a POSIX-compliant shell to run this script. If your /bin/sh is +# noncompliant, but you have some other compliant shell such as ksh or +# bash, then to run this script, type that shell name before the whole +# command line, like: +# +# ksh Gradle +# +# Busybox and similar reduced shells will NOT work, because this script +# requires all of these POSIX shell features: +# * functions; +# * expansions «$var», «${var}», «${var:-default}», «${var+SET}», +# «${var#prefix}», «${var%suffix}», and «$( cmd )»; +# * compound commands having a testable exit status, especially «case»; +# * various built-in commands including «command», «set», and «ulimit». +# +# Important for patching: +# +# (2) This script targets any POSIX shell, so it avoids extensions provided +# by Bash, Ksh, etc; in particular arrays are avoided. +# +# The "traditional" practice of packing multiple parameters into a +# space-separated string is a well documented source of bugs and security +# problems, so this is (mostly) avoided, by progressively accumulating +# options in "$@", and eventually passing that to Java. +# +# Where the inherited environment variables (DEFAULT_JVM_OPTS, JAVA_OPTS, +# and GRADLE_OPTS) rely on word-splitting, this is performed explicitly; +# see the in-line comments for details. +# +# There are tweaks for specific operating systems such as AIX, CygWin, +# Darwin, MinGW, and NonStop. +# +# (3) This script is generated from the Groovy template +# https://github.com/gradle/gradle/blob/HEAD/platforms/jvm/plugins-application/src/main/resources/org/gradle/api/internal/plugins/unixStartScript.txt +# within the Gradle project. +# +# You can find Gradle at https://github.com/gradle/gradle/. +# +############################################################################## + +# Attempt to set APP_HOME + +# Resolve links: $0 may be a link +app_path=$0 + +# Need this for daisy-chained symlinks. +while + APP_HOME=${app_path%"${app_path##*/}"} # leaves a trailing /; empty if no leading path + [ -h "$app_path" ] +do + ls=$( ls -ld "$app_path" ) + link=${ls#*' -> '} + case $link in #( + /*) app_path=$link ;; #( + *) app_path=$APP_HOME$link ;; + esac +done + +# This is normally unused +# shellcheck disable=SC2034 +APP_BASE_NAME=${0##*/} +# Discard cd standard output in case $CDPATH is set (https://github.com/gradle/gradle/issues/25036) +APP_HOME=$( cd -P "${APP_HOME:-./}" > /dev/null && printf '%s\n' "$PWD" ) || exit + +# Use the maximum available, or set MAX_FD != -1 to use that value. +MAX_FD=maximum + +warn () { + echo "$*" +} >&2 + +die () { + echo + echo "$*" + echo + exit 1 +} >&2 + +# OS specific support (must be 'true' or 'false'). +cygwin=false +msys=false +darwin=false +nonstop=false +case "$( uname )" in #( + CYGWIN* ) cygwin=true ;; #( + Darwin* ) darwin=true ;; #( + MSYS* | MINGW* ) msys=true ;; #( + NONSTOP* ) nonstop=true ;; +esac + + + +# Determine the Java command to use to start the JVM. +if [ -n "$JAVA_HOME" ] ; then + if [ -x "$JAVA_HOME/jre/sh/java" ] ; then + # IBM's JDK on AIX uses strange locations for the executables + JAVACMD=$JAVA_HOME/jre/sh/java + else + JAVACMD=$JAVA_HOME/bin/java + fi + if [ ! -x "$JAVACMD" ] ; then + die "ERROR: JAVA_HOME is set to an invalid directory: $JAVA_HOME + +Please set the JAVA_HOME variable in your environment to match the +location of your Java installation." + fi +else + JAVACMD=java + if ! command -v java >/dev/null 2>&1 + then + die "ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH. + +Please set the JAVA_HOME variable in your environment to match the +location of your Java installation." + fi +fi + +# Increase the maximum file descriptors if we can. +if ! "$cygwin" && ! "$darwin" && ! "$nonstop" ; then + case $MAX_FD in #( + max*) + # In POSIX sh, ulimit -H is undefined. That's why the result is checked to see if it worked. + # shellcheck disable=SC2039,SC3045 + MAX_FD=$( ulimit -H -n ) || + warn "Could not query maximum file descriptor limit" + esac + case $MAX_FD in #( + '' | soft) :;; #( + *) + # In POSIX sh, ulimit -n is undefined. That's why the result is checked to see if it worked. + # shellcheck disable=SC2039,SC3045 + ulimit -n "$MAX_FD" || + warn "Could not set maximum file descriptor limit to $MAX_FD" + esac +fi + +# Collect all arguments for the java command, stacking in reverse order: +# * args from the command line +# * the main class name +# * -classpath +# * -D...appname settings +# * --module-path (only if needed) +# * DEFAULT_JVM_OPTS, JAVA_OPTS, and GRADLE_OPTS environment variables. + +# For Cygwin or MSYS, switch paths to Windows format before running java +if "$cygwin" || "$msys" ; then + APP_HOME=$( cygpath --path --mixed "$APP_HOME" ) + + JAVACMD=$( cygpath --unix "$JAVACMD" ) + + # Now convert the arguments - kludge to limit ourselves to /bin/sh + for arg do + if + case $arg in #( + -*) false ;; # don't mess with options #( + /?*) t=${arg#/} t=/${t%%/*} # looks like a POSIX filepath + [ -e "$t" ] ;; #( + *) false ;; + esac + then + arg=$( cygpath --path --ignore --mixed "$arg" ) + fi + # Roll the args list around exactly as many times as the number of + # args, so each arg winds up back in the position where it started, but + # possibly modified. + # + # NB: a `for` loop captures its iteration list before it begins, so + # changing the positional parameters here affects neither the number of + # iterations, nor the values presented in `arg`. + shift # remove old arg + set -- "$@" "$arg" # push replacement arg + done +fi + + +# Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script. +DEFAULT_JVM_OPTS='"-Xmx64m" "-Xms64m"' + +# Collect all arguments for the java command: +# * DEFAULT_JVM_OPTS, JAVA_OPTS, and optsEnvironmentVar are not allowed to contain shell fragments, +# and any embedded shellness will be escaped. +# * For example: A user cannot expect ${Hostname} to be expanded, as it is an environment variable and will be +# treated as '${Hostname}' itself on the command line. + +set -- \ + "-Dorg.gradle.appname=$APP_BASE_NAME" \ + -jar "$APP_HOME/gradle/wrapper/gradle-wrapper.jar" \ + "$@" + +# Stop when "xargs" is not available. +if ! command -v xargs >/dev/null 2>&1 +then + die "xargs is not available" +fi + +# Use "xargs" to parse quoted args. +# +# With -n1 it outputs one arg per line, with the quotes and backslashes removed. +# +# In Bash we could simply go: +# +# readarray ARGS < <( xargs -n1 <<<"$var" ) && +# set -- "${ARGS[@]}" "$@" +# +# but POSIX shell has neither arrays nor command substitution, so instead we +# post-process each arg (as a line of input to sed) to backslash-escape any +# character that might be a shell metacharacter, then use eval to reverse +# that process (while maintaining the separation between arguments), and wrap +# the whole thing up as a single "set" statement. +# +# This will of course break if any of these variables contains a newline or +# an unmatched quote. +# + +eval "set -- $( + printf '%s\n' "$DEFAULT_JVM_OPTS $JAVA_OPTS $GRADLE_OPTS" | + xargs -n1 | + sed ' s~[^-[:alnum:]+,./:=@_]~\\&~g; ' | + tr '\n' ' ' + )" '"$@"' + +exec "$JAVACMD" "$@" diff --git a/tests/benchmark/kotlin/settings.gradle.kts b/tests/benchmark/kotlin/settings.gradle.kts new file mode 100644 index 000000000..137bf298e --- /dev/null +++ b/tests/benchmark/kotlin/settings.gradle.kts @@ -0,0 +1,21 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +rootProject.name = "opensandbox-pool-benchmark" + +// Build the Kotlin SDK from source (Gradle composite build) so the benchmark +// always runs the checked-out SDK code; no mavenLocal publish step needed. +includeBuild("../../../sdks/sandbox/kotlin") diff --git a/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Cli.kt b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Cli.kt new file mode 100644 index 000000000..ceb8ffade --- /dev/null +++ b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Cli.kt @@ -0,0 +1,181 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.alibaba.opensandbox.benchmark + +import java.time.Duration + +/** + * All benchmark knobs, parsed from `--key value` command-line arguments. + */ +data class BenchmarkConfig( + val mockBaseUrl: String, + val reportDir: String, + val scenarios: List, + val maxIdle: Int, + val warmupConcurrency: Int, + val reconcileIntervalMs: Long, + val idleTimeoutS: Long, + val acquireMinRemainingTtlS: Long, + val primaryLockTtlS: Long, + val degradedThreshold: Int, + val acquireReadyTimeoutMs: Long, + val warmupReadyTimeoutMs: Long, + val healthCheckPollingIntervalMs: Long, + val coldStartTimeoutMs: Long, + val warmWorkers: Int, + val warmRoundsPerWorker: Int, + val steadyWorkers: Int, + val steadyDurationS: Int, + val acquireRatePerMin: Int, + val steadyStartImmediately: Boolean, + val holdMinMs: Long, + val holdMaxMs: Long, + val replenishRounds: Int, + val replenishWaitTimeoutMs: Long, + val failureCreateRate: Double, + val failureAcquires: Int, + val staleAcquires: Int, + val staleRetries: Int, + val staleAcquireReadyTimeoutMs: Long, + val stalePoisonRate: Double, + val sharedConnectionPoolSize: Int, + val idleExpiryIdleTimeoutS: Long, + val idleExpiryDurationS: Int, +) { + val mockDomain: String + get() = mockBaseUrl.removePrefix("http://").removePrefix("https://").trimEnd('/') +} + +object Cli { + private val allKeys = + listOf( + "mock-base-url", + "report-dir", + "scenarios", + "max-idle", + "warmup-concurrency", + "reconcile-interval-ms", + "idle-timeout-s", + "acquire-min-remaining-ttl-s", + "primary-lock-ttl-s", + "degraded-threshold", + "acquire-ready-timeout-ms", + "warmup-ready-timeout-ms", + "health-check-polling-interval-ms", + "cold-start-timeout-ms", + "warm-workers", + "warm-rounds-per-worker", + "steady-workers", + "steady-duration-s", + "acquire-rate-per-min", + "steady-start-immediately", + "hold-min-ms", + "hold-max-ms", + "replenish-rounds", + "replenish-wait-timeout-ms", + "failure-create-rate", + "failure-acquires", + "stale-acquires", + "stale-retries", + "stale-acquire-ready-timeout-ms", + "stale-poison-rate", + "shared-connection-pool-size", + "idle-expiry-idle-timeout-s", + "idle-expiry-duration-s", + ) + + fun parse(args: Array): BenchmarkConfig { + val map = mutableMapOf() + var i = 0 + while (i < args.size) { + val key = args[i] + if (!key.startsWith("--")) { + throw IllegalArgumentException("unexpected argument: $key") + } + val name = key.removePrefix("--") + if (name !in allKeys) { + throw IllegalArgumentException("unknown option: $key") + } + // Boolean flags may be passed without a value (--flag == --flag true). + val next = args.getOrNull(i + 1) + if (next == null || next.startsWith("--")) { + map[name] = "true" + i += 1 + } else { + map[name] = next + i += 2 + } + } + return BenchmarkConfig( + mockBaseUrl = map["mock-base-url"] ?: "http://127.0.0.1:18080", + reportDir = map["report-dir"] ?: "results/run-${System.currentTimeMillis()}", + scenarios = + (map["scenarios"] ?: "all").split(",").map { it.trim() }.filter { it.isNotEmpty() }, + maxIdle = (map["max-idle"] ?: "20").toInt(), + warmupConcurrency = (map["warmup-concurrency"] ?: "4").toInt(), + reconcileIntervalMs = (map["reconcile-interval-ms"] ?: "1000").toLong(), + idleTimeoutS = (map["idle-timeout-s"] ?: "1800").toLong(), + // 0 = leave the SDK's auto-derived default (min(60s, idleTimeout/2)) + acquireMinRemainingTtlS = (map["acquire-min-remaining-ttl-s"] ?: "0").toLong(), + // 0 = leave the SDK default (60s) + primaryLockTtlS = (map["primary-lock-ttl-s"] ?: "0").toLong(), + // 0 = leave the SDK default (3) + degradedThreshold = (map["degraded-threshold"] ?: "0").toInt(), + acquireReadyTimeoutMs = (map["acquire-ready-timeout-ms"] ?: "15000").toLong(), + warmupReadyTimeoutMs = (map["warmup-ready-timeout-ms"] ?: "15000").toLong(), + healthCheckPollingIntervalMs = (map["health-check-polling-interval-ms"] ?: "200").toLong(), + coldStartTimeoutMs = (map["cold-start-timeout-ms"] ?: "120000").toLong(), + warmWorkers = (map["warm-workers"] ?: "16").toInt(), + warmRoundsPerWorker = (map["warm-rounds-per-worker"] ?: "150").toInt(), + steadyWorkers = (map["steady-workers"] ?: "16").toInt(), + steadyDurationS = (map["steady-duration-s"] ?: "60").toInt(), + // 0 = unlimited (workers run back-to-back); > 0 paces acquires + // evenly across each minute at this many acquires per minute. + acquireRatePerMin = (map["acquire-rate-per-min"] ?: "0").toInt(), + // Start loaders immediately after pool.start() instead of waiting + // for the idle buffer to fill (pool startup races high-frequency + // acquire — extreme cold-start-under-load scenario). + steadyStartImmediately = (map["steady-start-immediately"] ?: "false").toBoolean(), + holdMinMs = (map["hold-min-ms"] ?: "1000").toLong(), + holdMaxMs = (map["hold-max-ms"] ?: "5000").toLong(), + replenishRounds = (map["replenish-rounds"] ?: "20").toInt(), + replenishWaitTimeoutMs = (map["replenish-wait-timeout-ms"] ?: "15000").toLong(), + failureCreateRate = (map["failure-create-rate"] ?: "0.6").toDouble(), + failureAcquires = (map["failure-acquires"] ?: "60").toInt(), + staleAcquires = (map["stale-acquires"] ?: "100").toInt(), + staleRetries = (map["stale-retries"] ?: "3").toInt(), + staleAcquireReadyTimeoutMs = (map["stale-acquire-ready-timeout-ms"] ?: "3000").toLong(), + // fraction (0..1] of idle sandboxes to poison in stale-idle; 1.0 = poison all + stalePoisonRate = (map["stale-poison-rate"] ?: "1.0").toDouble(), + // Inject a shared OkHttp ConnectionPool across all sandbox + // clients (fixed default 500 idle slots; 0 = each sandbox keeps + // its own fresh connections, reproducing the connection-reset + // pathology at high concurrency; N = that many idle slots). + sharedConnectionPoolSize = (map["shared-connection-pool-size"] ?: "500").toInt(), + idleExpiryIdleTimeoutS = (map["idle-expiry-idle-timeout-s"] ?: "20").toLong(), + idleExpiryDurationS = (map["idle-expiry-duration-s"] ?: "40").toInt(), + ) + } + + fun usage(): String = + buildString { + appendLine("Usage: pool-benchmark [--key value ...]") + allKeys.forEach { appendLine(" --$it ") } + } +} + +val BenchmarkConfig.acquireTimeout: Duration get() = Duration.ofMinutes(10) diff --git a/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/FailingPoolStateStore.kt b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/FailingPoolStateStore.kt new file mode 100644 index 000000000..a77969909 --- /dev/null +++ b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/FailingPoolStateStore.kt @@ -0,0 +1,132 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.alibaba.opensandbox.benchmark + +import com.alibaba.opensandbox.sandbox.domain.exceptions.PoolStateStoreUnavailableException +import com.alibaba.opensandbox.sandbox.domain.pool.IdleEntry +import com.alibaba.opensandbox.sandbox.domain.pool.PoolDestroyState +import com.alibaba.opensandbox.sandbox.domain.pool.PoolStateStore +import com.alibaba.opensandbox.sandbox.domain.pool.StoreCounters +import com.alibaba.opensandbox.sandbox.domain.pool.TakeIdleResult +import com.alibaba.opensandbox.sandbox.infrastructure.pool.InMemoryPoolStateStore +import java.time.Duration +import java.time.Instant +import java.util.concurrent.atomic.AtomicLong + +/** + * A [PoolStateStore] wrapper that can be toggled into "outage" mode, where + * every operation throws [PoolStateStoreUnavailableException]. Used to test + * the pool's store-outage behavior (OSEP-0005): DIRECT_CREATE fallthrough + * keeps acquire available; FAIL_FAST fails closed. + */ +class FailingPoolStateStore( + private val delegate: PoolStateStore = InMemoryPoolStateStore(), +) : PoolStateStore { + @Volatile + private var failing = false + private val errorCount = AtomicLong() + + fun setFailing(f: Boolean) { + failing = f + } + + fun errorCount(): Long = errorCount.get() + + private fun gate(block: () -> T): T { + if (failing) { + errorCount.incrementAndGet() + throw PoolStateStoreUnavailableException("simulated store outage") + } + return block() + } + + override fun tryTakeIdle(poolName: String): String? = gate { delegate.tryTakeIdle(poolName) } + + override fun tryTakeIdle( + poolName: String, + minRemainingTtl: Duration, + ): TakeIdleResult = gate { delegate.tryTakeIdle(poolName, minRemainingTtl) } + + override fun putIdle( + poolName: String, + sandboxId: String, + ) = gate { delegate.putIdle(poolName, sandboxId) } + + override fun removeIdle( + poolName: String, + sandboxId: String, + ) = gate { delegate.removeIdle(poolName, sandboxId) } + + override fun tryAcquirePrimaryLock( + poolName: String, + ownerId: String, + ttl: Duration, + ): Boolean = gate { delegate.tryAcquirePrimaryLock(poolName, ownerId, ttl) } + + override fun renewPrimaryLock( + poolName: String, + ownerId: String, + ttl: Duration, + ): Boolean = gate { delegate.renewPrimaryLock(poolName, ownerId, ttl) } + + override fun releasePrimaryLock( + poolName: String, + ownerId: String, + ) = gate { delegate.releasePrimaryLock(poolName, ownerId) } + + override fun reapExpiredIdle( + poolName: String, + now: Instant, + ) = gate { delegate.reapExpiredIdle(poolName, now) } + + override fun reapExpiredIdle( + poolName: String, + now: Instant, + minRemainingTtl: Duration, + ): List = gate { delegate.reapExpiredIdle(poolName, now, minRemainingTtl) } + + override fun snapshotCounters(poolName: String): StoreCounters = gate { delegate.snapshotCounters(poolName) } + + override fun snapshotIdleEntries(poolName: String): List = gate { delegate.snapshotIdleEntries(poolName) } + + override fun getMaxIdle(poolName: String): Int? = gate { delegate.getMaxIdle(poolName) } + + override fun setMaxIdle( + poolName: String, + maxIdle: Int, + ) = gate { delegate.setMaxIdle(poolName, maxIdle) } + + override fun setIdleEntryTtl( + poolName: String, + idleTtl: Duration, + ) = gate { delegate.setIdleEntryTtl(poolName, idleTtl) } + + override fun getDestroyState(poolName: String): PoolDestroyState = gate { delegate.getDestroyState(poolName) } + + override fun beginDestroy( + poolName: String, + ownerId: String, + ) = gate { delegate.beginDestroy(poolName, ownerId) } + + override fun clearPoolState(poolName: String) = gate { delegate.clearPoolState(poolName) } + + override fun markDestroyed( + poolName: String, + ownerId: String, + tombstoneTtl: Duration?, + ) = gate { delegate.markDestroyed(poolName, ownerId, tombstoneTtl) } +} diff --git a/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Main.kt b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Main.kt new file mode 100644 index 000000000..5f0c38f32 --- /dev/null +++ b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Main.kt @@ -0,0 +1,277 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.alibaba.opensandbox.benchmark + +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonArray +import kotlinx.serialization.json.JsonElement +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.JsonPrimitive +import kotlinx.serialization.json.buildJsonObject +import kotlinx.serialization.json.put +import java.io.File +import java.time.Duration +import java.time.Instant + +fun main(args: Array) { + val cfg = + try { + Cli.parse(args) + } catch (e: Exception) { + System.err.println("error: ${e.message}") + System.err.println(Cli.usage()) + kotlin.system.exitProcess(2) + } + + val mock = MockControl(cfg.mockBaseUrl, Duration.ofSeconds(10)) + if (!mock.ping()) { + System.err.println( + "error: mock server not reachable at ${cfg.mockBaseUrl} " + + "(start it first, e.g. via tests/benchmark/run.sh)", + ) + kotlin.system.exitProcess(1) + } + + val scenarios = mutableListOf() + for (name in cfg.scenarios) { + if (name == "all") { + scenarios.addAll(Scenarios.ALL_SCENARIOS.keys) + } else { + if (name !in Scenarios.ALL_SCENARIOS) { + System.err.println("error: unknown scenario '$name' (available: ${Scenarios.ALL_SCENARIOS.keys})") + kotlin.system.exitProcess(2) + } + scenarios.add(name) + } + } + + println("== OpenSandbox pool benchmark ==") + println("mock: ${cfg.mockBaseUrl} scenarios: $scenarios") + println( + "config: maxIdle=${cfg.maxIdle} warmupConcurrency=${cfg.warmupConcurrency} " + + "reconcileIntervalMs=${cfg.reconcileIntervalMs} idleTimeoutS=${cfg.idleTimeoutS}", + ) + + val results = LinkedHashMap() + val perScenarioQps = LinkedHashMap() + for (name in scenarios) { + println("\n-- scenario: $name --") + val t0 = System.nanoTime() + val section = + try { + Scenarios.ALL_SCENARIOS.getValue(name)(cfg, mock) + } catch (t: Throwable) { + System.err.println("scenario $name failed: $t") + mapOf("error" to (t.message ?: t.toString())) + } + val elapsedMs = (System.nanoTime() - t0) / 1_000_000 + println(" completed in ${elapsedMs}ms") + section.forEach { (k, v) -> println(" $k=$v") } + results[name] = section + // Precise per-API QPS observed by the mock during this scenario. The + // mock ring keeps per-second counts; the series is included in the JSON + // report and exported as a per-scenario CSV for offline analysis. + val stats = mock.stats() + perScenarioQps[name] = stats["qps"] + writeTimeseriesCsv(File(cfg.reportDir, "server-timeseries-$name.csv"), name, stats) + } + results["perScenarioQps"] = perScenarioQps + + val endStats = mock.stats() + results["mockServerStats"] = endStats + + val report = + buildJsonObject { + put("runId", "${Instant.now().toEpochMilli()}") + put("timestamp", Instant.now().toString()) + put("config", toJsonElement(cfgValues(cfg))) + put("results", toJsonElement(results)) + } + + val outDir = File(cfg.reportDir) + outDir.mkdirs() + File(outDir, "mock-stats-end.json") + .writeText(Json { prettyPrint = true }.encodeToString(JsonObject.serializer(), endStats as JsonObject)) + File(outDir, "report.json").writeText(Json { prettyPrint = true }.encodeToString(JsonObject.serializer(), report)) + File(outDir, "report.md").writeText(renderMarkdown(cfg, results)) + println("\n== done: report written to ${outDir.absolutePath}/report.{json,md} ==") +} + +private fun cfgValues(cfg: BenchmarkConfig): Map = + linkedMapOf( + "mockBaseUrl" to cfg.mockBaseUrl, + "scenarios" to cfg.scenarios, + "maxIdle" to cfg.maxIdle, + "warmupConcurrency" to cfg.warmupConcurrency, + "reconcileIntervalMs" to cfg.reconcileIntervalMs, + "idleTimeoutS" to cfg.idleTimeoutS, + "acquireMinRemainingTtlS" to cfg.acquireMinRemainingTtlS, + "primaryLockTtlS" to cfg.primaryLockTtlS, + "degradedThreshold" to cfg.degradedThreshold, + "acquireReadyTimeoutMs" to cfg.acquireReadyTimeoutMs, + "warmupReadyTimeoutMs" to cfg.warmupReadyTimeoutMs, + "healthCheckPollingIntervalMs" to cfg.healthCheckPollingIntervalMs, + "warmWorkers" to cfg.warmWorkers, + "warmRoundsPerWorker" to cfg.warmRoundsPerWorker, + "steadyWorkers" to cfg.steadyWorkers, + "steadyDurationS" to cfg.steadyDurationS, + "acquireRatePerMin" to cfg.acquireRatePerMin, + "steadyStartImmediately" to cfg.steadyStartImmediately, + "holdMinMs" to cfg.holdMinMs, + "holdMaxMs" to cfg.holdMaxMs, + "failureCreateRate" to cfg.failureCreateRate, + "staleRetries" to cfg.staleRetries, + "stalePoisonRate" to cfg.stalePoisonRate, + "sharedConnectionPoolSize" to cfg.sharedConnectionPoolSize, + ) + +private fun toJsonElement(value: Any?): JsonElement = + when (value) { + is JsonElement -> value + is Map<*, *> -> JsonObject(value.entries.associate { (k, v) -> k.toString() to toJsonElement(v) }) + is Iterable<*> -> JsonArray(value.map { toJsonElement(it) }) + is Double -> JsonPrimitive(value) + is Float -> JsonPrimitive(value) + is Long -> JsonPrimitive(value) + is Int -> JsonPrimitive(value) + is Boolean -> JsonPrimitive(value) + is String -> JsonPrimitive(value) + is Number -> JsonPrimitive(value.toDouble()) + null -> JsonPrimitive("") + else -> JsonPrimitive(value.toString()) + } + +private fun renderMarkdown(cfg: BenchmarkConfig, results: Map): String { + val sb = StringBuilder() + sb.appendLine("# OpenSandbox Pool Benchmark") + sb.appendLine() + sb.appendLine("- runId: ${Instant.now().toEpochMilli()}") + sb.appendLine("- mock: ${cfg.mockBaseUrl}") + sb.appendLine( + "- maxIdle: ${cfg.maxIdle}, warmupConcurrency: ${cfg.warmupConcurrency}, " + + "reconcileIntervalMs: ${cfg.reconcileIntervalMs}, idleTimeoutS: ${cfg.idleTimeoutS}", + ) + sb.appendLine() + for ((scenario, section) in results) { + if (scenario == "mockServerStats") continue + sb.appendLine("## $scenario") + sb.appendLine() + if (section is Map<*, *>) { + renderMap(sb, section as Map, " ") + } + sb.appendLine() + } + sb.appendLine("## mock server stats (end of run)") + sb.appendLine() + if (results["mockServerStats"] is Map<*, *>) { + renderMap(sb, results["mockServerStats"] as Map, " ") + } + return sb.toString() +} + +private fun renderMap( + sb: StringBuilder, + map: Map<*, *>, + indent: String, +) { + for ((k, v) in map) { + if (k == "series") continue // full per-second series lives in report.json only + when (v) { + is Map<*, *> -> { + sb.appendLine("$indent$k:") + renderMap(sb, v, "$indent ") + } + is List<*> -> sb.appendLine("$indent$k: ${v.joinToString(",")}") + else -> sb.appendLine("$indent$k: $v") + } + } +} + +/** + * Merges the mock's per-second series (per-API QPS + alive gauge) for one + * scenario into a single CSV: `second,alive,create,delete,get,renew, + * endpoint,execd.ping,execd.other` (second = offset from the earliest second + * recorded in the scenario window). + */ +private data class Series( + val start: Long, + val values: List, +) + +private fun writeTimeseriesCsv( + file: File, + scenario: String, + stats: Map, +) { + fun extractSeries( + key: String, + value: Any?, + ): Series? { + if (value !is Map<*, *>) return null + val start = (value["seriesStartUnixSec"] as? JsonPrimitive)?.content?.toLongOrNull() ?: return null + val raw = value["series"] as? List<*> ?: return null + val values = raw.mapNotNull { (it as? JsonPrimitive)?.content?.toLongOrNull() } + return Series(start, values) + } + + val routes = + listOf( + "lifecycle.create", + "lifecycle.delete", + "lifecycle.get", + "lifecycle.renew", + "lifecycle.endpoint", + "execd.ping", + "execd.other", + ) + val qps = stats["qps"] as? Map<*, *> ?: return + val aliveStats = stats["aliveStats"] + val seriesByRoute = routes.associateWith { route -> extractSeries(route, qps[route]) } + val alive = extractSeries("alive", aliveStats) + + val allStarts = seriesByRoute.values.mapNotNull { it?.start } + (alive?.start ?: 0L) + if (allStarts.isEmpty()) return + val begin = allStarts.min() + val end = allStarts.maxOf { start -> + val s = seriesByRoute.values.firstOrNull { it?.start == start } ?: alive + start + (s?.values?.size?.toLong() ?: 1L) - 1 + } + + val sb = StringBuilder() + sb.appendLine("second,alive,create,delete,get,renew,endpoint,execd.ping,execd.other") + for (sec in begin..end) { + val offset = sec - begin + sb.append(offset) + sb.append(',').append(valueAt(alive, sec)) + for (route in routes) { + sb.append(',').append(valueAt(seriesByRoute[route], sec)) + } + sb.appendLine() + } + file.parentFile?.mkdirs() + file.writeText(sb.toString()) +} + +private fun valueAt( + series: Series?, + second: Long, +): Long { + if (series == null) return 0L + val idx = (second - series.start).toInt() + if (idx < 0 || idx >= series.values.size) return 0L + return series.values[idx] +} diff --git a/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Metrics.kt b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Metrics.kt new file mode 100644 index 000000000..3b63d4ad0 --- /dev/null +++ b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Metrics.kt @@ -0,0 +1,97 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.alibaba.opensandbox.benchmark + +/** + * Thread-safe latency collector. Contention is negligible at the benchmark's + * acquire rates, so a lock-protected ArrayList is sufficient. + */ +class LatencyCollector { + private val lock = Any() + private val samples = ArrayList() + private val failureCounts = LinkedHashMap() + + fun record(elapsedMs: Long) { + synchronized(lock) { + samples.add(elapsedMs) + } + } + + fun recordFailure(reason: String = "other") { + synchronized(lock) { + failureCounts[reason] = (failureCounts[reason] ?: 0L) + 1 + } + } + + fun snapshot(): LatencyStats { + val sorted: LongArray + val failures: Long + val failureByType: Map + synchronized(lock) { + sorted = LongArray(samples.size) + for ((i, v) in samples.withIndex()) sorted[i] = v + sorted.sort() + failures = failureCounts.values.sum() + failureByType = failureCounts.toMap() + } + return LatencyStats( + n = sorted.size.toLong(), + failures = failures, + failuresByType = failureByType, + meanMs = if (sorted.isEmpty()) 0.0 else sorted.average(), + p50 = percentile(sorted, 0.50), + p90 = percentile(sorted, 0.90), + p95 = percentile(sorted, 0.95), + p99 = percentile(sorted, 0.99), + p999 = percentile(sorted, 0.999), + maxMs = sorted.lastOrNull() ?: 0L, + ) + } + + private fun percentile(sorted: LongArray, p: Double): Long { + if (sorted.isEmpty()) return 0L + val idx = ((sorted.size - 1) * p).toInt() + return sorted[idx] + } +} + +data class LatencyStats( + val n: Long, + val failures: Long, + val failuresByType: Map = emptyMap(), + val meanMs: Double, + val p50: Long, + val p90: Long, + val p95: Long, + val p99: Long, + val p999: Long, + val maxMs: Long, +) { + fun toMap(): Map = + mapOf( + "count" to n, + "failures" to failures, + "failuresByType" to failuresByType, + "meanMs" to meanMs, + "p50Ms" to p50, + "p90Ms" to p90, + "p95Ms" to p95, + "p99Ms" to p99, + "p999Ms" to p999, + "maxMs" to maxMs, + ) +} diff --git a/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/MockControl.kt b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/MockControl.kt new file mode 100644 index 000000000..1b6da94ec --- /dev/null +++ b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/MockControl.kt @@ -0,0 +1,102 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.alibaba.opensandbox.benchmark + +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.buildJsonObject +import kotlinx.serialization.json.jsonObject +import kotlinx.serialization.json.put +import okhttp3.MediaType.Companion.toMediaType +import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.RequestBody.Companion.toRequestBody +import java.time.Duration + +/** + * Client for the mock server's control endpoints: + * `GET /__stats`, `POST /__config`, `POST /__reset`. + */ +class MockControl( + baseUrl: String, + requestTimeout: Duration, +) { + private val client = + OkHttpClient.Builder() + .connectTimeout(requestTimeout) + .readTimeout(requestTimeout) + .build() + private val base = baseUrl.trimEnd('/') + + private val json = Json { ignoreUnknownKeys = true } + + fun ping(): Boolean = + try { + client.newCall(Request.Builder().url("$base/__stats").get().build()).execute().use { it.isSuccessful } + } catch (e: Exception) { + false + } + + fun stats(): Map = get("__stats") + + fun reset() { + post("__reset", buildJsonObject {}) + } + + fun setFaults( + createFailureRate: Double? = null, + execdFailureRate: Double? = null, + poisonExisting: Boolean = false, + poisonRate: Double? = null, + ) { + val body = + buildJsonObject { + createFailureRate?.let { put("createFailureRate", it) } + execdFailureRate?.let { put("execdFailureRate", it) } + if (poisonExisting) put("poisonExisting", true) + poisonRate?.let { put("poisonRate", it) } + } + post("__config", body) + } + + private fun get(path: String): Map { + val response = client.newCall(Request.Builder().url("$base/$path").get().build()).execute() + response.use { + if (!it.isSuccessful) { + throw IllegalStateException("mock $path failed: HTTP ${it.code}") + } + val body = it.body?.string() ?: "{}" + return json.parseToJsonElement(body).jsonObject + } + } + + private fun post( + path: String, + body: JsonObject, + ) { + val request = + Request.Builder() + .url("$base/$path") + .post(body.toString().toRequestBody("application/json".toMediaType())) + .build() + client.newCall(request).execute().use { + if (!it.isSuccessful) { + throw IllegalStateException("mock $path failed: HTTP ${it.code} ${it.body?.string()}") + } + } + } +} diff --git a/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/PoolProbe.kt b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/PoolProbe.kt new file mode 100644 index 000000000..922455ec8 --- /dev/null +++ b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/PoolProbe.kt @@ -0,0 +1,194 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.alibaba.opensandbox.benchmark + +import com.alibaba.opensandbox.sandbox.domain.pool.PoolState +import com.alibaba.opensandbox.sandbox.pool.SandboxPool +import java.io.File +import java.lang.management.GarbageCollectorMXBean +import java.lang.management.ManagementFactory +import java.util.concurrent.atomic.AtomicBoolean + +/** + * Continuous client-side instrumentation: JVM threads, heap/GC, and pool + * health (snapshot-based) sampled every [intervalMs] during a scenario. + * Aggregates go into the report; the raw per-sample series are written to a + * CSV in the run directory for offline analysis. + */ +class PoolProbe( + private val pool: SandboxPool, + private val intervalMs: Long = 500, +) { + private val running = AtomicBoolean(true) + private val lock = Any() + private val sampleTimesMs = ArrayList() + private val threadCounts = ArrayList() + private val heapUsedMb = ArrayList() + private val idleSamples = ArrayList() + private val inFlightSamples = ArrayList() + private val degradedSamples = ArrayList() + private val backoffSamples = ArrayList() + private var thread: Thread? = null + + private val threadBean = ManagementFactory.getThreadMXBean() + private val memoryBean = ManagementFactory.getMemoryMXBean() + private val gcBeans: List = ManagementFactory.getGarbageCollectorMXBeans() + private val gcStartCount = gcBeans.sumOf { it.collectionCount } + private val gcStartTimeMs = gcBeans.sumOf { it.collectionTime } + + init { + // Peak thread count is reported relative to this probe's window. + threadBean.resetPeakThreadCount() + } + + fun start() { + val t = + Thread { + while (running.get()) { + sample() + try { + Thread.sleep(intervalMs) + } catch (_: InterruptedException) { + Thread.currentThread().interrupt() + break + } + } + } + t.isDaemon = true + t.name = "bench-probe" + thread = t + t.start() + } + + fun stop() { + running.set(false) + thread?.join(5000) + } + + private fun sample() { + val heap = memoryBean.heapMemoryUsage + val snap = + try { + pool.snapshot() + } catch (_: Exception) { + // Pool snapshot may fail while the state store is faulted. + return + } + val nowMs = System.currentTimeMillis() + synchronized(lock) { + sampleTimesMs.add(nowMs) + threadCounts.add(threadBean.threadCount) + heapUsedMb.add(heap.used / (1024.0 * 1024.0)) + idleSamples.add(snap.idleCount) + inFlightSamples.add(snap.inFlightOperations) + degradedSamples.add(snap.state == PoolState.DEGRADED) + backoffSamples.add(snap.backoffActive) + } + } + + fun report(): Map { + val threads: LongArray + val heap: DoubleArray + val idle: IntArray + val inFlight: IntArray + val degradedCount: Int + val backoffCount: Int + synchronized(lock) { + threads = LongArray(threadCounts.size) { threadCounts[it].toLong() } + heap = DoubleArray(heapUsedMb.size) { heapUsedMb[it] } + idle = IntArray(idleSamples.size) { idleSamples[it] } + inFlight = IntArray(inFlightSamples.size) { inFlightSamples[it] } + degradedCount = degradedSamples.count { it } + backoffCount = backoffSamples.count { it } + } + val zeroIdleRatio = + if (idle.isEmpty()) 0.0 else idle.count { it == 0 }.toDouble() / idle.size + return mapOf( + "threads" to stat(threads), + "threadPeakSinceProbeStart" to threadBean.peakThreadCount, + "heapUsedMb" to stat(heap), + "gcCollections" to (gcBeans.sumOf { it.collectionCount } - gcStartCount), + "gcTimeMs" to (gcBeans.sumOf { it.collectionTime } - gcStartTimeMs), + "poolIdleCount" to stat(idle.map { it.toLong() }.toLongArray()), + "poolIdleZeroRatio" to zeroIdleRatio, + "poolInFlight" to stat(inFlight.map { it.toLong() }.toLongArray()), + "poolDegradedSamples" to degradedCount, + "poolBackoffSamples" to backoffCount, + ) + } + + /** + * Writes the raw per-sample series as CSV (time offset ms, threads, + * heap MB, idle, inFlight, degraded, backoff) for offline analysis. + */ + fun writeCsv(file: File) { + val times: LongArray + val threads: IntArray + val heap: DoubleArray + val idle: IntArray + val inFlight: IntArray + val degraded: BooleanArray + val backoff: BooleanArray + synchronized(lock) { + times = sampleTimesMs.toLongArray() + threads = threadCounts.toIntArray() + heap = heapUsedMb.toDoubleArray() + idle = idleSamples.toIntArray() + inFlight = inFlightSamples.toIntArray() + degraded = degradedSamples.toBooleanArray() + backoff = backoffSamples.toBooleanArray() + } + if (times.isEmpty()) return + val start = times[0] + file.parentFile?.mkdirs() + file.writeText( + buildString { + appendLine("tMs,threads,heapUsedMb,idleCount,inFlight,degraded,backoff") + for (i in times.indices) { + append(times[i] - start) + append(',').append(threads[i]) + append(',').append(heap[i]) + append(',').append(idle[i]) + append(',').append(inFlight[i]) + append(',').append(if (degraded[i]) 1 else 0) + append(',').append(if (backoff[i]) 1 else 0) + appendLine() + } + }, + ) + } + + private fun stat(samples: LongArray): Map { + if (samples.isEmpty()) return mapOf("samples" to 0, "min" to 0L, "mean" to 0.0, "max" to 0L) + return mapOf( + "samples" to samples.size, + "min" to samples.min(), + "mean" to samples.average(), + "max" to samples.max(), + ) + } + + private fun stat(samples: DoubleArray): Map { + if (samples.isEmpty()) return mapOf("samples" to 0, "min" to 0.0, "mean" to 0.0, "max" to 0.0) + return mapOf( + "samples" to samples.size, + "min" to samples.min(), + "mean" to samples.average(), + "max" to samples.max(), + ) + } +} diff --git a/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/PoolRunner.kt b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/PoolRunner.kt new file mode 100644 index 000000000..f14973ed3 --- /dev/null +++ b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/PoolRunner.kt @@ -0,0 +1,142 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.alibaba.opensandbox.benchmark + +import com.alibaba.opensandbox.sandbox.pool.SandboxPool +import com.alibaba.opensandbox.sandbox.config.ConnectionConfig +import com.alibaba.opensandbox.sandbox.domain.pool.AcquirePolicy +import com.alibaba.opensandbox.sandbox.domain.pool.PoolCreationSpec +import com.alibaba.opensandbox.sandbox.domain.pool.PoolStateStore +import com.alibaba.opensandbox.sandbox.infrastructure.pool.InMemoryPoolStateStore +import java.time.Duration +import java.util.concurrent.TimeUnit + +/** + * Builds a [SandboxPool] wired to the mock server. Health checks stay enabled + * so the execd path (connect + readiness ping) is exercised. + */ +object PoolRunner { + fun build( + cfg: BenchmarkConfig, + poolName: String, + maxIdle: Int = cfg.maxIdle, + warmupConcurrency: Int = cfg.warmupConcurrency, + reconcileIntervalMs: Long = cfg.reconcileIntervalMs, + idleTimeoutS: Long = cfg.idleTimeoutS, + maxAcquireRetries: Int = cfg.staleRetries, + acquireReadyTimeoutMs: Long = cfg.acquireReadyTimeoutMs, + stateStore: PoolStateStore = InMemoryPoolStateStore(), + ): SandboxPool { + val connectionConfig = + ConnectionConfig.builder() + .domain(cfg.mockDomain) + .protocol("http") + .requestTimeout(Duration.ofSeconds(30)) + .disableMetrics() + .also { builder -> + // Shared connection pool by default (fixed 500 idle slots) + // so high-concurrency runs do not hit the per-sandbox + // fresh-connection churn documented in + // docs/guides/client-pool.md. Explicit 0 disables sharing. + if (cfg.sharedConnectionPoolSize > 0) { + builder.connectionPool( + okhttp3.ConnectionPool( + cfg.sharedConnectionPoolSize, + 5, + TimeUnit.MINUTES, + ), + ) + } + } + .build() + return SandboxPool.builder() + .poolName(poolName) + .ownerId("bench-owner-$poolName") + .maxIdle(maxIdle) + .stateStore(stateStore) + .connectionConfig(connectionConfig) + .creationSpec( + PoolCreationSpec.builder() + .image("benchmark:mock") + .entrypoint("tail", "-f", "/dev/null") + .build(), + ) + .warmupConcurrency(warmupConcurrency) + .reconcileInterval(Duration.ofMillis(reconcileIntervalMs)) + .acquireReadyTimeout(Duration.ofMillis(acquireReadyTimeoutMs)) + .warmupReadyTimeout(Duration.ofMillis(cfg.warmupReadyTimeoutMs)) + .acquireHealthCheckPollingInterval(Duration.ofMillis(cfg.healthCheckPollingIntervalMs)) + .warmupHealthCheckPollingInterval(Duration.ofMillis(cfg.healthCheckPollingIntervalMs)) + .idleTimeout(Duration.ofSeconds(idleTimeoutS)) + .maxAcquireRetries(maxAcquireRetries) + .also { builder -> + if (cfg.acquireMinRemainingTtlS > 0) { + builder.acquireMinRemainingTtl(Duration.ofSeconds(cfg.acquireMinRemainingTtlS)) + } + if (cfg.primaryLockTtlS > 0) { + builder.primaryLockTtl(Duration.ofSeconds(cfg.primaryLockTtlS)) + } + if (cfg.degradedThreshold > 0) { + builder.degradedThreshold(cfg.degradedThreshold) + } + } + .build() + } + + val DEFAULT_POLICY = AcquirePolicy.DIRECT_CREATE + val RETRY_POLICY = AcquirePolicy.RETRY_NEXT_IDLE_THEN_CREATE + + /** Polls snapshot until idleCount reaches [target]; returns elapsed ms or -1 on timeout. */ + fun waitForIdle( + pool: SandboxPool, + target: Int, + timeoutMs: Long, + ): Long { + val start = System.nanoTime() + val deadline = start + timeoutMs * 1_000_000 + while (true) { + val idle = pool.snapshot().idleCount + if (idle >= target) { + return (System.nanoTime() - start) / 1_000_000 + } + if (System.nanoTime() > deadline) { + return -1L + } + Thread.sleep(100) + } + } + + /** Polls snapshot until idleCount drops to [target] or below (shrink); returns elapsed ms or -1 on timeout. */ + fun waitForIdleBelow( + pool: SandboxPool, + target: Int, + timeoutMs: Long, + ): Long { + val start = System.nanoTime() + val deadline = start + timeoutMs * 1_000_000 + while (true) { + val idle = pool.snapshot().idleCount + if (idle <= target) { + return (System.nanoTime() - start) / 1_000_000 + } + if (System.nanoTime() > deadline) { + return -1L + } + Thread.sleep(100) + } + } +} diff --git a/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/RatePacer.kt b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/RatePacer.kt new file mode 100644 index 000000000..3e4d2ff75 --- /dev/null +++ b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/RatePacer.kt @@ -0,0 +1,64 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.alibaba.opensandbox.benchmark + +import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicLong + +/** + * Spreads acquires evenly over wall-clock time at a fixed per-minute rate. + * + * Every caller claims the next absolute slot (`start + k * interval`); slots + * are advanced with CAS so concurrent workers never share one. Slots that + * fall in the past are executed immediately, so a slow worker lets the rate + * catch up on the next slots — the long-run average stays at [ratePerMin] + * without drift. No-op when [ratePerMin] is <= 0. + */ +class RatePacer(private val ratePerMin: Int) { + private val intervalNanos = + if (ratePerMin <= 0) 0L else TimeUnit.MINUTES.toNanos(1) / ratePerMin + private val nextSlot = AtomicLong(0) + + /** Blocks until this caller's slot is due. */ + fun waitForSlot() { + if (intervalNanos <= 0) return + var slot = nextSlot.get() + while (true) { + if (slot == 0L) { + if (nextSlot.compareAndSet(0L, System.nanoTime())) { + slot = System.nanoTime() + break + } + slot = nextSlot.get() + continue + } + val next = slot + intervalNanos + if (nextSlot.compareAndSet(slot, next)) { + break + } + slot = nextSlot.get() + } + val waitMs = (slot - System.nanoTime()) / 1_000_000 + if (waitMs > 0) { + try { + Thread.sleep(waitMs) + } catch (_: InterruptedException) { + Thread.currentThread().interrupt() + } + } + } +} diff --git a/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Scenarios.kt b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Scenarios.kt new file mode 100644 index 000000000..f0994626c --- /dev/null +++ b/tests/benchmark/kotlin/src/main/kotlin/com/alibaba/opensandbox/benchmark/Scenarios.kt @@ -0,0 +1,645 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.alibaba.opensandbox.benchmark + +import com.alibaba.opensandbox.sandbox.Sandbox +import com.alibaba.opensandbox.sandbox.domain.exceptions.PoolAcquireFailedException +import com.alibaba.opensandbox.sandbox.domain.exceptions.PoolDestroyedException +import com.alibaba.opensandbox.sandbox.domain.exceptions.PoolEmptyException +import com.alibaba.opensandbox.sandbox.domain.exceptions.PoolNotRunningException +import com.alibaba.opensandbox.sandbox.domain.exceptions.PoolStateStoreUnavailableException +import com.alibaba.opensandbox.sandbox.domain.exceptions.SandboxReadyTimeoutException +import com.alibaba.opensandbox.sandbox.domain.pool.AcquirePolicy +import com.alibaba.opensandbox.sandbox.domain.pool.PoolState +import com.alibaba.opensandbox.sandbox.pool.PoolWarmupDiagnostics +import com.alibaba.opensandbox.sandbox.pool.SandboxPool +import kotlinx.serialization.json.JsonPrimitive +import java.io.File +import java.util.concurrent.CountDownLatch +import java.util.concurrent.CopyOnWriteArrayList +import java.util.concurrent.Executors +import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicBoolean +import java.util.concurrent.atomic.AtomicLong +import kotlin.random.Random + +/** + * Benchmark scenarios. Every scenario returns a flat-ish map that the report + * writer renders into JSON and Markdown. Each scenario owns a fresh pool and + * shuts it down before returning. + */ +object Scenarios { + + // ---------- cold-start ---------- + + fun coldStart(cfg: BenchmarkConfig, mock: MockControl): Map { + mock.reset() + PoolWarmupDiagnostics.reset() + val pool = PoolRunner.build(cfg, "cold-start") + pool.start() + val t0 = System.nanoTime() + val fillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + val stats = mock.stats() + pool.shutdown(graceful = false) + val diagnostics = PoolWarmupDiagnostics.snapshot() + + val created = num(stats, "stats.created") + return mapOf( + "fillTimeMs" to fillMs, + "timedOut" to (fillMs < 0), + "serverCreated" to created, + "serverAliveAtFill" to num(stats, "alive"), + "overCreationOvershoot" to (created - cfg.maxIdle).coerceAtLeast(0), + "warmupPipeline" to + mapOf( + "queueWaitMs" to phaseStats(diagnostics.queueWaitMs), + "createDurationMs" to phaseStats(diagnostics.createDurationMs), + "commitDurationMs" to phaseStats(diagnostics.commitDurationMs), + "tickIntervalMs" to phaseStats(diagnostics.tickIntervalMs), + "tickDurationMs" to phaseStats(diagnostics.tickDurationMs), + "submitBurst" to phaseStats(diagnostics.submitBurst), + "submitCalls" to diagnostics.submitCalls, + "inFlightPeak" to diagnostics.inFlightPeak, + "inFlightMean" to diagnostics.inFlightMean, + "createFailures" to diagnostics.createFailures, + ), + ) + } + + // ---------- warm-pool acquire latency ---------- + + fun warmLatency(cfg: BenchmarkConfig, mock: MockControl): Map { + mock.reset() + val pool = PoolRunner.build(cfg, "warm-latency") + pool.start() + val fillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + val createdBefore = num(mock.stats(), "stats.created") + + val latency = LatencyCollector() + val start = CountDownLatch(1) + val workers = cfg.warmWorkers + val rounds = cfg.warmRoundsPerWorker + val probe = PoolProbe(pool) + probe.start() + val threads = Executors.newFixedThreadPool(workers) + repeat(workers) { + threads.submit { + start.await() + repeat(rounds) { + val t0 = System.nanoTime() + try { + val sb = pool.acquire(cfg.acquireTimeout, PoolRunner.DEFAULT_POLICY) + latency.record((System.nanoTime() - t0) / 1_000_000) + killAndClose(sb) + } catch (t: Throwable) { + latency.recordFailure(classifyFailure(t)) + } + } + } + } + start.countDown() + threads.shutdown() + threads.awaitTermination(10, TimeUnit.MINUTES) + probe.stop() + probe.writeCsv(File(cfg.reportDir, "client-warm-latency.csv")) + + val createdDelta = num(mock.stats(), "stats.created") - createdBefore + val acquires = rounds * workers.toLong() + val latencyStats = latency.snapshot() + pool.shutdown(graceful = false) + + return mapOf( + "fillTimeMs" to fillMs, + "latency" to latencyStats.toMap(), + "acquires" to acquires, + "successRate" to successRate(latencyStats), + "serverCreatedDelta" to createdDelta, + "hitRatio" to (1.0 - createdDelta.toDouble() / acquires).coerceIn(0.0, 1.0), + "client" to probe.report(), + ) + } + + // ---------- steady-state throughput ---------- + + fun steadyState(cfg: BenchmarkConfig, mock: MockControl): Map { + mock.reset() + val pool = PoolRunner.build(cfg, "steady-state") + pool.start() + // Extreme-cold-start mode: loaders race the fill instead of waiting + // for the idle buffer (fillMs = -1 in that case). + val fillMs = + if (cfg.steadyStartImmediately) { + -1L + } else { + PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + } + val createdBefore = num(mock.stats(), "stats.created") + val killedBefore = num(mock.stats(), "stats.killed") + + val durationMs = cfg.steadyDurationS * 1000L + val deadline = System.nanoTime() + durationMs * 1_000_000 + val running = AtomicBoolean(true) + val latency = LatencyCollector() + val acquires = AtomicLong(0) + + val probe = PoolProbe(pool) + probe.start() + + val rng = Random(System.nanoTime()) + val pacer = RatePacer(cfg.acquireRatePerMin) + val threads = Executors.newFixedThreadPool(cfg.steadyWorkers) + val loaderStart = System.nanoTime() + repeat(cfg.steadyWorkers) { + threads.submit { + while (System.nanoTime() < deadline) { + pacer.waitForSlot() + val t0 = System.nanoTime() + try { + val sb = pool.acquire(cfg.acquireTimeout, PoolRunner.DEFAULT_POLICY) + latency.record((System.nanoTime() - t0) / 1_000_000) + acquires.incrementAndGet() + Thread.sleep(rng.nextLong(cfg.holdMinMs, cfg.holdMaxMs + 1)) + killAndClose(sb) + } catch (t: Throwable) { + latency.recordFailure(classifyFailure(t)) + Thread.sleep(200) + } + } + } + } + threads.shutdown() + // The loaders run for the full configured duration; the wait must not + // truncate them (a 15-min cap would silently halve a 30-min run). + threads.awaitTermination(durationMs / 1000 + 300, TimeUnit.SECONDS) + val loaderDurationMs = (System.nanoTime() - loaderStart) / 1_000_000 + running.set(false) + probe.stop() + probe.writeCsv(File(cfg.reportDir, "client-steady-state.csv")) + + val createdDelta = num(mock.stats(), "stats.created") - createdBefore + val killedDelta = num(mock.stats(), "stats.killed") - killedBefore + val latencyStats = latency.snapshot() + val client = probe.report() + pool.shutdown(graceful = false) + + val idleStat = client["poolIdleCount"] as Map + return mapOf( + "fillTimeMs" to fillMs, + "startImmediately" to cfg.steadyStartImmediately, + "durationMs" to durationMs, + "actualLoaderDurationMs" to loaderDurationMs, + "workers" to cfg.steadyWorkers, + "targetAcquiresPerMin" to cfg.acquireRatePerMin, + "acquiredCount" to acquires.get(), + "achievedAcquiresPerMin" to + (acquires.get().toDouble() * 60_000 / loaderDurationMs.coerceAtLeast(1)), + "throughputAcquiresPerSec" to (acquires.get().toDouble() * 1000 / loaderDurationMs.coerceAtLeast(1)), + "successRate" to successRate(latencyStats), + "latency" to latencyStats.toMap(), + "serverCreatedDelta" to createdDelta, + "serverKilledDelta" to killedDelta, + "replenishRatePerSec" to (createdDelta.toDouble() / cfg.steadyDurationS), + "killRatePerSec" to (killedDelta.toDouble() / cfg.steadyDurationS), + "directCreateRatio" to + ((createdDelta - killedDelta).coerceAtLeast(0).toDouble() / acquires.get().coerceAtLeast(1)), + "idleSamples" to (idleStat["samples"] as Int), + "idleMin" to (idleStat["min"] as Long), + "idleMean" to (idleStat["mean"] as Double), + "idleEmptyRatio" to (client["poolIdleZeroRatio"] as Double), + "client" to client, + ) + } + + // ---------- replenish lag ---------- + + fun replenishLag(cfg: BenchmarkConfig, mock: MockControl): Map { + mock.reset() + val pool = PoolRunner.build(cfg, "replenish-lag") + pool.start() + val fillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + + val lags = LatencyCollector() + var timedOut = 0L + repeat(cfg.replenishRounds) { + val sb = pool.acquire(cfg.acquireTimeout, PoolRunner.DEFAULT_POLICY) + killAndClose(sb) + val t0 = System.nanoTime() + val lag = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.replenishWaitTimeoutMs) + if (lag < 0) { + timedOut++ + } else { + lags.record(lag) + } + Thread.sleep(50) + } + pool.shutdown(graceful = false) + + return mapOf( + "fillTimeMs" to fillMs, + "rounds" to cfg.replenishRounds, + "replenishLagMs" to lags.snapshot().toMap(), + "timedOutRounds" to timedOut, + ) + } + + // ---------- creation-failure injection ---------- + + fun failureInjection(cfg: BenchmarkConfig, mock: MockControl): Map { + mock.reset() + val pool = PoolRunner.build(cfg, "failure-injection", maxAcquireRetries = 3) + pool.start() + val fillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + pool.releaseAllIdle() + + mock.setFaults(createFailureRate = cfg.failureCreateRate) + Thread.sleep(500) + + val latency = LatencyCollector() + repeat(cfg.failureAcquires) { + val t0 = System.nanoTime() + try { + val sb = pool.acquire(cfg.acquireTimeout, PoolRunner.RETRY_POLICY) + latency.record((System.nanoTime() - t0) / 1_000_000) + killAndClose(sb) + } catch (t: Throwable) { + latency.recordFailure(classifyFailure(t)) + } + } + Thread.sleep(500) + val snap = pool.snapshot() + + mock.setFaults(createFailureRate = 0.0) + val refillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + val stats = mock.stats() + pool.shutdown(graceful = false) + + return mapOf( + "fillTimeMs" to fillMs, + "createFailureRate" to cfg.failureCreateRate, + "acquires" to cfg.failureAcquires, + "latency" to latency.snapshot().toMap(), + "poolStateAfterBurst" to snap.state.name, + "backoffActive" to snap.backoffActive, + "failureCount" to snap.failureCount, + "lastError" to (snap.lastError ?: ""), + "serverCreateFailed" to num(stats, "stats.createFailed"), + "refillTimeMsAfterRecovery" to refillMs, + ) + } + + // ---------- stale idle sandboxes ---------- + + fun staleIdle(cfg: BenchmarkConfig, mock: MockControl): Map { + mock.reset() + val pool = + PoolRunner.build( + cfg, + "stale-idle", + maxAcquireRetries = cfg.staleRetries, + acquireReadyTimeoutMs = cfg.staleAcquireReadyTimeoutMs, + ) + pool.start() + val fillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + val createdBefore = num(mock.stats(), "stats.created") + + // Poison a fraction (default 1.0 = all) of the currently-alive + // sandboxes: their execd endpoints start failing, so idle candidates + // cannot be connected. Partial poisoning simulates real-world failure + // where the pool must skip bad candidates and return good ones. + mock.setFaults(poisonRate = cfg.stalePoisonRate) + + val latency = LatencyCollector() + repeat(cfg.staleAcquires) { + val t0 = System.nanoTime() + try { + val sb = pool.acquire(cfg.acquireTimeout, PoolRunner.RETRY_POLICY) + latency.record((System.nanoTime() - t0) / 1_000_000) + killAndClose(sb) + } catch (t: Throwable) { + latency.recordFailure(classifyFailure(t)) + } + } + val stats = mock.stats() + + // Pool must drain the stale idles and refill with fresh sandboxes. + val refillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + pool.shutdown(graceful = false) + + return mapOf( + "fillTimeMs" to fillMs, + "acquires" to cfg.staleAcquires, + "poisonRate" to cfg.stalePoisonRate, + "latency" to latency.snapshot().toMap(), + "successRate" to successRate(latency.snapshot()), + "serverExecdPoisoned" to num(stats, "stats.execdPoisoned"), + "serverCreatedDelta" to (num(stats, "stats.created") - createdBefore), + "serverAliveAfter" to num(stats, "alive"), + "refillTimeMs" to refillMs, + ) + } + + // ---------- server-side TTL self-heal ---------- + + fun idleExpiry(cfg: BenchmarkConfig, mock: MockControl): Map { + mock.reset() + val idleTimeoutS = cfg.idleExpiryIdleTimeoutS + val pool = + PoolRunner.build( + cfg, + "idle-expiry", + maxIdle = 10, + warmupConcurrency = 2, + reconcileIntervalMs = 500, + idleTimeoutS = idleTimeoutS, + ) + pool.start() + val fillMs = PoolRunner.waitForIdle(pool, 10, cfg.coldStartTimeoutMs) + val createdBefore = num(mock.stats(), "stats.created") + + val durationMs = cfg.idleExpiryDurationS * 1000L + val deadline = System.nanoTime() + durationMs * 1_000_000 + val running = AtomicBoolean(true) + val idleSamples = CopyOnWriteArrayList() + val sampler = Thread { + while (running.get()) { + idleSamples.add(pool.snapshot().idleCount) + Thread.sleep(250) + } + } + sampler.isDaemon = true + sampler.start() + + while (System.nanoTime() < deadline) { + Thread.sleep(100) + } + running.set(false) + + val stats = mock.stats() + val createdDelta = num(stats, "stats.created") - createdBefore + val idleMean = if (idleSamples.isEmpty()) 0.0 else idleSamples.average() + val idleMin = idleSamples.minOrNull() ?: 0 + pool.shutdown(graceful = false) + + return mapOf( + "idleTimeoutS" to idleTimeoutS, + "fillTimeMs" to fillMs, + "durationMs" to durationMs, + "serverCreatedDelta" to createdDelta, + "serverKilled" to num(stats, "stats.killed"), + "idleMean" to idleMean, + "idleMin" to idleMin, + "idleSamples" to idleSamples.size, + ) + } + + // ---------- resize (shrink + regrow) ---------- + + fun resize(cfg: BenchmarkConfig, mock: MockControl): Map { + mock.reset() + val pool = PoolRunner.build(cfg, "resize") + pool.start() + val fillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + val createdBefore = num(mock.stats(), "stats.created") + val killedBefore = num(mock.stats(), "stats.killed") + + // Shrink: excess idles must be drained and killed by reconcile. + val shrinkTarget = maxOf(1, cfg.maxIdle / 2) + pool.resize(shrinkTarget) + val shrinkMs = PoolRunner.waitForIdleBelow(pool, shrinkTarget, cfg.coldStartTimeoutMs) + Thread.sleep(1000) // let server-side kills settle + val killedDuringShrink = num(mock.stats(), "stats.killed") - killedBefore + + // Regrow back to the original target. + pool.resize(cfg.maxIdle) + val regrowMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + val stats = mock.stats() + pool.shutdown(graceful = false) + + return mapOf( + "fillTimeMs" to fillMs, + "shrinkTo" to shrinkTarget, + "shrinkTimeMs" to shrinkMs, + "killedDuringShrink" to killedDuringShrink, + "regrowTimeMs" to regrowMs, + "serverCreatedDelta" to (num(stats, "stats.created") - createdBefore), + "serverAliveAtEnd" to num(stats, "alive"), + ) + } + + // ---------- acquire racing graceful shutdown ---------- + + fun shutdownRace(cfg: BenchmarkConfig, mock: MockControl): Map { + mock.reset() + val pool = PoolRunner.build(cfg, "shutdown-race") + pool.start() + val fillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + + val shutdownStarted = AtomicBoolean(false) + val stop = AtomicBoolean(false) + // Collector 1: attempts while the pool is RUNNING. Collector 2: the + // race window itself (acquires racing DRAINING/STOPPED). + val runningLatency = LatencyCollector() + val raceLatency = LatencyCollector() + val runningAttempts = AtomicLong(0) + val raceAttempts = AtomicLong(0) + val threads = Executors.newFixedThreadPool(cfg.steadyWorkers) + repeat(cfg.steadyWorkers) { + threads.submit { + while (!stop.get()) { + val inRace = shutdownStarted.get() + if (inRace) raceAttempts.incrementAndGet() else runningAttempts.incrementAndGet() + val latency = if (inRace) raceLatency else runningLatency + try { + val sb = pool.acquire(cfg.acquireTimeout, PoolRunner.DEFAULT_POLICY) + latency.record(0) + killAndClose(sb) + } catch (t: Throwable) { + latency.recordFailure(classifyFailure(t)) + } + } + } + } + Thread.sleep(2000) // let workers hammer the warm pool + shutdownStarted.set(true) + val shutdownT0 = System.nanoTime() + pool.shutdown(graceful = true) + val shutdownMs = (System.nanoTime() - shutdownT0) / 1_000_000 + stop.set(true) + threads.shutdownNow() + threads.awaitTermination(30, TimeUnit.SECONDS) + + val runningStats = runningLatency.snapshot() + val raceStats = raceLatency.snapshot() + return mapOf( + "fillTimeMs" to fillMs, + "runningPhaseAttempts" to runningAttempts.get(), + "runningPhase" to + mapOf( + "successRate" to successRate(runningStats), + "latency" to runningStats.toMap(), + ), + "raceWindowAttempts" to raceAttempts.get(), + "raceWindow" to + mapOf( + "successRate" to successRate(raceStats), + "latency" to raceStats.toMap(), + "rejectedDuringDraining" to (raceStats.failuresByType["poolNotRunning"] ?: 0L), + ), + "shutdownMs" to shutdownMs, + ) + } + + // ---------- state-store outage (OSEP-0005 fallthrough) ---------- + + fun storeOutage(cfg: BenchmarkConfig, mock: MockControl): Map { + mock.reset() + val store = FailingPoolStateStore() + val pool = + PoolRunner.build( + cfg, + "store-outage", + maxAcquireRetries = 1, + stateStore = store, + ) + pool.start() + val fillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + val createdBefore = num(mock.stats(), "stats.created") + + // Drain part of the idle buffer before the outage so recovery has + // real refill work (the pool cannot commit warmups while the store + // is down, and fallthrough acquires never touch the store). + val drainCount = maxOf(1, cfg.maxIdle / 3) + repeat(drainCount) { + val sb = pool.acquire(cfg.acquireTimeout, PoolRunner.DEFAULT_POLICY) + killAndClose(sb) + } + + // Phase 1: store down, DIRECT_CREATE policy must fall through to + // direct create (OSEP-0005) and keep acquire available. + store.setFailing(true) + Thread.sleep(500) + val phase1 = LatencyCollector() + repeat(cfg.failureAcquires) { + val t0 = System.nanoTime() + try { + val sb = pool.acquire(cfg.acquireTimeout, PoolRunner.DEFAULT_POLICY) + phase1.record((System.nanoTime() - t0) / 1_000_000) + killAndClose(sb) + } catch (t: Throwable) { + phase1.recordFailure(classifyFailure(t)) + } + } + val phase1Stats = phase1.snapshot() + val errorsPhase1 = store.errorCount() + + // Phase 2: store still down, FAIL_FAST must fail closed and surface + // PoolStateStoreUnavailableException. + val phase2 = LatencyCollector() + repeat(10) { + try { + val sb = pool.acquire(cfg.acquireTimeout, AcquirePolicy.FAIL_FAST) + phase2.record(0) + killAndClose(sb) + } catch (t: Throwable) { + phase2.recordFailure(classifyFailure(t)) + } + } + val phase2Stats = phase2.snapshot() + + // Phase 3: store recovers; the pool must refill. + store.setFailing(false) + val refillMs = PoolRunner.waitForIdle(pool, cfg.maxIdle, cfg.coldStartTimeoutMs) + val stats = mock.stats() + pool.shutdown(graceful = false) + + return mapOf( + "fillTimeMs" to fillMs, + "phase1DirectCreateFallthrough" to phase1Stats.toMap(), + "phase1StoreErrorCount" to errorsPhase1, + "phase2FailFastFailClosed" to phase2Stats.toMap(), + "refillTimeMsAfterRecovery" to refillMs, + "serverCreatedDelta" to (num(stats, "stats.created") - createdBefore), + ) + } + + // ---------- helpers ---------- + + private fun phaseStats(s: PoolWarmupDiagnostics.PhaseStats): Map = + mapOf( + "count" to s.count, + "meanMs" to s.meanMs, + "p50Ms" to s.p50Ms, + "p95Ms" to s.p95Ms, + "maxMs" to s.maxMs, + ) + + private fun successRate(stats: LatencyStats): Double { + val total = stats.n + stats.failures + return if (total == 0L) 0.0 else (stats.n.toDouble() / total) + } + + private fun classifyFailure(t: Throwable): String = + when (t) { + is SandboxReadyTimeoutException -> "readyTimeout" + is PoolNotRunningException -> "poolNotRunning" + is PoolEmptyException -> "poolEmpty" + is PoolAcquireFailedException -> "acquireFailed" + is PoolDestroyedException -> "poolDestroyed" + is PoolStateStoreUnavailableException -> "storeUnavailable" + else -> "other" + } + + private fun killAndClose(sandbox: Sandbox) { + try { + sandbox.kill() + } finally { + try { + sandbox.close() + } catch (_: Exception) { + // ignore + } + } + } + + private fun num(stats: Map, dottedKey: String): Long { + var cur: Any? = stats + for (part in dottedKey.split(".")) { + if (cur !is Map<*, *>) return 0L + cur = cur[part] + } + return when (cur) { + is Number -> cur.toLong() + is String -> cur.toLongOrNull() ?: 0L + is JsonPrimitive -> cur.content.toLongOrNull() ?: 0L + else -> 0L + } + } + + val ALL_SCENARIOS = + mapOf( + "cold-start" to ::coldStart, + "warm-latency" to ::warmLatency, + "steady-state" to ::steadyState, + "replenish-lag" to ::replenishLag, + "failure-injection" to ::failureInjection, + "stale-idle" to ::staleIdle, + "idle-expiry" to ::idleExpiry, + "resize" to ::resize, + "shutdown-race" to ::shutdownRace, + "store-outage" to ::storeOutage, + ) +} diff --git a/tests/benchmark/kotlin/src/main/resources/simplelogger.properties b/tests/benchmark/kotlin/src/main/resources/simplelogger.properties new file mode 100644 index 000000000..66abf87b5 --- /dev/null +++ b/tests/benchmark/kotlin/src/main/resources/simplelogger.properties @@ -0,0 +1,4 @@ +org.slf4j.simpleLogger.defaultLogLevel=error +org.slf4j.simpleLogger.showThreadName=false +org.slf4j.simpleLogger.showDateTime=false +org.slf4j.simpleLogger.showLogName=false diff --git a/tests/benchmark/mockserver/config.go b/tests/benchmark/mockserver/config.go new file mode 100644 index 000000000..5b43153c8 --- /dev/null +++ b/tests/benchmark/mockserver/config.go @@ -0,0 +1,191 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package main + +import ( + "encoding/json" + "fmt" + "math" + "math/rand" + "os" + "time" +) + +// Config models one mock-server run. It is loaded from a JSON file at startup +// and can be mutated at runtime through POST /__config (see FaultConfig). +type Config struct { + // CreateLatencyMs controls how long POST /v1/sandboxes takes before + // returning, mimicking real sandbox provisioning time. + CreateLatencyMs LatencySpec `json:"createLatencyMs"` + // CreateFailureRate is the probability [0,1] that POST /v1/sandboxes + // returns HTTP 500 instead of creating a sandbox. + CreateFailureRate float64 `json:"createFailureRate"` + // BootDelayMs is how long a freshly created sandbox stays in Pending + // state. Its execd endpoint (and therefore the SDK readiness probe) + // only succeeds after the boot delay elapses. + BootDelayMs int64 `json:"bootDelayMs"` + // ExecdFailureRate is the probability [0,1] that any execd request + // returns HTTP 500. + ExecdFailureRate float64 `json:"execdFailureRate"` + // DefaultTtl is the server-side lifetime assigned to created sandboxes + // when the request does not carry a timeout. Sandboxes are reaped at + // expiry: lifecycle lookups and execd pings start failing. + DefaultTtl time.Duration `json:"-"` + DefaultTtlSeconds int64 `json:"defaultTtlSeconds"` + // LatencyOverrides controls the response time of every other API route + // (lifecycle.get, lifecycle.delete, lifecycle.renew, lifecycle.endpoint, + // execd.ping, execd.other). Routes without an override respond + // immediately. An entry for lifecycle.create overrides CreateLatencyMs. + LatencyOverrides map[string]LatencySpec `json:"latencyOverrides"` +} + +type LatencySpec struct { + // Distribution: "lognormal" (default), "uniform", or "fixed". + Distribution string `json:"distribution"` + MeanMs float64 `json:"meanMs"` + StddevMs float64 `json:"stddevMs"` + MinMs float64 `json:"minMs"` + // MaxMs is the upper bound for "uniform" sampling. + MaxMs float64 `json:"maxMs"` +} + +// FaultConfig is the mutable runtime subset of Config, updated via +// POST /__config. Fields are optional; absent fields keep current values. +type FaultConfig struct { + CreateFailureRate *float64 `json:"createFailureRate"` + ExecdFailureRate *float64 `json:"execdFailureRate"` + BootDelayMs *int64 `json:"bootDelayMs"` + CreateLatencyMs *LatencySpec `json:"createLatencyMs"` + // LatencyOverrides replaces the whole per-route response-time map when + // present (even when empty). + LatencyOverrides *map[string]LatencySpec `json:"latencyOverrides"` + // PoisonExisting flips every currently-alive sandbox into a poisoned + // state: its execd endpoint starts failing so SDK connects to it break. + // Newly created sandboxes are unaffected. Used to simulate stale idle + // sandboxes (e.g. sandboxes that died server-side). + PoisonExisting bool `json:"poisonExisting"` + // PoisonRate flips a random subset (probability [0,1]) of currently-alive + // sandboxes into the poisoned state, simulating partial failure. + PoisonRate *float64 `json:"poisonRate"` +} + +func loadConfig(path string) (*Config, error) { + cfg := &Config{ + // Default response-time profile: create and delete take a uniform + // 300-800ms, execd ping a fixed 100ms, every other API a uniform + // 50-100ms. + CreateLatencyMs: LatencySpec{ + Distribution: "uniform", + MinMs: 300, + MaxMs: 800, + }, + LatencyOverrides: map[string]LatencySpec{ + "lifecycle.delete": {Distribution: "uniform", MinMs: 300, MaxMs: 800}, + "lifecycle.get": {Distribution: "uniform", MinMs: 50, MaxMs: 100}, + "lifecycle.renew": {Distribution: "uniform", MinMs: 50, MaxMs: 100}, + "lifecycle.endpoint": {Distribution: "uniform", MinMs: 50, MaxMs: 100}, + "execd.ping": {Distribution: "fixed", MeanMs: 100}, + "execd.other": {Distribution: "uniform", MinMs: 50, MaxMs: 100}, + }, + DefaultTtlSeconds: 3600, + } + if path != "" { + raw, err := os.ReadFile(path) + if err != nil { + return nil, fmt.Errorf("read config: %w", err) + } + // The default latencyOverrides map is pre-seeded, and encoding/json + // merges into existing maps instead of replacing them. Detect whether + // the file mentions the key and drop the defaults first so an empty + // map in the file means "no route latency" rather than "keep defaults". + var keys map[string]json.RawMessage + if err := json.Unmarshal(raw, &keys); err != nil { + return nil, fmt.Errorf("parse config: %w", err) + } + if _, ok := keys["latencyOverrides"]; ok { + cfg.LatencyOverrides = nil + } + if err := json.Unmarshal(raw, cfg); err != nil { + return nil, fmt.Errorf("parse config: %w", err) + } + } + if cfg.CreateLatencyMs.Distribution == "" { + cfg.CreateLatencyMs.Distribution = "lognormal" + } + cfg.DefaultTtl = time.Duration(cfg.DefaultTtlSeconds) * time.Second + return cfg, nil +} + +func (cfg *Config) applyFault(f FaultConfig) { + if f.CreateFailureRate != nil { + cfg.CreateFailureRate = *f.CreateFailureRate + } + if f.ExecdFailureRate != nil { + cfg.ExecdFailureRate = *f.ExecdFailureRate + } + if f.BootDelayMs != nil { + cfg.BootDelayMs = *f.BootDelayMs + } + if f.CreateLatencyMs != nil { + cfg.CreateLatencyMs = *f.CreateLatencyMs + } + if f.LatencyOverrides != nil { + cfg.LatencyOverrides = *f.LatencyOverrides + } +} + +// latencyFor returns the latency spec for one API route, or nil when the +// route should respond immediately. +func (cfg *Config) latencyFor(route string) *LatencySpec { + if cfg.LatencyOverrides != nil { + if spec, ok := cfg.LatencyOverrides[route]; ok { + return &spec + } + } + return nil +} + +// sample returns a latency duration drawn from the configured distribution. +// Uses math/rand's global functions, which are goroutine-safe. +func (s *LatencySpec) sample() time.Duration { + switch s.Distribution { + case "fixed": + return time.Duration(s.MeanMs) * time.Millisecond + case "uniform": + lo := s.MinMs + hi := s.MaxMs + if hi <= lo { + hi = s.MeanMs + } + if hi < lo { + hi, lo = lo, hi + } + ms := lo + rand.Float64()*(hi-lo) + return time.Duration(ms) * time.Millisecond + default: // lognormal + mean := math.Max(s.MeanMs, 0.001) + stddev := math.Max(s.StddevMs, 0.001) + mu := math.Log(mean * mean / math.Sqrt(mean*mean+stddev*stddev)) + sigma := math.Sqrt(math.Log(1 + stddev*stddev/(mean*mean))) + v := mu + sigma*rand.NormFloat64() + ms := math.Exp(v) + if ms < s.MinMs { + ms = s.MinMs + } + return time.Duration(ms) * time.Millisecond + } +} diff --git a/tests/benchmark/mockserver/go.mod b/tests/benchmark/mockserver/go.mod new file mode 100644 index 000000000..82aa3ebae --- /dev/null +++ b/tests/benchmark/mockserver/go.mod @@ -0,0 +1,17 @@ +// Copyright 2026 Alibaba Group Holding Ltd. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +module github.com/alibaba/OpenSandbox/tests/benchmark/mockserver + +go 1.20 diff --git a/tests/benchmark/mockserver/main.go b/tests/benchmark/mockserver/main.go new file mode 100644 index 000000000..5d5488da8 --- /dev/null +++ b/tests/benchmark/mockserver/main.go @@ -0,0 +1,84 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package main + +import ( + "flag" + "fmt" + "log" + "net" + "net/http" + "strconv" + "time" +) + +func main() { + lifecycleAddr := flag.String("lifecycle-addr", "127.0.0.1:18080", "lifecycle API listen address") + execdAddr := flag.String("execd-addr", "127.0.0.1:18081", "execd API listen address") + configPath := flag.String("config", "", "path to a mock server config JSON file") + statsWindowSec := flag.Int("stats-window-sec", DefaultStatsWindowSec, "per-API QPS history window in seconds") + flag.Parse() + + cfg, err := loadConfig(*configPath) + if err != nil { + log.Fatalf("load config: %v", err) + } + + host, port, err := splitHostPort(*execdAddr) + if err != nil { + log.Fatalf("invalid execd addr: %v", err) + } + mock := newMockServer(cfg, host, port, *statsWindowSec) + + lifecycleMux := http.NewServeMux() + lifecycleMux.HandleFunc("/", mock.handleLifecycle) + execdMux := http.NewServeMux() + execdMux.HandleFunc("/", mock.handleExecd) + + lifecycleSrv := &http.Server{ + Addr: *lifecycleAddr, + Handler: lifecycleMux, + ReadHeaderTimeout: 10 * time.Second, + } + execdSrv := &http.Server{ + Addr: *execdAddr, + Handler: execdMux, + ReadHeaderTimeout: 10 * time.Second, + } + + log.Printf("mock lifecycle server listening on http://%s", *lifecycleAddr) + log.Printf("mock execd server listening on http://%s", *execdAddr) + + go mock.startAliveTicker() + + errCh := make(chan error, 2) + go func() { errCh <- lifecycleSrv.ListenAndServe() }() + go func() { errCh <- execdSrv.ListenAndServe() }() + log.Fatal(<-errCh) +} + +func splitHostPort(addr string) (string, int, error) { + host, portStr, err := net.SplitHostPort(addr) + if err != nil { + return "", 0, fmt.Errorf("invalid addr %q: %w", addr, err) + } + port, err := strconv.Atoi(portStr) + if err != nil { + return "", 0, fmt.Errorf("invalid port %q: %w", portStr, err) + } + return host, port, nil +} diff --git a/tests/benchmark/mockserver/server.go b/tests/benchmark/mockserver/server.go new file mode 100644 index 000000000..0d9d6eb42 --- /dev/null +++ b/tests/benchmark/mockserver/server.go @@ -0,0 +1,557 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package main + +import ( + "encoding/json" + "fmt" + "io" + "math/rand" + "net/http" + "strings" + "sync" + "sync/atomic" + "time" +) + +// MockServer implements the OpenSandbox lifecycle API surface used by the +// sandbox SDKs plus a per-sandbox execd listener. It is a standalone, +// language-agnostic mock intended for pool benchmarks: it simulates +// provisioning latency, sandbox boot time, server-side TTL expiry, and +// per-endpoint faults. +type MockServer struct { + cfgMu sync.RWMutex + cfg *Config + mu sync.RWMutex + nextSeq uint64 + sandboxes map[string]*Sandbox + + execdHost string + execdPort int + + stats Stats + qps *QpsRegistry + alive *gaugeTracker +} + +// Sandbox is the mock's view of a sandbox on the lifecycle side. +type Sandbox struct { + ID string + CreatedAt time.Time + ExpiresAt time.Time + State string // "Pending" | "Running" | "Terminated" + Poisoned bool +} + +func (s *Sandbox) alive(now time.Time) bool { + return s.State != "Terminated" && (s.ExpiresAt.IsZero() || s.ExpiresAt.After(now)) +} + +func (s *Sandbox) running(now time.Time) bool { + return s.alive(now) && s.State == "Running" +} + +func newMockServer(cfg *Config, execdHost string, execdPort int, windowSec int) *MockServer { + return &MockServer{ + cfg: cfg, + sandboxes: make(map[string]*Sandbox), + execdHost: execdHost, + execdPort: execdPort, + qps: newQpsRegistry(windowSec), + alive: newGaugeTracker(windowSec), + } +} + +// startAliveTicker records the alive-sandbox count once per second so the +// driver can see the pool-size trajectory (over-creation, shrink, drift). +func (m *MockServer) startAliveTicker() { + record := func(now time.Time) { + m.mu.RLock() + alive := 0 + for _, sb := range m.sandboxes { + if sb.alive(now) { + alive++ + } + } + m.mu.RUnlock() + m.alive.record(now, int64(alive)) + } + record(time.Now()) + ticker := time.NewTicker(time.Second) + for range ticker.C { + record(time.Now()) + } +} + +// recordQps attributes one request to a route. Handlers defer this at entry so +// every outcome (including faults) is counted. +func (m *MockServer) recordQps(route string, start time.Time) { + m.qps.record(route, time.Now(), time.Since(start)) +} + +// applyRouteLatency sleeps for the configured response time of [route], +// if any override exists. +func (m *MockServer) applyRouteLatency(route string) { + cfg := m.cfgSnapshot() + if spec := cfg.latencyFor(route); spec != nil { + time.Sleep(spec.sample()) + } +} + +// ---------- lifecycle handlers ---------- + +func (m *MockServer) handleLifecycle(w http.ResponseWriter, r *http.Request) { + path := r.URL.Path + switch { + case path == "/__stats": + m.handleStats(w, r) + case path == "/__config": + m.handleConfig(w, r) + case path == "/__reset": + m.handleReset(w, r) + case r.Method == http.MethodPost && path == "/v1/sandboxes": + m.handleCreate(w, r) + case r.Method == http.MethodDelete && strings.HasPrefix(path, "/v1/sandboxes/"): + m.handleDelete(w, r, strings.TrimPrefix(path, "/v1/sandboxes/")) + case r.Method == http.MethodPost && strings.HasSuffix(path, "/renew-expiration"): + id := strings.TrimSuffix(strings.TrimPrefix(path, "/v1/sandboxes/"), "/renew-expiration") + m.handleRenew(w, r, id) + case r.Method == http.MethodGet && strings.HasPrefix(path, "/v1/sandboxes/") && strings.Contains(path, "/endpoints/"): + parts := strings.Split(strings.TrimPrefix(path, "/v1/sandboxes/"), "/endpoints/") + m.handleEndpoint(w, r, parts[0]) + case r.Method == http.MethodGet && strings.HasPrefix(path, "/v1/sandboxes/"): + m.handleGet(w, r, strings.TrimPrefix(path, "/v1/sandboxes/")) + default: + writeError(w, http.StatusNotFound, "NOT_FOUND", "no such route: "+path) + } +} + +func (m *MockServer) handleCreate(w http.ResponseWriter, r *http.Request) { + start := time.Now() + defer m.recordQps("lifecycle.create", start) + _, _ = io.Copy(io.Discard, r.Body) + + cfg := m.cfgSnapshot() + if cfg.CreateFailureRate > 0 && rand.Float64() < cfg.CreateFailureRate { + m.stats.incCreateFailed() + writeError(w, http.StatusInternalServerError, "INTERNAL_ERROR", "simulated provisioning failure") + return + } + + latency := cfg.CreateLatencyMs.sample() + if override := cfg.latencyFor("lifecycle.create"); override != nil { + latency = override.sample() + } + time.Sleep(latency) + m.stats.recordCreateLatency(latency) + + now := time.Now().UTC() + expiresAt := now.Add(cfg.DefaultTtl) + id := m.newSandboxID() + sandbox := &Sandbox{ + ID: id, + CreatedAt: now, + ExpiresAt: expiresAt, + State: "Pending", + } + m.mu.Lock() + m.sandboxes[id] = sandbox + m.mu.Unlock() + + go m.bootSandbox(id, cfg.BootDelayMs) + + m.stats.incCreated() + writeJSON(w, http.StatusCreated, map[string]any{ + "id": id, + "status": sandboxStatus(sandbox, now), + "createdAt": now.Format(time.RFC3339), + "expiresAt": expiresAt.Format(time.RFC3339), + "entrypoint": []string{"tail", "-f", "/dev/null"}, + }) +} + +func (m *MockServer) bootSandbox(id string, delayMs int64) { + if delayMs <= 0 { + delayMs = 1 + } + time.Sleep(time.Duration(delayMs) * time.Millisecond) + m.mu.Lock() + if sb, ok := m.sandboxes[id]; ok && sb.State == "Pending" { + sb.State = "Running" + } + m.mu.Unlock() +} + +func (m *MockServer) handleGet(w http.ResponseWriter, r *http.Request, id string) { + start := time.Now() + defer m.recordQps("lifecycle.get", start) + m.applyRouteLatency("lifecycle.get") + _ = r.Body.Close() + m.stats.incSandboxGets() + m.mu.RLock() + sb := m.sandboxes[id] + m.mu.RUnlock() + if sb == nil || !sb.alive(time.Now()) { + writeError(w, http.StatusNotFound, "NOT_FOUND", "sandbox not found: "+id) + return + } + writeJSON(w, http.StatusOK, m.sandboxInfo(sb)) +} + +func (m *MockServer) handleDelete(w http.ResponseWriter, r *http.Request, id string) { + start := time.Now() + defer m.recordQps("lifecycle.delete", start) + m.applyRouteLatency("lifecycle.delete") + _ = r.Body.Close() + m.mu.Lock() + sb := m.sandboxes[id] + if sb != nil { + sb.State = "Terminated" + } + m.mu.Unlock() + m.stats.incKilled() + // DELETE of an unknown sandbox still succeeds (best-effort semantics); + // a killed sandbox's execd endpoint stops responding. + w.WriteHeader(http.StatusNoContent) +} + +func (m *MockServer) handleRenew(w http.ResponseWriter, r *http.Request, id string) { + start := time.Now() + defer m.recordQps("lifecycle.renew", start) + m.applyRouteLatency("lifecycle.renew") + defer r.Body.Close() + m.stats.incRenews() + m.mu.RLock() + sb := m.sandboxes[id] + m.mu.RUnlock() + if sb == nil || !sb.alive(time.Now()) { + writeError(w, http.StatusNotFound, "NOT_FOUND", "sandbox not found: "+id) + return + } + var body struct { + ExpiresAt string `json:"expiresAt"` + } + if err := json.NewDecoder(r.Body).Decode(&body); err != nil || body.ExpiresAt == "" { + writeError(w, http.StatusBadRequest, "INVALID_REQUEST", "expiresAt is required") + return + } + expiresAt, err := time.Parse(time.RFC3339, body.ExpiresAt) + if err != nil { + writeError(w, http.StatusBadRequest, "INVALID_REQUEST", "expiresAt must be RFC3339: "+err.Error()) + return + } + m.mu.Lock() + sb.ExpiresAt = expiresAt + m.mu.Unlock() + writeJSON(w, http.StatusOK, map[string]any{"expiresAt": expiresAt.Format(time.RFC3339)}) +} + +func (m *MockServer) handleEndpoint(w http.ResponseWriter, r *http.Request, id string) { + start := time.Now() + defer m.recordQps("lifecycle.endpoint", start) + m.applyRouteLatency("lifecycle.endpoint") + _ = r.Body.Close() + m.stats.incEndpointGets() + m.mu.RLock() + sb := m.sandboxes[id] + m.mu.RUnlock() + if sb == nil || !sb.alive(time.Now()) { + writeError(w, http.StatusNotFound, "NOT_FOUND", "sandbox not found: "+id) + return + } + // The endpoint URL and token let the execd listener attribute requests + // to this sandbox so boot state and poisoning can be enforced. + writeJSON(w, http.StatusOK, map[string]any{ + "endpoint": fmt.Sprintf("%s:%d", m.execdHost, m.execdPort), + "headers": map[string]string{"X-EXECD-ACCESS-TOKEN": execdToken(id)}, + }) +} + +func (m *MockServer) sandboxInfo(sb *Sandbox) map[string]any { + return map[string]any{ + "id": sb.ID, + "status": sandboxStatus(sb, time.Now()), + "createdAt": sb.CreatedAt.Format(time.RFC3339), + "expiresAt": sb.ExpiresAt.Format(time.RFC3339), + "entrypoint": []string{"tail", "-f", "/dev/null"}, + } +} + +// ---------- execd handlers ---------- + +func (m *MockServer) handleExecd(w http.ResponseWriter, r *http.Request) { + start := time.Now() + route := "execd.other" + if r.URL.Path == "/ping" { + route = "execd.ping" + } + defer m.recordQps(route, start) + m.stats.incExecdRequests() + token := r.Header.Get("X-EXECD-ACCESS-TOKEN") + id := strings.TrimPrefix(token, "mock-token-") + if token != "" && token == execdToken(id) { + m.mu.RLock() + sb := m.sandboxes[id] + m.mu.RUnlock() + if sb == nil || !sb.running(time.Now()) { + // Not booted yet, expired, or killed. Fail fast and answer with + // a non-retryable status (404): the SDK's retry interceptor + // retries 5xx and transport errors with backoff, which would + // pollute the readiness-poll timing this mock is meant to + // measure. 404 makes the SDK poll at its configured interval + // until ready. The readiness probe only starts paying the route + // latency once the sandbox is actually up. + writeError(w, http.StatusNotFound, "NOT_READY", "sandbox not ready: "+id) + return + } + if sb.Poisoned { + m.stats.incExecdPoisoned() + writeError(w, http.StatusNotFound, "POISONED", "sandbox endpoint poisoned") + return + } + } + // A ready sandbox's requests pay the configured route latency. + m.applyRouteLatency(route) + cfg := m.cfgSnapshot() + if cfg.ExecdFailureRate > 0 && rand.Float64() < cfg.ExecdFailureRate { + m.stats.incExecdFailures() + writeError(w, http.StatusInternalServerError, "INTERNAL_ERROR", "simulated execd failure") + return + } + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(`{"status":"ok"}`)) +} + +// ---------- control handlers ---------- + +func (m *MockServer) handleStats(w http.ResponseWriter, r *http.Request) { + _ = r.Body.Close() + stats := m.stats.snapshot() + m.mu.RLock() + alive := 0 + poisoned := 0 + now := time.Now() + for _, sb := range m.sandboxes { + if sb.alive(now) { + alive++ + } + if sb.Poisoned { + poisoned++ + } + } + m.mu.RUnlock() + writeJSON(w, http.StatusOK, map[string]any{ + "stats": stats, + "alive": alive, + "aliveStats": m.alive.snapshot(time.Now()), + "poisoned": poisoned, + "config": m.cfgSnapshot(), + "qps": m.qps.snapshot(time.Now()), + "serverTime": time.Now().UTC().Format(time.RFC3339), + }) +} + +func (m *MockServer) handleConfig(w http.ResponseWriter, r *http.Request) { + defer r.Body.Close() + var f FaultConfig + if err := json.NewDecoder(r.Body).Decode(&f); err != nil { + writeError(w, http.StatusBadRequest, "INVALID_REQUEST", err.Error()) + return + } + m.cfgMu.Lock() + m.cfg.applyFault(f) + m.cfgMu.Unlock() + if f.PoisonExisting { + m.mu.Lock() + for _, sb := range m.sandboxes { + if sb.alive(time.Now()) { + sb.Poisoned = true + } + } + m.mu.Unlock() + } + if f.PoisonRate != nil { + rate := *f.PoisonRate + m.mu.Lock() + for _, sb := range m.sandboxes { + if sb.alive(time.Now()) && rand.Float64() < rate { + sb.Poisoned = true + } + } + m.mu.Unlock() + } + writeJSON(w, http.StatusOK, map[string]any{"config": m.cfgSnapshot()}) +} + +func (m *MockServer) handleReset(w http.ResponseWriter, r *http.Request) { + _ = r.Body.Close() + m.stats.reset() + m.qps.reset(time.Now()) + m.alive.reset(time.Now()) + writeJSON(w, http.StatusOK, map[string]any{"reset": true}) +} + +// ---------- helpers ---------- + +func (m *MockServer) newSandboxID() string { + m.mu.Lock() + defer m.mu.Unlock() + m.nextSeq++ + return fmt.Sprintf("sbx-mock-%d-%d", time.Now().UnixNano(), m.nextSeq) +} + +func (m *MockServer) cfgSnapshot() Config { + m.cfgMu.RLock() + defer m.cfgMu.RUnlock() + return *m.cfg +} + +func sandboxStatus(sb *Sandbox, now time.Time) map[string]any { + return map[string]any{ + "state": sb.State, + "lastTransitionAt": sb.CreatedAt.Format(time.RFC3339), + } +} + +func execdToken(id string) string { return "mock-token-" + id } + +func writeJSON(w http.ResponseWriter, status int, v any) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(v) +} + +func writeError(w http.ResponseWriter, status int, code, message string) { + writeJSON(w, status, map[string]any{"code": code, "message": message}) +} + +// Stats are the server-side counters exposed via /__stats. They let the +// benchmark driver validate pool behavior (no over-creation, stale cleanup, +// hit rate) and cross-check client-observed numbers. +type Stats struct { + created atomic.Int64 + createFailed atomic.Int64 + killed atomic.Int64 + renews atomic.Int64 + sandboxGets atomic.Int64 + endpointGets atomic.Int64 + execdRequests atomic.Int64 + execdFailures atomic.Int64 + execdPoisoned atomic.Int64 + latencyMu sync.Mutex + createLatencyMs []int64 + maxCreateLatency atomic.Int64 +} + +type StatsSnapshot struct { + Created int64 `json:"created"` + CreateFailed int64 `json:"createFailed"` + Killed int64 `json:"killed"` + Renews int64 `json:"renews"` + SandboxGets int64 `json:"sandboxGets"` + EndpointGets int64 `json:"endpointGets"` + ExecdRequests int64 `json:"execdRequests"` + ExecdFailures int64 `json:"execdFailures"` + ExecdPoisoned int64 `json:"execdPoisoned"` + CreateLatencyMsAvg float64 `json:"createLatencyMsAvg"` + CreateLatencyMsP50 int64 `json:"createLatencyMsP50"` + CreateLatencyMsP95 int64 `json:"createLatencyMsP95"` + CreateLatencyMsP99 int64 `json:"createLatencyMsP99"` + MaxCreateLatencyMs int64 `json:"maxCreateLatencyMs"` +} + +func (s *Stats) incCreated() { s.created.Add(1) } +func (s *Stats) incCreateFailed() { s.createFailed.Add(1) } +func (s *Stats) incKilled() { s.killed.Add(1) } +func (s *Stats) incRenews() { s.renews.Add(1) } +func (s *Stats) incSandboxGets() { s.sandboxGets.Add(1) } +func (s *Stats) incEndpointGets() { s.endpointGets.Add(1) } +func (s *Stats) incExecdRequests() { s.execdRequests.Add(1) } +func (s *Stats) incExecdFailures() { s.execdFailures.Add(1) } +func (s *Stats) incExecdPoisoned() { s.execdPoisoned.Add(1) } + +func (s *Stats) recordCreateLatency(d time.Duration) { + ms := d.Milliseconds() + s.latencyMu.Lock() + s.createLatencyMs = append(s.createLatencyMs, ms) + s.latencyMu.Unlock() + if cur := s.maxCreateLatency.Load(); ms > cur { + s.maxCreateLatency.CompareAndSwap(cur, ms) + } +} + +func (s *Stats) snapshot() StatsSnapshot { + s.latencyMu.Lock() + samples := append([]int64(nil), s.createLatencyMs...) + s.latencyMu.Unlock() + sorted := make([]int64, len(samples)) + copy(sorted, samples) + // insertion sort: sample counts stay small for benchmark runs + for i := 1; i < len(sorted); i++ { + for j := i; j > 0 && sorted[j] < sorted[j-1]; j-- { + sorted[j], sorted[j-1] = sorted[j-1], sorted[j] + } + } + percentile := func(p float64) int64 { + if len(sorted) == 0 { + return 0 + } + idx := int(float64(len(sorted)-1) * p) + return sorted[idx] + } + var sum int64 + for _, v := range sorted { + sum += v + } + avg := 0.0 + if len(sorted) > 0 { + avg = float64(sum) / float64(len(sorted)) + } + return StatsSnapshot{ + Created: s.created.Load(), + CreateFailed: s.createFailed.Load(), + Killed: s.killed.Load(), + Renews: s.renews.Load(), + SandboxGets: s.sandboxGets.Load(), + EndpointGets: s.endpointGets.Load(), + ExecdRequests: s.execdRequests.Load(), + ExecdFailures: s.execdFailures.Load(), + ExecdPoisoned: s.execdPoisoned.Load(), + CreateLatencyMsAvg: avg, + CreateLatencyMsP50: percentile(0.50), + CreateLatencyMsP95: percentile(0.95), + CreateLatencyMsP99: percentile(0.99), + MaxCreateLatencyMs: s.maxCreateLatency.Load(), + } +} + +func (s *Stats) reset() { + s.created.Store(0) + s.createFailed.Store(0) + s.killed.Store(0) + s.renews.Store(0) + s.sandboxGets.Store(0) + s.endpointGets.Store(0) + s.execdRequests.Store(0) + s.execdFailures.Store(0) + s.execdPoisoned.Store(0) + s.maxCreateLatency.Store(0) + s.latencyMu.Lock() + s.createLatencyMs = nil + s.latencyMu.Unlock() +} diff --git a/tests/benchmark/mockserver/stats.go b/tests/benchmark/mockserver/stats.go new file mode 100644 index 000000000..79af61437 --- /dev/null +++ b/tests/benchmark/mockserver/stats.go @@ -0,0 +1,289 @@ +/* + * Copyright 2026 Alibaba Group Holding Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package main + +import ( + "sync" + "time" +) + +// DefaultStatsWindowSec is how many seconds of per-API request history the +// mock keeps for QPS analysis. 30 minutes covers a full default benchmark +// run; raise it with -stats-window-sec for longer soaks. +const DefaultStatsWindowSec = 1800 + +// qpsTracker records exact per-second request counts for one API route in a +// ring buffer. Snapshots expose totals, recent-window rates, and the full +// per-second series so the driver can plot pool load over time. +type qpsTracker struct { + mu sync.Mutex + window int64 + buckets []int64 // ring of per-second counts + secs []int64 // wall-clock second each bucket covers + startSec int64 // first recorded second (0 until first request/reset) + total int64 + latencySumMs int64 + maxMs int64 +} + +func newQpsTracker(windowSec int) *qpsTracker { + return &qpsTracker{ + window: int64(windowSec), + buckets: make([]int64, windowSec), + secs: make([]int64, windowSec), + } +} + +func (t *qpsTracker) record(now time.Time, elapsed time.Duration) { + sec := now.Unix() + idx := sec % t.window + t.mu.Lock() + if t.secs[idx] != sec { + // Lazy bucket reset: first request in this second. + t.secs[idx] = sec + t.buckets[idx] = 0 + } + t.buckets[idx]++ + t.total++ + if t.startSec == 0 { + t.startSec = sec + } + ms := elapsed.Milliseconds() + t.latencySumMs += ms + if ms > t.maxMs { + t.maxMs = ms + } + t.mu.Unlock() +} + +func (t *qpsTracker) reset(now time.Time) { + t.mu.Lock() + defer t.mu.Unlock() + for i := range t.buckets { + t.buckets[i] = 0 + t.secs[i] = 0 + } + t.total = 0 + t.latencySumMs = 0 + t.maxMs = 0 + t.startSec = now.Unix() +} + +type QpsSnapshot struct { + Total int64 `json:"total"` + Qps1s float64 `json:"qps1s"` + Qps5s float64 `json:"qps5s"` + Qps60s float64 `json:"qps60s"` + SeriesStart int64 `json:"seriesStartUnixSec"` + Series []int64 `json:"series"` + // AvgMs is the average handler latency for this route since the last reset. + AvgMs float64 `json:"avgMs"` + MaxMs int64 `json:"maxMs"` +} + +// snapshot returns totals, rates over the trailing 1s/5s/60s windows (the +// current, possibly partial second counts toward qps1s), and the per-second +// series covering the retained window. Requests older than the ring are +// dropped, so for runs longer than the window the driver should poll /__stats +// and accumulate externally. +func (t *qpsTracker) snapshot(now time.Time) QpsSnapshot { + sec := now.Unix() + t.mu.Lock() + defer t.mu.Unlock() + + if t.startSec == 0 { + t.startSec = sec + } + begin := t.startSec + if sec-begin+1 > t.window { + begin = sec - t.window + 1 + } + + avgLatency := 0.0 + if t.total > 0 { + avgLatency = float64(t.latencySumMs) / float64(t.total) + } + + series := make([]int64, 0, sec-begin+1) + for s := begin; s <= sec; s++ { + idx := s % t.window + count := t.buckets[idx] + if t.secs[idx] != s { + count = 0 + } + series = append(series, count) + } + + rate := func(windowSec int64) float64 { + if windowSec <= 0 { + return 0 + } + start := sec - windowSec + 1 + if start < begin { + start = begin + } + if sec < start { + return 0 + } + var sum int64 + for s := start; s <= sec; s++ { + idx := s % t.window + if t.secs[idx] == s { + sum += t.buckets[idx] + } + } + return float64(sum) / float64(sec-start+1) + } + + return QpsSnapshot{ + Total: t.total, + Qps1s: rate(1), + Qps5s: rate(5), + Qps60s: rate(60), + SeriesStart: begin, + Series: series, + AvgMs: avgLatency, + MaxMs: t.maxMs, + } +} + +// QpsRegistry tracks one tracker per API route. +type QpsRegistry struct { + windowSec int + mu sync.RWMutex + trackers map[string]*qpsTracker +} + +func newQpsRegistry(windowSec int) *QpsRegistry { + return &QpsRegistry{ + windowSec: windowSec, + trackers: make(map[string]*qpsTracker), + } +} + +func (r *QpsRegistry) record(route string, now time.Time, elapsed time.Duration) { + r.mu.RLock() + t := r.trackers[route] + r.mu.RUnlock() + if t == nil { + r.mu.Lock() + t = r.trackers[route] + if t == nil { + t = newQpsTracker(r.windowSec) + r.trackers[route] = t + } + r.mu.Unlock() + } + t.record(now, elapsed) +} + +func (r *QpsRegistry) reset(now time.Time) { + r.mu.RLock() + defer r.mu.RUnlock() + for _, t := range r.trackers { + t.reset(now) + } +} + +func (r *QpsRegistry) snapshot(now time.Time) map[string]QpsSnapshot { + r.mu.RLock() + defer r.mu.RUnlock() + out := make(map[string]QpsSnapshot, len(r.trackers)) + for route, t := range r.trackers { + out[route] = t.snapshot(now) + } + return out +} + +// GaugeSnapshot is the per-second view of a gauge (e.g. alive sandboxes). +type GaugeSnapshot struct { + Max int64 `json:"max"` + SeriesStart int64 `json:"seriesStartUnixSec"` + Series []int64 `json:"series"` +} + +// gaugeTracker keeps the per-second peak of a gauge (ring buffer) plus the +// all-time max since the last reset. +type gaugeTracker struct { + mu sync.Mutex + window int64 + buckets []int64 + secs []int64 + startSec int64 + maxAll int64 +} + +func newGaugeTracker(windowSec int) *gaugeTracker { + return &gaugeTracker{ + window: int64(windowSec), + buckets: make([]int64, windowSec), + secs: make([]int64, windowSec), + } +} + +func (g *gaugeTracker) record(now time.Time, value int64) { + sec := now.Unix() + idx := sec % g.window + g.mu.Lock() + if g.secs[idx] != sec { + g.secs[idx] = sec + g.buckets[idx] = value + } else if value > g.buckets[idx] { + g.buckets[idx] = value + } + if value > g.maxAll { + g.maxAll = value + } + if g.startSec == 0 { + g.startSec = sec + } + g.mu.Unlock() +} + +func (g *gaugeTracker) reset(now time.Time) { + g.mu.Lock() + defer g.mu.Unlock() + for i := range g.buckets { + g.buckets[i] = 0 + g.secs[i] = 0 + } + g.maxAll = 0 + g.startSec = now.Unix() +} + +func (g *gaugeTracker) snapshot(now time.Time) GaugeSnapshot { + sec := now.Unix() + g.mu.Lock() + defer g.mu.Unlock() + if g.startSec == 0 { + g.startSec = sec + } + begin := g.startSec + if sec-begin+1 > g.window { + begin = sec - g.window + 1 + } + series := make([]int64, 0, sec-begin+1) + for s := begin; s <= sec; s++ { + idx := s % g.window + value := g.buckets[idx] + if g.secs[idx] != s { + value = 0 + } + series = append(series, value) + } + return GaugeSnapshot{Max: g.maxAll, SeriesStart: begin, Series: series} +} diff --git a/tests/benchmark/run.sh b/tests/benchmark/run.sh new file mode 100755 index 000000000..b73fdde5c --- /dev/null +++ b/tests/benchmark/run.sh @@ -0,0 +1,126 @@ +#!/usr/bin/env bash +# Copyright 2026 Alibaba Group Holding Ltd. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +set -euo pipefail + +BENCH_DIR="$(cd "$(dirname "$0")" && pwd)" +REPO_ROOT="$(cd "${BENCH_DIR}/../.." && pwd)" + +# Gradle 9 requires JVM 17+. Override JAVA_HOME when it is unset or points at +# an older JDK; otherwise pick the highest installed major >= 17. +NEEDS_OVERRIDE=false +if [ -z "${JAVA_HOME:-}" ]; then + NEEDS_OVERRIDE=true +elif [ -x "${JAVA_HOME}/bin/java" ]; then + MAJOR="$("${JAVA_HOME}/bin/java" -version 2>&1 | sed -nE 's/.*version "([0-9]+).*/\1/p' | head -1)" + case "${MAJOR}" in + 8|9|10|11|12|13|14|15|16) NEEDS_OVERRIDE=true ;; + esac +fi +if [ "${NEEDS_OVERRIDE}" = "true" ] && command -v /usr/libexec/java_home > /dev/null 2>&1; then + # Highest installed major >= 17, e.g. "17.0.7" or "21.0.4". + VERSION="$( + /usr/libexec/java_home -V 2>&1 \ + | sed -nE 's/^[[:space:]]*([0-9]+)\.([0-9]+).*/\1.\2/p' \ + | awk -F. '$1 >= 17' \ + | sort -t. -k1,1n -k2,2n | tail -1 + )" + if [ -n "${VERSION}" ]; then + JH="$(/usr/libexec/java_home -v "${VERSION}" 2>/dev/null || true)" + if [ -n "${JH}" ]; then + export JAVA_HOME="${JH}" + fi + fi +fi + +LIFECYCLE_ADDR="${LIFECYCLE_ADDR:-127.0.0.1:18080}" +EXECD_ADDR="${EXECD_ADDR:-127.0.0.1:18081}" +MOCK_CONFIG="${MOCK_CONFIG:-${BENCH_DIR}/configs/default.json}" +MOCK_PID="" + +cleanup() { + if [ -n "${MOCK_PID}" ]; then + kill "${MOCK_PID}" 2>/dev/null || true + wait "${MOCK_PID}" 2>/dev/null || true + fi +} +trap cleanup EXIT + +usage() { + cat <<'EOF' +Usage: run.sh [--mock-config ] [-- ] + +Environment: + LIFECYCLE_ADDR lifecycle mock listen address (default 127.0.0.1:18080) + EXECD_ADDR execd mock listen address (default 127.0.0.1:18081) + MOCK_CONFIG mock server config JSON (default configs/default.json) + +Driver args (after --) are forwarded to the benchmark driver, e.g.: + ./run.sh -- --max-idle 50 --scenarios warm-latency,steady-state +EOF +} + +# parse flags +DRIVER_ARGS=() +while [ $# -gt 0 ]; do + case "$1" in + --mock-config) shift; MOCK_CONFIG="$1" ;; + --) shift; DRIVER_ARGS=("$@"); break ;; + -h|--help) usage; exit 0 ;; + *) DRIVER_ARGS=("$@"); break ;; + esac + shift +done + +# 1. build the mock server +echo "== building mock server ==" +(cd "${BENCH_DIR}/mockserver" && go build -o "${BENCH_DIR}/bin/mockserver" .) + +# 2. one run directory holds every artifact of this run (reports, CSVs, +# mock config, driver args, mock log) for convenient offline analysis. +RUN_DIR="${BENCH_DIR}/results/run-$(date +%Y%m%d-%H%M%S)" +mkdir -p "${RUN_DIR}" +printf '%s\n' "${DRIVER_ARGS[@]}" > "${RUN_DIR}/driver-args.txt" +cp "${MOCK_CONFIG}" "${RUN_DIR}/mock-config.json" + +# 3. start the mock server +echo "== starting mock server (lifecycle=${LIFECYCLE_ADDR}, execd=${EXECD_ADDR}, config=${MOCK_CONFIG}) ==" +echo "== run directory: ${RUN_DIR} ==" +"${BENCH_DIR}/bin/mockserver" \ + -lifecycle-addr "${LIFECYCLE_ADDR}" \ + -execd-addr "${EXECD_ADDR}" \ + -config "${MOCK_CONFIG}" \ + > "${RUN_DIR}/mockserver.log" 2>&1 & +MOCK_PID=$! + +for _ in $(seq 1 50); do + if curl -fsS "http://${LIFECYCLE_ADDR}/__stats" > /dev/null 2>&1; then + break + fi + sleep 0.2 +done +if ! curl -fsS "http://${LIFECYCLE_ADDR}/__stats" > /dev/null 2>&1; then + echo "error: mock server did not come up" >&2 + cat "${RUN_DIR}/mockserver.log" >&2 + exit 1 +fi + +# 4. run the driver (Kotlin SDK is built from source via composite build, +# see kotlin/settings.gradle.kts) +echo "== running benchmark driver ==" +DRIVER_ARGS+=("--report-dir" "${RUN_DIR}") +(cd "${BENCH_DIR}/kotlin" && ./gradlew --console=plain run --args="${DRIVER_ARGS[*]}") + +echo "== done: ${RUN_DIR} =="