fix(harness): bound benchmark provenance memory and document measurement limits

This commit is contained in:
Yujong Lee 2026-09-05 13:22:44 -07:00
parent e9dbe047a2
commit 75ae4d591a
5 changed files with 74 additions and 100 deletions

View file

@ -1 +1,33 @@
Measures Python and Rust SDK latency, CPU time, and process memory against deterministic local provider replays, outside correctness-check overhead
# What this is
Measures Python and Rust SDK latency, process CPU time, and sampled RSS against local provider replays. Initial coverage is sync/async Mistral OCR at concurrency one. Other SDK functions remain explicitly unimplemented. Run locally, with no CI integration
# How it works
Derive five profiles from the existing `e2e_parity` recording without editing it. `small` uses a 32 KiB PDF and one response page. `request_medium` and `request_large` increase only the PDF to 256 KiB and 2 MiB. `response_medium` and `response_large` increase only the response to 16 and 128 pages. Insert PDF comment padding before the original final `startxref` and EOF trailer, preserving object offsets. Response pages are synthetic repetitions, independent of actual PDF content
Generate fixtures in the controller, outside SDK worker imports. Each backend, route, profile, and repeat gets fresh timing and memory workers against a separate local HTTP provider process. Run Python and Rust sequentially, reversing their order on alternating repeats. The provider drains request bytes and serves preloaded responses without JSON parsing or capture. Its CPU and RSS are excluded
Warmup, preflight response checks, bounded-memory native hashing, and garbage collection precede readiness. Require matching Python/Rust preflight digests and matching timing/memory worker digests. Every request checks the parity harness's User-Agent convention to reject Python fallback during a Rust run. Use `e2e_parity` for complete request and response semantics
Time each SDK call until its returned result is discarded. Async calls share a persistent event loop. Process CPU and batch elapsed time include sample collection overhead and work on Python/native threads. Startup, fixture loading, warmup, preflight, final garbage collection, and reporting are excluded. Deferred callbacks can outlive the measured batch
Sample only the SDK worker's RSS in the separate memory pass. Baseline follows warmup and garbage collection. Peak includes baseline, periodic samples, and the final sample. After RSS follows another garbage collection with input/client state resident. RSS includes native allocations and shared pages. Sampling can miss short peaks, and retained RSS does not prove a leak
Publish readiness/results atomically in temporary JSON files, reserving stdout/stderr for diagnostics. Bound readiness, measurement, and shutdown waits. Terminate and reap workers on failure or interruption. Reject missing extensions, backend mismatches, exceptions, and incomplete samples. Preserve completed pairs in partial reports when a worker fails
Report pooled p50/p95/p99 latency, CPU milliseconds per call, sequential calls per second, baseline/peak/after RSS, and Python p50 divided by backend p50. JSON retains raw per-repeat samples, fixture/extension hashes, Python version, options, platform, Git revision, and working-tree state. A pass confirms valid measurements without imposing performance thresholds. Use an idle host, inspect repeat variation, and treat short-run tail estimates cautiously. Loopback HTTP, allocator behavior, and deferred work affect results. Streaming, gateway overhead, live provider latency, and concurrent throughput are outside scope
Editable `uv sync` builds a development extension. Build release explicitly and retain `--no-sync`:
```sh
uv sync --frozen --python 3.12
VIRTUAL_ENV="$PWD/.venv" uvx --from maturin==1.15.0 maturin develop --release
uv run --no-sync python -m tests.rust-python-harness run e2e_benchmark \
--surface sdk --function ocr \
--benchmark-arg=--output=/tmp/e2e-benchmark.json
```
Forward each option through `--benchmark-arg=...`. Defaults: `--iterations=100`, `--warmup=10`, `--repeats=3`, `--timeout=120` seconds, `--sample-interval-ms=5`. Repeat `--profile=NAME` or `--route=ocr|aocr` to select subsets. `--output=PATH` exports JSON. Invalid values and unwritable destinations produce handled CLI errors. `run all` includes this strategy. A smoke run can select `small`, `ocr`, 10 iterations, 2 warmups, and 1 repeat
Run focused checks with `uv run --no-sync pytest -o consider_namespace_packages=true tests/rust-python-harness/strategies/e2e_benchmark tests/rust-python-harness/cli -q`

View file

@ -1,87 +0,0 @@
# End-to-end SDK benchmark
Compare `LITELLM_RUST=0` and `LITELLM_RUST=1` against a separate local HTTP provider process, with no real provider calls, credentials, or Docker required
The initial workload covers synchronous and asynchronous Mistral OCR using the existing `e2e_parity` recording. Other SDK functions are explicitly unimplemented. This measures SDK calls including loopback HTTP transport and response construction. It does not measure gateway overhead, streaming, or concurrent load
## Run
Build the Rust extension in release mode first. An editable `uv sync` normally builds the development profile, which is unsuitable for a Python/Rust performance comparison
```sh
uv sync --frozen --python 3.12
VIRTUAL_ENV="$PWD/.venv" uvx --from maturin==1.15.0 maturin develop --release
uv run --no-sync python -m tests.rust-python-harness run e2e_benchmark \
--surface sdk --function ocr \
--benchmark-arg=--output=/tmp/e2e-benchmark.json
```
Keep `--no-sync` on the benchmark command so it uses the extension you just built. Use an otherwise idle machine and run the same command on both revisions when evaluating a change
For a short smoke run:
```sh
uv run --no-sync python -m tests.rust-python-harness run e2e_benchmark \
--function ocr \
--benchmark-arg=--profile=small \
--benchmark-arg=--route=ocr \
--benchmark-arg=--iterations=10 \
--benchmark-arg=--warmup=2 \
--benchmark-arg=--repeats=1 \
--benchmark-arg=--output=/tmp/e2e-benchmark-smoke.json
```
`run all` also runs this strategy with its defaults. No CI integration is added
## Workloads
The seed cassette stays under `e2e_parity/sdk/ocr/fixtures/data`. The benchmark derives synthetic size variants in memory; it never edits or re-records the correctness fixtures
| Profile | Inline PDF bytes | Response pages |
| --- | ---: | ---: |
| small | 32 KiB | 1 |
| request_medium | 256 KiB | 1 |
| request_large | 2 MiB | 1 |
| response_medium | 32 KiB | 16 |
| response_large | 32 KiB | 128 |
Request variants add PDF comment padding before the final `startxref` marker, preserving existing object offsets and the EOF trailer. The SDK sends base64 plus JSON framing, so wire request sizes exceed the document sizes above. Response variants repeat recorded pages with contiguous indexes and adjusted usage. They exercise realistic response structure, but their page count intentionally varies independently of the input PDF's content
## Measurements
Each backend, route, size, and repeat gets a fresh SDK process for timing and another for memory. Python and Rust execute sequentially, with their order reversed on alternating repeats. The local provider serves preloaded bytes without parsing or capturing request JSON. Its CPU and RSS are outside the SDK measurements
Workers warm up their clients and run an untimed response check before measuring. Python and Rust response digests must match. Every provider request also checks the existing parity harness's User-Agent convention: a Rust run using Python's HTTP path fails instead of reporting a comparison between two Python runs. Missing native extensions, SDK exceptions, timeouts, and incomplete samples fail the run. Workers publish readiness and results atomically in temporary JSON files; stdout and stderr go to a diagnostic log. The controller waits for timing workers to exit and samples RSS only for memory workers. A timeout or interruption terminates and reaps the SDK worker, with bounded shutdown waits for both SDK and provider processes
Latency starts immediately before the SDK call and ends when its result has been returned and discarded. Async calls are awaited on a persistent event loop. CPU is process CPU time during the timed batch, including Python and native threads. Fixture loading, process startup, warmup, preflight serialization, and report generation are excluded. Default SDK behavior is retained, so deferred background work can extend beyond a call's return; these metrics describe the measurement window, not the eventual cost of every callback
The memory controller uses `psutil` to sample only the SDK worker's RSS during a separate run, avoiding polling overhead in latency results. Baseline RSS is taken after warmup and garbage collection. Peak is the highest sampled RSS, including the baseline and final sample. After RSS is measured after the workload and another garbage collection, with input/client state still resident. RSS includes native allocations and shared resident pages, so it is not equivalent to Python heap size or uniquely owned memory. Sampling can miss brief peaks; these values are not an exact allocator high-water mark or proof of a leak
The terminal reports pooled p50/p95/p99 latency, CPU milliseconds per call, sequential calls per second, baseline/peak/after RSS, and speedup (`Python p50 / backend p50`). Throughput is at concurrency one, not saturation capacity. Short runs cannot estimate tail latency reliably. The JSON retains each repeat's raw latency samples, CPU and memory measurements, input/response sizes, seed hash, Python version, native extension hash, settings, platform, Git revision, and whether the working tree has changes
## Options
Pass each option through `--benchmark-arg=...`
| Option | Default | Meaning |
| --- | --- | --- |
| `--iterations=N` | 100 | Measured calls per worker |
| `--warmup=N` | 10 | Warmup calls, followed by one preflight call |
| `--repeats=N` | 3 | Fresh paired runs per workload |
| `--profile=NAME` | All five | Select a size profile; repeat for several |
| `--route=ocr` or `--route=aocr` | Both | Select SDK entrypoint; repeat for both |
| `--timeout=SECONDS` | 120 | Worker readiness and measurement deadline |
| `--sample-interval-ms=N` | 5 | Memory sampling interval, at least 1 ms |
| `--output=PATH` | None | Write a JSON report, including partial results on worker failure |
Run the strategy tests and existing harness checks with:
```sh
uv run --no-sync pytest -o consider_namespace_packages=true \
tests/rust-python-harness/strategies/e2e_benchmark \
tests/rust-python-harness/shared tests/rust-python-harness/cli \
tests/rust-python-harness/strategies/unit_tests_mapping \
tests/rust-python-harness/strategies/unit_tests_parity \
tests/rust-python-harness/strategies/unit_tests_rust \
tests/test_rust_python_harness.py -q
```

View file

@ -38,6 +38,6 @@ STRATEGY: Final = StrategyDefinition(
surfaces=("sdk",),
runner_argument=RunnerArgumentDefinition(
option="--benchmark-arg",
help="benchmark option, e.g. --benchmark-arg=--iterations=100; see the strategy README",
help="benchmark option, e.g. --benchmark-arg=--iterations=100; see the strategy AGENTS.md",
),
)

View file

@ -2,8 +2,10 @@ from __future__ import annotations
import asyncio
import base64
import hashlib
import subprocess
import sys
import tracemalloc
from pathlib import Path
from time import monotonic, sleep
from typing import Final
@ -18,11 +20,11 @@ from litellm.llms.base_llm.ocr.transformation import OCRResponse
from ...cli import main
from .execution import execute_phase, sdk_process, wait_for_output
from .constants import PYTHON_SENTINEL
from .models import Invocation, Options
from .models import Invocation, Options, Route
from .provider import provider_process
from .reporting import percentile, render_measurements
from .runner import Report, parse_options
from .worker import measure_async, measure_sync
from .worker import file_sha256, measure_async, measure_sync
from .workloads import JSON_OBJECT, JSON_PAGES, ocr_workload, padded_pdf
REPO_ROOT: Final = Path(__file__).resolve().parents[4]
@ -162,10 +164,13 @@ def test_replay_rejects_python_fallback_during_rust_measurement() -> None:
assert "backend mismatch" in response.text
def test_worker_errors_are_reported_instead_of_counted_as_fast_calls() -> None:
@pytest.mark.parametrize("route", ("ocr", "aocr"))
def test_worker_errors_are_reported_instead_of_counted_as_fast_calls(route: Route) -> None:
workload: Final = ocr_workload("small")
with provider_process(workload.response, "rust") as url:
request: Final = invocation().model_copy(update={"provider_url": url, "document_url": workload.document_url})
request: Final = invocation().model_copy(
update={"provider_url": url, "document_url": workload.document_url, "route": route}
)
with pytest.raises(RuntimeError, match="backend mismatch"):
execute_phase(request, "python", Options(iterations=3, warmup=1), REPO_ROOT)
@ -184,7 +189,7 @@ def test_cli_runs_both_backends_and_exports_measurements(tmp_path: Path, capsys:
"--benchmark-arg=--route=aocr",
"--benchmark-arg=--iterations=3",
"--benchmark-arg=--warmup=1",
"--benchmark-arg=--repeats=1",
"--benchmark-arg=--repeats=2",
f"--benchmark-arg=--output={output}",
)
)
@ -192,7 +197,12 @@ def test_cli_runs_both_backends_and_exports_measurements(tmp_path: Path, capsys:
assert exit_code == 0, captured.out + captured.err
assert "Result: PASSED" in captured.out
report: Final = Report.model_validate_json(output.read_bytes())
assert {value.backend for value in report.measurements} == {"python", "rust"}
assert tuple((value.repeat, value.backend) for value in report.measurements) == (
(0, "python"),
(0, "rust"),
(1, "rust"),
(1, "python"),
)
assert len({value.ready.response_digest for value in report.measurements}) == 1
for value in report.measurements:
assert len(value.timing.latency_ms) == 3
@ -242,3 +252,17 @@ def test_cli_reports_unsupported_functions_without_measurements(
assert "not implemented" in captured.out
report: Final = Report.model_validate_json(output.read_bytes())
assert report.measurements == ()
def test_native_provenance_hash_uses_bounded_memory(tmp_path: Path) -> None:
native: Final = tmp_path / "native.so"
payload: Final = b"native-extension-data" * (1024 * 1024)
native.write_bytes(payload)
expected: Final = hashlib.sha256(payload).hexdigest()
tracemalloc.start()
try:
assert file_sha256(native) == expected
_, peak = tracemalloc.get_traced_memory()
assert peak < 1024 * 1024
finally:
tracemalloc.stop()

View file

@ -16,6 +16,11 @@ from litellm.llms.base_llm.ocr.transformation import OCRResponse
from .models import Invocation, Ready, Timing
def file_sha256(path: Path) -> str:
with path.open("rb") as source:
return hashlib.file_digest(source, "sha256").hexdigest()
def _ready(response: OCRResponse) -> Ready:
from litellm.rust_bridge import get_native_bridge
from litellm.rust_bridge.configuration import rust_enabled
@ -27,7 +32,7 @@ def _ready(response: OCRResponse) -> Ready:
return Ready(
response_digest=hashlib.sha256(json.dumps(response.model_dump(), sort_keys=True).encode()).hexdigest(),
python_version=platform.python_version(),
native_sha256=hashlib.sha256(Path(native_path).read_bytes()).hexdigest() if native_path else None,
native_sha256=file_sha256(Path(native_path)) if native_path else None,
)
@ -63,24 +68,24 @@ async def _async_sample(call: Callable[[], Awaitable[OCRResponse]]) -> float:
def measure_sync(call: Callable[[], OCRResponse], invocation: Invocation) -> Timing:
cpu_start: Final = process_time_ns()
wall_start: Final = perf_counter_ns()
if invocation.phase == "memory":
for _ in range(invocation.iterations):
call()
return Timing(latency_ms=(), cpu_ms=0, elapsed_ms=0)
cpu_start: Final = process_time_ns()
wall_start: Final = perf_counter_ns()
samples: Final = tuple(_sync_sample(call) for _ in range(invocation.iterations))
elapsed: Final = perf_counter_ns() - wall_start
return Timing(latency_ms=samples, cpu_ms=(process_time_ns() - cpu_start) / 1e6, elapsed_ms=elapsed / 1e6)
async def measure_async(call: Callable[[], Awaitable[OCRResponse]], invocation: Invocation) -> Timing:
cpu_start: Final = process_time_ns()
wall_start: Final = perf_counter_ns()
if invocation.phase == "memory":
for _ in range(invocation.iterations):
await call()
return Timing(latency_ms=(), cpu_ms=0, elapsed_ms=0)
cpu_start: Final = process_time_ns()
wall_start: Final = perf_counter_ns()
samples: Final = tuple([await _async_sample(call) for _ in range(invocation.iterations)])
elapsed: Final = perf_counter_ns() - wall_start
return Timing(latency_ms=samples, cpu_ms=(process_time_ns() - cpu_start) / 1e6, elapsed_ms=elapsed / 1e6)