"""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()