Source code for stormlog.infer.cache_state

"""Requested and verified prefix-cache state for each workload case.

A benchmark can ask for a cold cache and can call a reset endpoint before
each case, such as vLLM's ``/reset_prefix_cache`` (available when the server
runs in development mode) or SGLang's ``/flush_cache``. Asking is not
proof: until an engine adapter can read the cache's contents, every case
records its cache state as unverified and says why.

A 2xx answer is not proof of a reset either. vLLM answers HTTP 200 with
``{"success": false}`` while blocks are still held, so the body is read: a
reset is acknowledged only when the server says ``success: true``, refused
when it keeps saying anything else, and accepted but unconfirmed when a 2xx
body has no ``success`` field.
"""

from __future__ import annotations

import json
import time
import urllib.error
import urllib.request
from dataclasses import dataclass
from typing import Any

# Re-exported: callers imported redact_url from here before it moved.
from ..scrub import redact_url
from .openai_client import inference_opener

UNSPECIFIED = "unspecified"
COLD = "cold"
CACHE_STATES = (UNSPECIFIED, COLD)
UNVERIFIED = "unverified"

ACKNOWLEDGED = "acknowledged"
REFUSED = "refused"
ACCEPTED_UNVERIFIED = "accepted_unverified"
# vLLM refuses a reset while blocks are held and says callers may retry.
RESET_RETRY_SECONDS = 10.0
RESET_RETRY_INTERVAL_SECONDS = 0.5
_RESET_BODY_LIMIT = 64 * 1024


[docs] @dataclass(frozen=True) class CacheReset: """The outcome of a cache reset, over every attempt it took. ``success`` is True when the body said ``"success": true``, False when it had a ``success`` field with any other value, and None without one. ``at_ns`` is when the first attempt was sent; ``answered_at_ns`` is when the answer recorded here, the last attempt's, came back. """ url: str at_ns: int status: int | None = None error: str | None = None success: bool | None = None attempts: int = 1 answered_at_ns: int | None = None @property def answer(self) -> str | None: """``acknowledged``, ``refused`` or ``accepted_unverified``; None if no 2xx.""" if self.status is None or not 200 <= self.status < 300: return None if self.success is None: return ACCEPTED_UNVERIFIED return ACKNOWLEDGED if self.success else REFUSED @property def succeeded(self) -> bool: """The server accepted the reset, whether or not it confirmed it.""" return self.answer in (ACKNOWLEDGED, ACCEPTED_UNVERIFIED) @property def acknowledged(self) -> bool: """The server said the reset happened.""" return self.answer == ACKNOWLEDGED
[docs] def to_record(self) -> dict[str, Any]: return { "url": self.url, "at_ns": self.at_ns, "status": self.status, "error": self.error, "success": self.success, "answer": self.answer, "attempts": self.attempts, "answered_at_ns": self.answered_at_ns, }
[docs] def reset_cache( url: str, *, timeout_seconds: float, api_key: str | None = None, retry_seconds: float = RESET_RETRY_SECONDS, ) -> CacheReset: """POST to a reset endpoint and record what happened; never raises. A reset the server refuses (``success: false``) is tried again every half second for up to ``retry_seconds``. No attempt starts after that, though the last one may take up to ``timeout_seconds`` to answer. The API key, when there is one, goes along as it does with every request, since the reset route usually sits behind the same server. """ at_ns = time.time_ns() deadline = time.monotonic() + retry_seconds attempts = 1 reset = _post_reset(url, at_ns, timeout_seconds=timeout_seconds, api_key=api_key) while reset.answer == REFUSED and ( deadline - time.monotonic() >= RESET_RETRY_INTERVAL_SECONDS ): time.sleep(RESET_RETRY_INTERVAL_SECONDS) attempts += 1 reset = _post_reset( url, at_ns, timeout_seconds=timeout_seconds, api_key=api_key ) if reset.answer == REFUSED: error = f"refused: success false on {attempts} attempts" return _with(reset, error=error, attempts=attempts) return _with(reset, error=reset.error, attempts=attempts)
def _post_reset( url: str, at_ns: int, *, timeout_seconds: float, api_key: str | None ) -> CacheReset: recorded = redact_url(url) headers = {"Authorization": f"Bearer {api_key}"} if api_key else {} request = urllib.request.Request(url, data=b"", headers=headers, method="POST") try: # As inference requests do: no redirects, and no environment proxy. with inference_opener().open(request, timeout=timeout_seconds) as response: success = _success_field(response.read(_RESET_BODY_LIMIT)) return CacheReset( recorded, at_ns, status=int(response.status), success=success, answered_at_ns=time.time_ns(), ) except urllib.error.HTTPError as exc: return CacheReset( recorded, at_ns, status=exc.code, error=f"HTTP {exc.code}", answered_at_ns=time.time_ns(), ) except OSError as exc: return CacheReset( recorded, at_ns, error=f"{type(exc).__name__}: {exc}", answered_at_ns=time.time_ns(), ) def _success_field(body: bytes) -> bool | None: """vLLM's ``{"success": bool}``: only ``true`` is a yes; None without one.""" try: payload = json.loads(body.decode("utf-8")) except (UnicodeDecodeError, ValueError): return None if not isinstance(payload, dict) or "success" not in payload: return None return payload["success"] is True def _with(reset: CacheReset, *, error: str | None, attempts: int) -> CacheReset: return CacheReset( reset.url, reset.at_ns, status=reset.status, error=error, success=reset.success, attempts=attempts, answered_at_ns=reset.answered_at_ns, )
[docs] def run_kind( requested: str, warmup_requests: int, reset: CacheReset | None = None ) -> str: """Whether a case was designed as a cold start or a warmed-up steady state. This names the run's design, not evidence about the cache. A cold start whose reset failed is unspecified: the failure shows the cache was not cleared. """ if warmup_requests > 0: return "steady_state" if requested == COLD and (reset is None or reset.succeeded): return "cold_start" return UNSPECIFIED
[docs] def cache_state_record( *, session_id: str, case_id: str, requested: str, reset: CacheReset | None, warmup_requests: int, ) -> dict[str, Any]: """The ``infer.cache_state`` record written before a case runs.""" return { "schema_version": 1, "event_type": "infer.cache_state", "session_id": session_id, "case_id": case_id, "requested": requested, "reset": reset.to_record() if reset is not None else None, "attempted": reset is not None, "acknowledged": reset is not None and reset.acknowledged, "verified": UNVERIFIED, "reason": _reason(requested, reset), "run_kind": run_kind(requested, warmup_requests, reset), }
def _reason(requested: str, reset: CacheReset | None) -> str: if reset is not None and not reset.succeeded: return f"the cache reset failed ({reset.error or reset.status})" if reset is not None and reset.acknowledged: return ( "the server acknowledged the cache reset, but no engine adapter " "can confirm it" ) if reset is not None: return ( f"the cache reset returned HTTP {reset.status} without saying " "whether it succeeded, and no engine adapter can confirm it" ) if requested == COLD: return "nothing reset the cache, and no engine adapter can read it" return ( "no cache state was requested; earlier traffic, including an earlier " "run with the same seed, decides what the cache holds" )
[docs] def cache_summary(record: dict[str, Any] | None) -> dict[str, Any]: """The cache block of a case report; older artifacts did not record one.""" if record is None: return { "requested": UNSPECIFIED, "reset": None, "attempted": None, "acknowledged": None, "verified": UNVERIFIED, "reason": "the artifact does not record a cache state", "run_kind": None, } return { "requested": record.get("requested"), "reset": record.get("reset"), # Artifacts written before resets were parsed do not say. "attempted": record.get("attempted"), "acknowledged": record.get("acknowledged"), "verified": record.get("verified"), "reason": record.get("reason"), "run_kind": record.get("run_kind"), }
[docs] def cache_lines(cache: Any) -> list[str]: """A text-report line when a cache state was requested or reset.""" if not isinstance(cache, dict): return [] if cache.get("requested") != COLD and cache.get("reset") is None: return [] return [ f" cache: {cache.get('requested')} requested, {cache.get('verified')} " f"({cache.get('reason')}); run kind {cache.get('run_kind')}" ]