Files
EvoScientist/scripts/benchmark_usage_spool.py

128 lines
3.9 KiB
Python

"""Measure durable UsageEvent enqueue latency for the release platform matrix."""
from __future__ import annotations
import json
import os
import platform
import tempfile
import time
from concurrent.futures import ThreadPoolExecutor
from datetime import UTC, datetime
from pathlib import Path
from EvoScientist.usage.schema import UsageEventV1
from EvoScientist.usage.spool import UsageSpool
WORKER_COUNTS = (1, 4, 16)
SAMPLES_PER_WORKER = 32
P95_LIMIT_MS = 20.0
P99_LIMIT_MS = 100.0
class BenchmarkSpool(UsageSpool):
def _run(self) -> None:
self._stop.wait()
def _percentile(samples: list[float], percentile: float) -> float:
ordered = sorted(samples)
index = min(len(ordered) - 1, max(0, int(len(ordered) * percentile) - 1))
return ordered[index]
def _event(index: int) -> UsageEventV1:
model_call_id = f"benchmark-{index:08d}"
now = datetime.now(UTC)
deployment_id = "11111111-1111-4111-8111-111111111111"
return UsageEventV1(
schema_version=1,
event_id=f"{deployment_id}:{model_call_id}:callback_final:1",
event_type="usage_observed",
source="callback_final",
authority_class="observed_final",
revision=1,
deployment_id=deployment_id,
workspace_id="ws1_benchmark",
model_call_id=model_call_id,
parent_run_id=None,
provider_request_id=None,
thread_id="benchmark-thread",
source_session_id=None,
turn_id="benchmark-turn",
workspace_dir=None,
scope="main",
source_agent="EvoScientist",
provider_profile_id="benchmark",
provider_revision=None,
provider_adapter="openai",
model_alias="benchmark",
upstream_model_id="benchmark",
usage_status="confirmed",
input_tokens=1,
output_tokens=1,
provider_total_tokens=2,
input_token_details={},
output_token_details={},
started_at=now,
observed_at=now,
completed_at=now,
)
def _run_case(root: Path, workers: int) -> dict[str, float | int]:
spool_root = root / str(workers)
os.environ.update(
{
"EVOSCIENTIST_USAGE_SINK_URL": "http://127.0.0.1:1/api/usage/events",
"EVOSCIENTIST_USAGE_SINK_TOKEN": "benchmark",
"EVOSCIENTIST_DEPLOYMENT_ID": "11111111-1111-4111-8111-111111111111",
"EVOSCIENTIST_WORKSPACE_ID": "ws1_benchmark",
"EVOSCIENTIST_USAGE_SPOOL_DIR": str(spool_root),
"EVOSCIENTIST_USAGE_SPOOL_MAX_FILES": "100000",
}
)
spool = BenchmarkSpool()
def enqueue(index: int) -> float:
started = time.perf_counter_ns()
spool.enqueue(_event(workers * 1_000_000 + index))
return (time.perf_counter_ns() - started) / 1_000_000
try:
with ThreadPoolExecutor(max_workers=workers) as pool:
samples = list(pool.map(enqueue, range(workers * SAMPLES_PER_WORKER)))
finally:
spool.close()
return {
"workers": workers,
"samples": len(samples),
"p50_ms": round(_percentile(samples, 0.50), 3),
"p95_ms": round(_percentile(samples, 0.95), 3),
"p99_ms": round(_percentile(samples, 0.99), 3),
"max_ms": round(max(samples), 3),
}
def main() -> None:
with tempfile.TemporaryDirectory(prefix="evosci-usage-benchmark-") as temporary:
results = [_run_case(Path(temporary), workers) for workers in WORKER_COUNTS]
report = {
"platform": platform.platform(),
"python": platform.python_version(),
"limits_ms": {"p95": P95_LIMIT_MS, "p99": P99_LIMIT_MS},
"results": results,
}
print(json.dumps(report, indent=2, sort_keys=True))
failed = [
result
for result in results
if result["p95_ms"] > P95_LIMIT_MS or result["p99_ms"] > P99_LIMIT_MS
]
if failed:
raise SystemExit("usage spool latency exceeds the release threshold")
if __name__ == "__main__":
main()