Merge pull request #37360 from BerriAI/litellm_lit_5729_e2e_record_replay_seam

feat(e2e): add record/replay transport seam and fixture bundle format
This commit is contained in:
Mateo Wang 2026-08-19 14:28:03 -07:00 • committed by GitHub
commit 59c7e7a17d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
10 changed files with 1713 additions and 16 deletions

1
.gitignore vendored
View file

@ -1,5 +1,6 @@
.python-version
.venv
tests/e2e/.fixtures/
.venv-typecheck
.venv_policy_test
.env

View file

@ -71,6 +71,16 @@ Request and response bodies are typed pydantic models in `models.py`; only the f
Mark live tests with `@pytest.mark.e2e` (on the class or the module). Pure coverage of the harness itself carries no marker and runs regardless. Use `scoped_key` for a fresh all-models key that auto-deletes, `resources` when you need to create and tear down more than a key, and `unique_marker()` from `e2e_config` to keep prompts, tags, and customer ids from colliding across concurrent runs and the shared response cache
## Record and replay fixtures
`E2E_FIXTURE_MODE` selects the transport every client is built on: `live` (the default, and what an unset variable means: nothing changes), `record` (run against the live proxy and write every interaction to a fixture bundle), or `replay` (serve every interaction back from the bundle with no HTTP at all, so a replay run needs no proxy and cannot bill a provider). The seam is `select_transport` in `fixture_transport.py`, applied inside `build_proxy_client`; both transports fulfil the same `Transport` protocol, so no test or client changes shape in any mode
A bundle (default `tests/e2e/.fixtures`, override with `E2E_FIXTURE_DIR`) is a directory: `manifest.json` carries the record timestamp, harness git version, and format version, and each test gets a subdirectory holding one JSON file per transport call in call order (`0000-post-chat-completions.json`). Auth header values are redacted on write, and file uploads store a sha256 digest instead of the bytes; response bodies are stored verbatim (a /key/generate response keeps the ephemeral virtual key it minted), which is part of why bundles are gitignored. `fixture_bundle.py` owns the format
Replay matches calls per test by transport verb and path in recorded order and raises `ReplayMiss` on any drift, naming the recorded and the actual call; a passed test must also consume its whole recording, or teardown fails it naming the first leftover interaction. Either way the fix is always to re-record with `E2E_FIXTURE_MODE=record`. Record starts fresh every time: it wipes the previous bundle (refusing to wipe a directory that is not a bundle) and never reads it. A replay bundle whose manifest is older than seven days hard-fails at collection time naming the bundle's age, so replay can never certify against fixtures that have drifted more than a week from the live proxy
Deliberately not here yet: canonical content-based match keys (LIT-5741), streaming chunk fidelity (LIT-5742), and scoping record/replay to provider-bound traffic (LIT-5745)
## Typing
The harness is fully typed with no error budget: `make lint-e2e-basedpyright` must report zero basedpyright errors, and CI enforces that on any PR touching `tests/e2e/**/*.py`. When a response field is untyped, model it in `models.py` (just the fields you read) and let pydantic validate it, rather than threading a `dict` or `Any` through the test

View file

@ -52,6 +52,17 @@ The suites run against a live proxy, so bring one up first by running the litell
Some suites need extra services the bare proxy does not start. The `logging/` OTEL trace-completeness tests read spans back from a jaeger query API at `http://localhost:16686` (override with `E2E_OTEL_QUERY_URL`); run a `jaegertracing/all-in-one` and point `PHOENIX_COLLECTOR_HTTP_ENDPOINT` at its OTLP ingest. The `mcp/` suite needs the deterministic upstream MCP server in `mcp_tests/mcp_e2e_upstream_server.py` reachable by the proxy
### Record and replay
`E2E_FIXTURE_MODE=record` runs a suite against the live proxy as usual while writing every request/response pair to a fixture bundle (default `tests/e2e/.fixtures`, override with `E2E_FIXTURE_DIR`); `E2E_FIXTURE_MODE=replay` then runs the same suite entirely from that bundle, with no proxy traffic and no provider spend; the proxy liveness gate is skipped, so replay runs with no proxy up at all. Unset (or `live`) behaves exactly as before the knob existed
```bash
E2E_FIXTURE_MODE=record uv run pytest tests/e2e/llm_translation/ -v
E2E_FIXTURE_MODE=replay uv run pytest tests/e2e/llm_translation/ -v
```
Replay fails hard (`ReplayMiss`) when the tests drift from the recording, and a bundle older than seven days fails at collection time naming its age; either way the fix is to re-record. See `CLAUDE.md` in this directory for the bundle format and the transport seam
Tests marked `@pytest.mark.e2e` hard-fail when no proxy answers `/health/liveliness`, so a run that goes red with `No live proxy` at setup means the proxy isn't up; they never skip for a missing proxy, so an absent proxy can't be mistaken for a pass
## What a complete test looks like

View file

@ -15,19 +15,27 @@ shared fixtures build on it.
import functools
import os
from collections.abc import Iterator
from collections.abc import Generator, Iterator
from datetime import datetime, timezone
import pytest
import requests
from e2e_config import CONTROL_PLANE_BASE_URL, PROXY_BASE_URL
from e2e_config import CONTROL_PLANE_BASE_URL, FIXTURE_DIR, FIXTURE_MODE_RAW, PROXY_BASE_URL
from e2e_db import RESET_OPT_IN_ENV, reset_spend_logs, run_spend_log_cleanup
from fixture_transport import (
fixture_mode_collection_error,
fixture_report_lines,
parse_fixture_mode,
replay_leftover_error,
)
from junit_properties import attach_result_properties
from lifecycle import ProxyClientProvider, ResourceManager
from proxy_client import ProxyClient, build_proxy_client
_E2E_TEST_RAN = pytest.StashKey[bool]()
_CALL_PASSED = pytest.StashKey[bool]()
def pytest_configure(config: pytest.Config) -> None:
@ -49,6 +57,21 @@ def pytest_configure(config: pytest.Config) -> None:
)
def pytest_sessionstart(session: pytest.Session) -> None:
"""Abort before collection when E2E_FIXTURE_MODE can never work: an unknown
mode value, or replay against a missing, unreadable, or stale bundle (the
stale message names the bundle's age). Live and record modes pass through."""
reason = fixture_mode_collection_error(
FIXTURE_MODE_RAW, FIXTURE_DIR, now=datetime.now(timezone.utc)
)
if reason is not None:
raise pytest.UsageError(reason)
def pytest_report_header(config: pytest.Config) -> list[str]:
return fixture_report_lines(FIXTURE_MODE_RAW, FIXTURE_DIR, now=datetime.now(timezone.utc))
def pytest_collection_modifyitems(items: list[pytest.Item]) -> None:
"""Attach the two custom signals (suite package and covered cell ids) to every
test's user_properties so the standard JUnit report (`--junitxml`) records them
@ -91,9 +114,12 @@ def _proxy_fail_reason() -> str | None:
def pytest_runtest_setup(item: pytest.Item) -> None:
"""Hard-fail `e2e`-marked tests unless a proxy answers its liveness probe.
Unmarked tests (unit coverage of the harness) don't touch the proxy, so they
run even when none is up. Never skip for a missing proxy."""
run even when none is up. Never skip for a missing proxy. Replay mode serves
every call from the fixture bundle, so it needs no live proxy either."""
if item.get_closest_marker("e2e") is None:
return
if parse_fixture_mode(FIXTURE_MODE_RAW) == "replay":
return
reason = _proxy_fail_reason()
if reason is not None:
pytest.fail(reason)
@ -110,6 +136,36 @@ def pytest_runtest_call(item: pytest.Item) -> None:
item.session.stash[_E2E_TEST_RAN] = True
@pytest.hookimpl(wrapper=True)
def pytest_runtest_makereport(
item: pytest.Item, call: pytest.CallInfo[None]
) -> Generator[None, pytest.TestReport, pytest.TestReport]:
"""Stash the call-phase outcome so teardown can tell a passed test from a
failed one without re-deriving it."""
report = yield
if report.when == "call":
item.stash[_CALL_PASSED] = report.passed
return report
@pytest.hookimpl(wrapper=True)
def pytest_runtest_teardown(item: pytest.Item) -> Generator[None, None, None]:
"""In replay mode a passing test must consume its whole recording: leftover
interactions mean the test now makes fewer calls than it did at record time,
so the replay proved less than the bundle claims. The check runs after the
yield so fixture finalizers replay their recorded calls first. Failed tests
are left alone - their own failure already explains any unconsumed tail."""
result = yield
if not item.stash.get(_CALL_PASSED, False):
return result
reason = replay_leftover_error(
mode_raw=FIXTURE_MODE_RAW, bundle_dir=FIXTURE_DIR, test_key=item.nodeid
)
if reason is not None:
pytest.fail(reason)
return result
def pytest_sessionfinish(session: pytest.Session, exitstatus: int) -> None:
"""Once the whole e2e session is done (all suites), optionally truncate the
spend logs so the DB doesn't accumulate test rows. The truncate is destructive

View file

@ -13,6 +13,8 @@ from pathlib import Path
from dotenv import load_dotenv
from fixture_transport import deterministic_marker, parse_fixture_mode
# Local runs keep provider / DataDog keys in tests/e2e/.env (see CONTRIBUTING.md).
# Compose injects them into the proxy container, but pytest on the host does not
# inherit that file unless we load it. override=False so a real shell export wins.
@ -90,6 +92,15 @@ PROPAGATION_TIMEOUT = float(os.environ.get("E2E_PROPAGATION_TIMEOUT", "15"))
EXPECT_RUST = os.environ.get("E2E_EXPECT_RUST", "").strip().lower() in ("1", "true", "yes")
# Record/replay fixture selection (see fixture_transport.py). The raw mode value
# is parsed and validated there; "live" (the default, also for empty values)
# means the harness behaves exactly as before this knob existed.
FIXTURE_MODE_RAW = os.environ.get("E2E_FIXTURE_MODE", "live")
FIXTURE_DIR = Path(
os.environ.get("E2E_FIXTURE_DIR", "").strip()
or str(Path(__file__).resolve().parent / ".fixtures")
)
# Deliberately modest concurrency. The suite shares its proxy with every other
# suite in the run, and 750 users at spawn rate 50 saturated the request path hard
# enough to distort latency-sensitive neighbours (and to spend real provider money
@ -148,7 +159,11 @@ def datadog_mcp_url(*, toolsets: str = "core") -> str:
def unique_marker() -> str:
"""A short unique token per call/run, so concurrent runs and the shared
response cache never collide on prompts, tags, or customer ids."""
response cache never collide on prompts, tags, or customer ids. In record
and replay modes the token is deterministic per test instead, so a replay
run regenerates the exact requests the record run sent."""
if parse_fixture_mode(FIXTURE_MODE_RAW) in ("record", "replay"):
return deterministic_marker()
return uuid.uuid4().hex[:12]

314
tests/e2e/fixture_bundle.py Normal file
View file

@ -0,0 +1,314 @@
"""On-disk fixture bundle format for record/replay e2e runs (LIT-5729).
A bundle is a directory: one ``manifest.json`` (record timestamp + harness
version + format version) plus one subdirectory per test, holding one JSON file
per transport interaction in call order. Bundles older than
``MAX_BUNDLE_AGE`` hard-fail replay at collection time (see conftest), so a
green replay run can never certify against fixtures that have drifted more than
a week from the live proxy.
This module owns the format only. The transports that produce and consume it
live in fixture_transport.py; canonical request matching, streaming chunk
fidelity, and provider-scoping are follow-ups (LIT-5741/5742/5745) and are
deliberately absent here, which is why every interaction file stores the full
redacted request even though replay today matches by call order.
"""
from __future__ import annotations
import hashlib
import re
import shutil
import subprocess
from dataclasses import dataclass, field
from datetime import datetime, timedelta, timezone
from pathlib import Path
from typing import Annotated, Final, Literal
from pydantic import BaseModel, Field, JsonValue, TypeAdapter
from e2e_http import (
BinaryStream,
NetworkError,
ProbeResult,
RateLimitedError,
Result,
StreamingResponse,
Success,
UnauthorizedError,
UnknownApiError,
ValidationError,
)
BUNDLE_FORMAT_VERSION: Final = 1
MAX_BUNDLE_AGE: Final = timedelta(days=7)
MANIFEST_FILENAME: Final = "manifest.json"
_JSON: Final[TypeAdapter[JsonValue]] = TypeAdapter(JsonValue)
class Manifest(BaseModel):
format_version: int
recorded_at: datetime
harness_version: str
class RecordedRequest(BaseModel):
"""The request as the transport saw it, auth header values redacted.
Replay today only matches ``method`` (the transport verb, not the HTTP verb)
and ``path`` in call order; the rest is stored so LIT-5741 can move to
content-based match keys without re-recording. File uploads store a content
digest instead of the bytes."""
method: str
path: str
headers: dict[str, str]
params: dict[str, str] = {}
body: JsonValue | None = None
form: dict[str, str] | None = None
file_name: str | None = None
file_sha256: str | None = None
file_bytes: int | None = None
class RecordedResult(BaseModel):
"""A ``Result[R]`` flattened for disk. ``data`` holds the success payload as
raw JSON; replay re-validates it against the ``response_type`` the caller
passes, exactly like a live response body."""
shape: Literal["result"] = "result"
kind: Literal["success", "network", "unauthorized", "rate_limited", "validation", "unknown"]
status_code: int | None = None
data: JsonValue | None = None
message: str | None = None
body: str | None = None
retry_after_seconds: int | None = None
class RecordedStreaming(BaseModel):
shape: Literal["streaming"] = "streaming"
payload: StreamingResponse
class RecordedBinary(BaseModel):
shape: Literal["binary"] = "binary"
payload: BinaryStream
class RecordedProbe(BaseModel):
shape: Literal["probe"] = "probe"
payload: ProbeResult
type RecordedResponse = RecordedResult | RecordedStreaming | RecordedBinary | RecordedProbe
class Interaction(BaseModel):
request: RecordedRequest
response: Annotated[
RecordedResult | RecordedStreaming | RecordedBinary | RecordedProbe,
Field(discriminator="shape"),
]
def to_json_value(model: BaseModel) -> JsonValue:
return _JSON.validate_json(model.model_dump_json(by_alias=True))
def from_result[R: BaseModel](result: Result[R]) -> RecordedResult:
match result:
case Success(status_code=status_code, data=data):
return RecordedResult(kind="success", status_code=status_code, data=to_json_value(data))
case NetworkError(message=message):
return RecordedResult(kind="network", message=message)
case UnauthorizedError():
return RecordedResult(kind="unauthorized")
case RateLimitedError(retry_after_seconds=retry_after_seconds, body=body):
return RecordedResult(kind="rate_limited", retry_after_seconds=retry_after_seconds, body=body)
case ValidationError(message=message):
return RecordedResult(kind="validation", message=message)
case UnknownApiError(status_code=status_code, body=body):
return RecordedResult(kind="unknown", status_code=status_code, body=body)
def to_result[R: BaseModel](recorded: RecordedResult, response_type: type[R]) -> Result[R]:
match recorded.kind:
case "success":
return Success(
status_code=recorded.status_code or 200,
data=response_type.model_validate(recorded.data),
)
case "network":
return NetworkError(message=recorded.message or "")
case "unauthorized":
return UnauthorizedError()
case "rate_limited":
return RateLimitedError(
retry_after_seconds=recorded.retry_after_seconds, body=recorded.body or ""
)
case "validation":
return ValidationError(message=recorded.message or "")
case "unknown":
return UnknownApiError(status_code=recorded.status_code or 0, body=recorded.body or "")
def slugify(raw: str, *, limit: int = 60) -> str:
clean = re.sub(r"[^A-Za-z0-9_.-]+", "-", raw).strip("-")
return clean[:limit].rstrip("-")
def slug_for_test(test_key: str) -> str:
"""Directory name for one test's interactions: a readable tail plus a short
digest of the full node id, so same-named methods in different classes or
files never collide."""
digest = hashlib.sha1(test_key.encode()).hexdigest()[:8]
tail = slugify(test_key.rsplit("::", 1)[-1])
return f"{tail}-{digest}" if tail else digest
def interaction_filename(ordinal: int, request: RecordedRequest) -> str:
path_part = slugify(request.path, limit=40) or "root"
return f"{ordinal:04d}-{request.method}-{path_part}.json"
def harness_version() -> str:
try:
proc = subprocess.run(
("git", "rev-parse", "--short", "HEAD"),
cwd=Path(__file__).resolve().parent,
capture_output=True,
text=True,
timeout=10,
check=False,
)
except (OSError, subprocess.SubprocessError):
return "unknown"
return proc.stdout.strip() or "unknown"
@dataclass(slots=True)
class BundleRecorder:
"""Appends interaction files under ``root``, one subdirectory per test, with
a per-test ordinal that fixes replay order. ``prepare_bundle`` is the only
constructor: it guarantees the directory started empty with a fresh
manifest, so record mode never reads (or merges into) an existing bundle."""
root: Path
_ordinals: dict[str, int] = field(default_factory=dict)
def record(self, *, test_key: str, request: RecordedRequest, response: RecordedResponse) -> None:
slug = slug_for_test(test_key)
ordinal = self._ordinals.get(slug, 0)
self._ordinals[slug] = ordinal + 1
directory = self.root / slug
directory.mkdir(parents=True, exist_ok=True)
interaction = Interaction(request=request, response=response)
target = directory / interaction_filename(ordinal, request)
target.write_text(interaction.model_dump_json(indent=2), encoding="utf-8")
@dataclass(frozen=True, slots=True)
class UnsafeBundleDir:
path: Path
reason: str
def prepare_bundle(root: Path) -> BundleRecorder | UnsafeBundleDir:
"""Start a fresh bundle at ``root`` for record mode: wipe whatever bundle is
there and write a new manifest. Refuses to wipe a directory that is neither
empty nor a bundle (no manifest.json), so a mistyped E2E_FIXTURE_DIR can
never delete unrelated files."""
if root.exists():
if not root.is_dir():
return UnsafeBundleDir(path=root, reason="exists and is not a directory")
entries = tuple(root.iterdir())
if entries and not (root / MANIFEST_FILENAME).is_file():
return UnsafeBundleDir(
path=root,
reason=f"is not empty and has no {MANIFEST_FILENAME}; refusing to wipe a non-bundle directory",
)
shutil.rmtree(root)
root.mkdir(parents=True)
manifest = Manifest(
format_version=BUNDLE_FORMAT_VERSION,
recorded_at=datetime.now(timezone.utc),
harness_version=harness_version(),
)
(root / MANIFEST_FILENAME).write_text(manifest.model_dump_json(indent=2), encoding="utf-8")
return BundleRecorder(root=root)
@dataclass(frozen=True, slots=True)
class FreshBundle:
manifest: Manifest
@dataclass(frozen=True, slots=True)
class StaleBundle:
recorded_at: datetime
age: timedelta
limit: timedelta
@dataclass(frozen=True, slots=True)
class UnreadableBundle:
reason: str
type BundleFreshness = FreshBundle | StaleBundle | UnreadableBundle
def _read_manifest(root: Path) -> Manifest | UnreadableBundle:
manifest_path = root / MANIFEST_FILENAME
if not manifest_path.is_file():
return UnreadableBundle(reason=f"no {MANIFEST_FILENAME} found (record one with E2E_FIXTURE_MODE=record)")
try:
return Manifest.model_validate_json(manifest_path.read_text(encoding="utf-8"))
except ValueError as exc:
return UnreadableBundle(reason=f"{MANIFEST_FILENAME} is invalid: {exc}")
def check_freshness(root: Path, *, now: datetime) -> BundleFreshness:
manifest = _read_manifest(root)
if isinstance(manifest, UnreadableBundle):
return manifest
if manifest.format_version != BUNDLE_FORMAT_VERSION:
return UnreadableBundle(
reason=f"format_version {manifest.format_version} != supported {BUNDLE_FORMAT_VERSION}"
)
recorded_at = (
manifest.recorded_at
if manifest.recorded_at.tzinfo is not None
else manifest.recorded_at.replace(tzinfo=timezone.utc)
)
age = now - recorded_at
if age > MAX_BUNDLE_AGE:
return StaleBundle(recorded_at=recorded_at, age=age, limit=MAX_BUNDLE_AGE)
return FreshBundle(manifest=manifest)
def format_age(age: timedelta) -> str:
total_hours = int(age.total_seconds()) // 3600
return f"{total_hours // 24}d{total_hours % 24}h"
@dataclass(frozen=True, slots=True)
class LoadedBundle:
manifest: Manifest
interactions: dict[str, tuple[Interaction, ...]]
def load_bundle(root: Path) -> LoadedBundle | UnreadableBundle:
manifest = _read_manifest(root)
if isinstance(manifest, UnreadableBundle):
return manifest
interactions = {
directory.name: tuple(
Interaction.model_validate_json(file.read_text(encoding="utf-8"))
for file in sorted(directory.glob("*.json"))
)
for directory in sorted(root.iterdir())
if directory.is_dir()
}
return LoadedBundle(manifest=manifest, interactions=interactions)

View file

@ -0,0 +1,574 @@
"""Record/replay transports behind the same ``Transport`` protocol (LIT-5729).
``RecordingTransport`` decorates the live transport: every call passes through
unchanged and its request/response pair is appended to the fixture bundle.
``ReplayTransport`` implements the protocol from a recorded bundle alone: no
HTTP, no proxy, no provider spend. Because both fulfil ``Transport``, no test
or client changes shape; ``build_proxy_client`` picks the transport from
``E2E_FIXTURE_MODE`` (live | record | replay, default live).
Replay matches each call by test node id and call order, verifying transport
verb + path and failing hard on any drift (``ReplayMiss``). Canonical
content-based match keys are LIT-5741; streaming chunk fidelity is LIT-5742;
scoping record/replay to provider-bound traffic is LIT-5745.
"""
from __future__ import annotations
import functools
import hashlib
import os
from dataclasses import dataclass, field
from datetime import datetime
from pathlib import Path
from typing import Final, Literal, assert_never
from pydantic import BaseModel
from e2e_http import AuthHeaders, BinaryStream, ProbeResult, Result, StreamingResponse
from fixture_bundle import (
BundleRecorder,
FreshBundle,
Interaction,
LoadedBundle,
RecordedBinary,
RecordedProbe,
RecordedRequest,
RecordedResponse,
RecordedResult,
RecordedStreaming,
StaleBundle,
UnreadableBundle,
UnsafeBundleDir,
check_freshness,
format_age,
from_result,
load_bundle,
prepare_bundle,
slug_for_test,
to_json_value,
to_result,
)
from transport import Transport
type FixtureMode = Literal["live", "record", "replay"]
FIXTURE_MODES: Final[tuple[FixtureMode, ...]] = ("live", "record", "replay")
SESSION_TEST_KEY: Final = "session"
REDACTED_HEADER_NAMES: Final[frozenset[str]] = frozenset({"authorization", "x-litellm-api-key"})
REDACTED_VALUE: Final = "<redacted>"
@dataclass(frozen=True, slots=True)
class InvalidFixtureMode:
value: str
def parse_fixture_mode(raw: str) -> FixtureMode | InvalidFixtureMode:
normalized = raw.strip().lower() or "live"
match normalized:
case "live" | "record" | "replay":
return normalized
case _:
return InvalidFixtureMode(value=raw)
def current_test_key() -> str:
"""The pytest node id of the running test, from the PYTEST_CURRENT_TEST env
var pytest maintains (``<nodeid> (setup|call|teardown)``); ``session`` for
calls outside any test (e.g. session-finish cleanup)."""
raw = os.environ.get("PYTEST_CURRENT_TEST", "")
if not raw:
return SESSION_TEST_KEY
return raw.rsplit(" (", 1)[0]
class ReplayMiss(AssertionError):
"""Replay had no recorded interaction for a call the suite made. The test
drifted from the bundle (or the bundle from the suite): re-record."""
_marker_ordinals: Final[dict[str, int]] = {}
def deterministic_marker() -> str:
"""Stable stand-in for uuid-based unique markers in record and replay modes:
the Nth marker of a test is a pure function of the test's node id and N, so a
replay run regenerates exactly the model names, prompts, and tags the record
run sent and every recorded poll response still satisfies its predicate."""
test_key = current_test_key()
ordinal = _marker_ordinals.get(test_key, 0)
_marker_ordinals[test_key] = ordinal + 1
return hashlib.sha1(f"{test_key}#{ordinal}".encode()).hexdigest()[:12]
def _dump_flat(model: BaseModel | None) -> dict[str, str]:
if model is None:
return {}
dumped: dict[str, object] = model.model_dump(by_alias=True, exclude_none=True)
return {key: str(value) for key, value in dumped.items()}
def _redact(headers: dict[str, str]) -> dict[str, str]:
return {
name: REDACTED_VALUE if name.lower() in REDACTED_HEADER_NAMES else value
for name, value in headers.items()
}
def recorded_request(
method: str,
path: str,
*,
headers: BaseModel,
body: BaseModel | None = None,
params: BaseModel | None = None,
form: BaseModel | None = None,
file_name: str | None = None,
file_content: bytes | None = None,
) -> RecordedRequest:
return RecordedRequest(
method=method,
path=path,
headers=_redact(_dump_flat(headers)),
params=_dump_flat(params),
body=None if body is None else to_json_value(body),
form=None if form is None else _dump_flat(form),
file_name=file_name,
file_sha256=None if file_content is None else hashlib.sha256(file_content).hexdigest(),
file_bytes=None if file_content is None else len(file_content),
)
@dataclass(frozen=True, slots=True)
class RecordingTransport:
"""Decorator over the live transport: forwards every call and appends the
interaction to the bundle, so a green live run leaves behind exactly the
traffic replay needs."""
inner: Transport
recorder: BundleRecorder
def _record(self, request: RecordedRequest, response: RecordedResponse) -> None:
self.recorder.record(test_key=current_test_key(), request=request, response=response)
def bearer(self, key: str) -> AuthHeaders:
return self.inner.bearer(key)
@property
def master(self) -> AuthHeaders:
return self.inner.master
def post[R: BaseModel](
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
) -> Result[R]:
result = self.inner.post(path, headers=headers, json=json, response_type=response_type)
self._record(recorded_request("post", path, headers=headers, body=json), from_result(result))
return result
def get[R: BaseModel](
self,
path: str,
*,
headers: BaseModel,
params: BaseModel,
response_type: type[R],
timeout: float | None = None,
) -> Result[R]:
result = self.inner.get(
path, headers=headers, params=params, response_type=response_type, timeout=timeout
)
self._record(recorded_request("get", path, headers=headers, params=params), from_result(result))
return result
def delete[R: BaseModel](
self,
path: str,
*,
headers: BaseModel,
json: BaseModel,
response_type: type[R],
params: BaseModel | None = None,
) -> Result[R]:
result = self.inner.delete(
path, headers=headers, json=json, response_type=response_type, params=params
)
self._record(
recorded_request("delete", path, headers=headers, body=json, params=params),
from_result(result),
)
return result
def patch[R: BaseModel](
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
) -> Result[R]:
result = self.inner.patch(path, headers=headers, json=json, response_type=response_type)
self._record(recorded_request("patch", path, headers=headers, body=json), from_result(result))
return result
def put[R: BaseModel](
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
) -> Result[R]:
result = self.inner.put(path, headers=headers, json=json, response_type=response_type)
self._record(recorded_request("put", path, headers=headers, body=json), from_result(result))
return result
def stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamingResponse:
response = self.inner.stream(path, headers=headers, json=json)
self._record(
recorded_request("stream", path, headers=headers, body=json),
RecordedStreaming(payload=response),
)
return response
def stream_binary(
self, path: str, *, headers: BaseModel, json: BaseModel, chunk_size: int = 8192
) -> BinaryStream:
response = self.inner.stream_binary(path, headers=headers, json=json, chunk_size=chunk_size)
self._record(
recorded_request("stream_binary", path, headers=headers, body=json),
RecordedBinary(payload=response),
)
return response
def send(
self,
path: str,
*,
headers: BaseModel,
json: BaseModel,
params: BaseModel | None = None,
stream: bool = False,
) -> StreamingResponse:
response = self.inner.send(path, headers=headers, json=json, params=params, stream=stream)
self._record(
recorded_request("send", path, headers=headers, body=json, params=params),
RecordedStreaming(payload=response),
)
return response
def probe(self, path: str, *, params: BaseModel) -> ProbeResult:
response = self.inner.probe(path, params=params)
self._record(
recorded_request("probe", path, headers=self.master, params=params),
RecordedProbe(payload=response),
)
return response
def upload[R: BaseModel](
self,
path: str,
*,
headers: BaseModel,
form: BaseModel,
filename: str,
content: bytes,
file_content_type: str = "application/jsonl",
file_field: str = "file",
params: BaseModel | None = None,
response_type: type[R],
) -> Result[R]:
result = self.inner.upload(
path,
headers=headers,
form=form,
filename=filename,
content=content,
file_content_type=file_content_type,
file_field=file_field,
params=params,
response_type=response_type,
)
self._record(
recorded_request(
"upload",
path,
headers=headers,
params=params,
form=form,
file_name=filename,
file_content=content,
),
from_result(result),
)
return result
def download(self, path: str, *, headers: BaseModel) -> StreamingResponse:
response = self.inner.download(path, headers=headers)
self._record(
recorded_request("download", path, headers=headers),
RecordedStreaming(payload=response),
)
return response
@dataclass(slots=True)
class ReplaySource:
"""One shared cursor set over a loaded bundle, so every client built in the
session consumes the same recorded sequence per test."""
bundle: LoadedBundle
_cursors: dict[str, int] = field(default_factory=dict)
def next_interaction(self, method: str, path: str) -> Interaction:
test_key = current_test_key()
slug = slug_for_test(test_key)
recorded = self.bundle.interactions.get(slug, ())
index = self._cursors.get(slug, 0)
if index >= len(recorded):
raise ReplayMiss(
f"replay exhausted for {test_key}: call #{index + 1} ({method} {path}) has no recorded "
f"interaction ({len(recorded)} recorded under {slug}); re-record with E2E_FIXTURE_MODE=record"
)
interaction = recorded[index]
if interaction.request.method != method or interaction.request.path != path:
raise ReplayMiss(
f"replay mismatch for {test_key} at call #{index + 1}: recorded "
f"{interaction.request.method} {interaction.request.path}, test made {method} {path}; "
"re-record with E2E_FIXTURE_MODE=record"
)
self._cursors[slug] = index + 1
return interaction
def leftover_error(self, test_key: str) -> str | None:
"""Non-None when the test consumed fewer interactions than were recorded,
meaning a passing replay proved less than the bundle claims."""
slug = slug_for_test(test_key)
recorded = self.bundle.interactions.get(slug, ())
consumed = self._cursors.get(slug, 0)
if consumed >= len(recorded):
return None
pending = recorded[consumed]
return (
f"replay incomplete for {test_key}: {len(recorded) - consumed} of {len(recorded)} recorded "
f"interactions never consumed, next is {pending.request.method} {pending.request.path}; "
"re-record with E2E_FIXTURE_MODE=record"
)
def _expect_result(interaction: Interaction) -> RecordedResult:
match interaction.response:
case RecordedResult() as recorded:
return recorded
case RecordedStreaming() | RecordedBinary() | RecordedProbe():
raise ReplayMiss(
f"recorded {interaction.request.method} {interaction.request.path} is not a typed result"
)
def _expect_streaming(interaction: Interaction) -> StreamingResponse:
match interaction.response:
case RecordedStreaming(payload=payload):
return payload
case RecordedResult() | RecordedBinary() | RecordedProbe():
raise ReplayMiss(
f"recorded {interaction.request.method} {interaction.request.path} is not a streaming response"
)
@dataclass(frozen=True, slots=True)
class ReplayTransport:
"""A ``Transport`` served entirely from a recorded bundle: never opens a
connection, so a replay run cannot bill a provider."""
source: ReplaySource
master_key: str
def bearer(self, key: str) -> AuthHeaders:
return AuthHeaders(authorization=f"Bearer {key}")
@property
def master(self) -> AuthHeaders:
return self.bearer(self.master_key)
def post[R: BaseModel](
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
) -> Result[R]:
return to_result(_expect_result(self.source.next_interaction("post", path)), response_type)
def get[R: BaseModel](
self,
path: str,
*,
headers: BaseModel,
params: BaseModel,
response_type: type[R],
timeout: float | None = None,
) -> Result[R]:
return to_result(_expect_result(self.source.next_interaction("get", path)), response_type)
def delete[R: BaseModel](
self,
path: str,
*,
headers: BaseModel,
json: BaseModel,
response_type: type[R],
params: BaseModel | None = None,
) -> Result[R]:
return to_result(_expect_result(self.source.next_interaction("delete", path)), response_type)
def patch[R: BaseModel](
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
) -> Result[R]:
return to_result(_expect_result(self.source.next_interaction("patch", path)), response_type)
def put[R: BaseModel](
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
) -> Result[R]:
return to_result(_expect_result(self.source.next_interaction("put", path)), response_type)
def stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamingResponse:
return _expect_streaming(self.source.next_interaction("stream", path))
def stream_binary(
self, path: str, *, headers: BaseModel, json: BaseModel, chunk_size: int = 8192
) -> BinaryStream:
interaction = self.source.next_interaction("stream_binary", path)
match interaction.response:
case RecordedBinary(payload=payload):
return payload
case RecordedResult() | RecordedStreaming() | RecordedProbe():
raise ReplayMiss(
f"recorded stream_binary {interaction.request.path} is not a binary stream"
)
def send(
self,
path: str,
*,
headers: BaseModel,
json: BaseModel,
params: BaseModel | None = None,
stream: bool = False,
) -> StreamingResponse:
return _expect_streaming(self.source.next_interaction("send", path))
def probe(self, path: str, *, params: BaseModel) -> ProbeResult:
interaction = self.source.next_interaction("probe", path)
match interaction.response:
case RecordedProbe(payload=payload):
return payload
case RecordedResult() | RecordedStreaming() | RecordedBinary():
raise ReplayMiss(f"recorded probe {interaction.request.path} is not a probe result")
def upload[R: BaseModel](
self,
path: str,
*,
headers: BaseModel,
form: BaseModel,
filename: str,
content: bytes,
file_content_type: str = "application/jsonl",
file_field: str = "file",
params: BaseModel | None = None,
response_type: type[R],
) -> Result[R]:
return to_result(_expect_result(self.source.next_interaction("upload", path)), response_type)
def download(self, path: str, *, headers: BaseModel) -> StreamingResponse:
return _expect_streaming(self.source.next_interaction("download", path))
@functools.lru_cache(maxsize=8)
def _shared_recorder(root: Path) -> BundleRecorder:
prepared = prepare_bundle(root)
if isinstance(prepared, UnsafeBundleDir):
raise ValueError(f"E2E_FIXTURE_DIR {prepared.path} {prepared.reason}")
return prepared
@functools.lru_cache(maxsize=8)
def _shared_replay_source(root: Path) -> ReplaySource:
loaded = load_bundle(root)
if isinstance(loaded, UnreadableBundle):
raise ValueError(f"cannot replay from {root}: {loaded.reason}")
return ReplaySource(bundle=loaded)
def replay_leftover_error(*, mode_raw: str, bundle_dir: Path, test_key: str) -> str | None:
"""Teardown-time completeness check: in replay mode a passed test with
unconsumed recorded interactions must fail instead of passing against a
recording it no longer matches. Inert in every other mode."""
if parse_fixture_mode(mode_raw) != "replay":
return None
return _shared_replay_source(bundle_dir).leftover_error(test_key)
def select_transport(
live: Transport, *, mode_raw: str, bundle_dir: Path, master_key: str
) -> Transport:
"""The one seam every client build goes through: wraps (record), replaces
(replay), or passes through (live) the transport per E2E_FIXTURE_MODE. The
recorder and replay cursors are process-wide singletons per bundle dir, so
every client in a session shares one bundle and one recorded sequence."""
mode = parse_fixture_mode(mode_raw)
match mode:
case InvalidFixtureMode(value=value):
raise ValueError(f"E2E_FIXTURE_MODE={value!r} is not one of {', '.join(FIXTURE_MODES)}")
case "live":
return live
case "record":
return RecordingTransport(inner=live, recorder=_shared_recorder(bundle_dir))
case "replay":
return ReplayTransport(source=_shared_replay_source(bundle_dir), master_key=master_key)
case _:
assert_never(mode)
def fixture_mode_collection_error(mode_raw: str, bundle_dir: Path, *, now: datetime) -> str | None:
"""Session-abort reason for a fixture-mode setup that can never work, or None.
Called at collection time (conftest pytest_sessionstart) so a stale or missing
bundle fails the whole run up front, naming the bundle age, instead of failing
every test individually."""
mode = parse_fixture_mode(mode_raw)
match mode:
case InvalidFixtureMode(value=value):
return f"E2E_FIXTURE_MODE={value!r} is not one of {', '.join(FIXTURE_MODES)}"
case "live" | "record":
return None
case "replay":
freshness = check_freshness(bundle_dir, now=now)
match freshness:
case FreshBundle():
return None
case StaleBundle(recorded_at=recorded_at, age=age, limit=limit):
return (
f"fixture bundle at {bundle_dir} is stale: recorded {recorded_at.isoformat()}, "
f"age {format_age(age)} exceeds the {limit.days}-day limit; "
"re-record with E2E_FIXTURE_MODE=record"
)
case UnreadableBundle(reason=reason):
return f"E2E_FIXTURE_MODE=replay cannot use bundle at {bundle_dir}: {reason}"
case _:
assert_never(freshness)
case _:
assert_never(mode)
def fixture_report_lines(mode_raw: str, bundle_dir: Path, *, now: datetime) -> list[str]:
"""pytest report-header lines; empty in live mode so an unset
E2E_FIXTURE_MODE keeps today's output byte-identical."""
mode = parse_fixture_mode(mode_raw)
match mode:
case InvalidFixtureMode() | "live":
return []
case "record":
return [f"e2e fixture mode: record -> {bundle_dir}"]
case "replay":
freshness = check_freshness(bundle_dir, now=now)
match freshness:
case FreshBundle(manifest=manifest):
return [
f"e2e fixture mode: replay <- {bundle_dir} "
f"(recorded {manifest.recorded_at.isoformat()}, harness {manifest.harness_version})"
]
case StaleBundle() | UnreadableBundle():
return [f"e2e fixture mode: replay <- {bundle_dir}"]
case _:
assert_never(freshness)
case _:
assert_never(mode)

View file

@ -65,6 +65,8 @@ from models import (
)
from e2e_config import (
CONTROL_PLANE_BASE_URL,
FIXTURE_DIR,
FIXTURE_MODE_RAW,
MASTER_KEY,
POLL_INTERVAL,
POLL_TIMEOUT,
@ -72,6 +74,7 @@ from e2e_config import (
REQUEST_TIMEOUT,
settle_propagation,
)
from fixture_transport import select_transport
from transport import HttpTransport, SplitTransport, Transport
RowsPredicate = Callable[[list[SpendLogRow]], bool]
@ -531,19 +534,29 @@ def build_proxy_client(
The endpoints are injectable for callers that resolve the proxy some other
way than ``e2e_config``'s env names (see ``claude_code/_env.py``); they must
pass all three together, since a caller that overrides only the data plane
would leave management calls pointed at the env default."""
would leave management calls pointed at the env default.
E2E_FIXTURE_MODE wraps (record) or replaces (replay) the transport here, so
every client built from this seam records or replays without changing shape;
unset it stays the plain SplitTransport (see fixture_transport.py)."""
split = SplitTransport(
data=HttpTransport(
base_url=base_url,
master_key=master_key,
request_timeout=REQUEST_TIMEOUT,
),
control=HttpTransport(
base_url=control_plane_base_url,
master_key=master_key,
request_timeout=REQUEST_TIMEOUT,
),
)
return ProxyClient(
transport=SplitTransport(
data=HttpTransport(
base_url=base_url,
master_key=master_key,
request_timeout=REQUEST_TIMEOUT,
),
control=HttpTransport(
base_url=control_plane_base_url,
master_key=master_key,
request_timeout=REQUEST_TIMEOUT,
),
transport=select_transport(
split,
mode_raw=FIXTURE_MODE_RAW,
bundle_dir=FIXTURE_DIR,
master_key=master_key,
),
poll_timeout=POLL_TIMEOUT,
poll_interval=POLL_INTERVAL,

View file

@ -0,0 +1,218 @@
"""Harness coverage for the on-disk fixture bundle format (LIT-5729).
No proxy and no ``e2e`` marker: these pin the bundle CONTRACT - the seven-day
freshness gate that names the bundle's age, record mode's wipe safety (never
delete a directory that is not a bundle), collision-free per-test slugs, and
lossless Result round-trips - so replay can never silently drift from what
record wrote.
"""
from __future__ import annotations
from datetime import datetime, timedelta, timezone
from pathlib import Path
import pytest
from pydantic import BaseModel
from e2e_http import (
NetworkError,
RateLimitedError,
Result,
Success,
UnauthorizedError,
UnknownApiError,
ValidationError,
)
from fixture_bundle import (
BUNDLE_FORMAT_VERSION,
MANIFEST_FILENAME,
MAX_BUNDLE_AGE,
BundleRecorder,
FreshBundle,
LoadedBundle,
Manifest,
RecordedRequest,
RecordedResult,
StaleBundle,
UnreadableBundle,
UnsafeBundleDir,
check_freshness,
format_age,
from_result,
interaction_filename,
load_bundle,
prepare_bundle,
slug_for_test,
to_result,
)
NOW = datetime(2026, 8, 18, 12, 0, 0, tzinfo=timezone.utc)
class Payload(BaseModel):
value: str
def write_manifest(
root: Path, recorded_at: datetime, *, format_version: int = BUNDLE_FORMAT_VERSION
) -> None:
root.mkdir(parents=True, exist_ok=True)
manifest = Manifest(
format_version=format_version, recorded_at=recorded_at, harness_version="abc1234"
)
(root / MANIFEST_FILENAME).write_text(manifest.model_dump_json(), encoding="utf-8")
def prepared(root: Path) -> BundleRecorder:
recorder = prepare_bundle(root)
assert isinstance(recorder, BundleRecorder)
return recorder
def plain_request(path: str) -> RecordedRequest:
return RecordedRequest(method="post", path=path, headers={})
class TestResultRoundTrip:
@pytest.mark.parametrize(
"result",
[
Success(status_code=201, data=Payload(value="ok")),
NetworkError(message="connection refused"),
UnauthorizedError(),
RateLimitedError(retry_after_seconds=7, body="slow down"),
ValidationError(message="bad shape"),
UnknownApiError(status_code=502, body="upstream exploded"),
],
)
def test_every_result_kind_survives_disk_and_back(self, result: Result[Payload]) -> None:
assert to_result(from_result(result), Payload) == result
class TestFreshness:
def test_bundle_at_the_limit_is_still_fresh(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
write_manifest(root, NOW - MAX_BUNDLE_AGE)
assert isinstance(check_freshness(root, now=NOW), FreshBundle)
def test_stale_bundle_reports_age_and_limit(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
write_manifest(root, NOW - timedelta(days=8, hours=3))
freshness = check_freshness(root, now=NOW)
assert isinstance(freshness, StaleBundle)
assert freshness.age == timedelta(days=8, hours=3)
assert format_age(freshness.age) == "8d3h"
assert freshness.limit == MAX_BUNDLE_AGE
def test_naive_recorded_at_is_read_as_utc(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
write_manifest(root, (NOW - timedelta(days=1)).replace(tzinfo=None))
assert isinstance(check_freshness(root, now=NOW), FreshBundle)
def test_missing_manifest_is_unreadable_with_recording_hint(self, tmp_path: Path) -> None:
freshness = check_freshness(tmp_path / "absent", now=NOW)
assert isinstance(freshness, UnreadableBundle)
assert MANIFEST_FILENAME in freshness.reason
assert "E2E_FIXTURE_MODE=record" in freshness.reason
def test_corrupt_manifest_is_unreadable(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
root.mkdir()
(root / MANIFEST_FILENAME).write_text("{not json", encoding="utf-8")
assert isinstance(check_freshness(root, now=NOW), UnreadableBundle)
def test_unknown_format_version_is_unreadable(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
write_manifest(root, NOW, format_version=BUNDLE_FORMAT_VERSION + 1)
freshness = check_freshness(root, now=NOW)
assert isinstance(freshness, UnreadableBundle)
assert f"format_version {BUNDLE_FORMAT_VERSION + 1}" in freshness.reason
class TestPrepareBundle:
def test_fresh_directory_gets_a_fresh_manifest(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
prepared(root)
freshness = check_freshness(root, now=datetime.now(timezone.utc))
assert isinstance(freshness, FreshBundle)
assert freshness.manifest.format_version == BUNDLE_FORMAT_VERSION
assert freshness.manifest.harness_version
def test_record_wipes_the_previous_bundle_instead_of_reading_it(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
prepared(root).record(
test_key="old.py::test_old",
request=plain_request("/stale"),
response=RecordedResult(kind="unauthorized"),
)
assert any(entry.is_dir() for entry in root.iterdir())
prepared(root)
assert {entry.name for entry in root.iterdir()} == {MANIFEST_FILENAME}
def test_refuses_to_wipe_a_directory_that_is_not_a_bundle(self, tmp_path: Path) -> None:
root = tmp_path / "precious"
root.mkdir()
(root / "notes.txt").write_text("keep me", encoding="utf-8")
outcome = prepare_bundle(root)
assert isinstance(outcome, UnsafeBundleDir)
assert MANIFEST_FILENAME in outcome.reason
assert (root / "notes.txt").read_text(encoding="utf-8") == "keep me"
def test_refuses_a_path_that_is_a_file(self, tmp_path: Path) -> None:
target = tmp_path / "not-a-dir"
target.write_text("x", encoding="utf-8")
outcome = prepare_bundle(target)
assert isinstance(outcome, UnsafeBundleDir)
assert "not a directory" in outcome.reason
class TestSlugs:
def test_slug_for_test_is_deterministic(self) -> None:
key = "tests/e2e/suite/test_mod.py::TestX::test_case"
assert slug_for_test(key) == slug_for_test(key)
def test_same_tail_in_different_files_never_collides(self) -> None:
first = slug_for_test("tests/e2e/a/test_a.py::test_case")
second = slug_for_test("tests/e2e/b/test_b.py::test_case")
assert first != second
assert first.startswith("test_case-")
assert second.startswith("test_case-")
def test_interaction_filename_orders_and_slugs(self) -> None:
request = RecordedRequest(method="post", path="/chat/completions", headers={})
assert interaction_filename(3, request) == "0003-post-chat-completions.json"
class TestRecordAndLoad:
def test_load_returns_interactions_in_recorded_order(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
recorder = prepared(root)
key = "suite/test_mod.py::test_ordered"
for path in ("/first", "/second", "/third"):
recorder.record(
test_key=key,
request=plain_request(path),
response=RecordedResult(kind="unauthorized"),
)
loaded = load_bundle(root)
assert isinstance(loaded, LoadedBundle)
assert [
interaction.request.path for interaction in loaded.interactions[slug_for_test(key)]
] == ["/first", "/second", "/third"]
def test_interactions_group_per_test(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
recorder = prepared(root)
for key in ("suite/test_a.py::test_one", "suite/test_b.py::test_two"):
recorder.record(
test_key=key,
request=plain_request(f"/{key[-3:]}"),
response=RecordedResult(kind="unauthorized"),
)
loaded = load_bundle(root)
assert isinstance(loaded, LoadedBundle)
assert set(loaded.interactions) == {
slug_for_test("suite/test_a.py::test_one"),
slug_for_test("suite/test_b.py::test_two"),
}

View file

@ -0,0 +1,485 @@
"""Harness coverage for the record/replay transports (LIT-5729).
No proxy and no ``e2e`` marker. A fake in-memory ``Transport`` stands in for
the live one (dependency injection, no monkeypatching): recording must pass
every value through unchanged while writing one redacted interaction file per
call, and replay must serve identical values from the bundle alone - the
fake's call log proves nothing reaches the inner transport - failing hard
(``ReplayMiss``) on any drift in order, verb, or path. The collection-time
gate and report header are pinned here too, including the stale message that
names the bundle's age.
"""
from __future__ import annotations
import hashlib
from dataclasses import dataclass, field
from datetime import datetime, timedelta, timezone
from pathlib import Path
import pytest
from pydantic import BaseModel
from e2e_http import (
AuthHeaders,
BinaryStream,
ProbeResult,
Result,
StreamingResponse,
Success,
)
from fixture_bundle import (
BUNDLE_FORMAT_VERSION,
MANIFEST_FILENAME,
BundleRecorder,
Interaction,
LoadedBundle,
Manifest,
load_bundle,
prepare_bundle,
slug_for_test,
)
from fixture_transport import (
InvalidFixtureMode,
RecordingTransport,
ReplayMiss,
ReplaySource,
ReplayTransport,
current_test_key,
deterministic_marker,
fixture_mode_collection_error,
fixture_report_lines,
parse_fixture_mode,
replay_leftover_error,
select_transport,
)
from transport import Transport
NOW = datetime(2026, 8, 18, 12, 0, 0, tzinfo=timezone.utc)
class Payload(BaseModel):
value: str
class Body(BaseModel):
prompt: str
class Query(BaseModel):
q: str
STREAMING = StreamingResponse(
status_code=200,
body="",
content_type="text/event-stream",
chunks=2,
stream_events=["one", "two"],
stream_done=True,
)
BINARY = BinaryStream(status_code=200, content_type="audio/mpeg", chunk_count=3, total_bytes=42)
PROBE = ProbeResult(status_code=200, body="alive")
@dataclass
class FakeTransport:
calls: list[str] = field(default_factory=list)
def bearer(self, key: str) -> AuthHeaders:
return AuthHeaders(authorization=f"Bearer {key}")
@property
def master(self) -> AuthHeaders:
return self.bearer("sk-fake-master")
def _success[R: BaseModel](self, response_type: type[R]) -> Result[R]:
return Success(status_code=200, data=response_type.model_validate({"value": "live"}))
def post[R: BaseModel](
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
) -> Result[R]:
self.calls.append(f"post {path}")
return self._success(response_type)
def get[R: BaseModel](
self,
path: str,
*,
headers: BaseModel,
params: BaseModel,
response_type: type[R],
timeout: float | None = None,
) -> Result[R]:
self.calls.append(f"get {path}")
return self._success(response_type)
def delete[R: BaseModel](
self,
path: str,
*,
headers: BaseModel,
json: BaseModel,
response_type: type[R],
params: BaseModel | None = None,
) -> Result[R]:
self.calls.append(f"delete {path}")
return self._success(response_type)
def patch[R: BaseModel](
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
) -> Result[R]:
self.calls.append(f"patch {path}")
return self._success(response_type)
def put[R: BaseModel](
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
) -> Result[R]:
self.calls.append(f"put {path}")
return self._success(response_type)
def stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamingResponse:
self.calls.append(f"stream {path}")
return STREAMING
def stream_binary(
self, path: str, *, headers: BaseModel, json: BaseModel, chunk_size: int = 8192
) -> BinaryStream:
self.calls.append(f"stream_binary {path}")
return BINARY
def send(
self,
path: str,
*,
headers: BaseModel,
json: BaseModel,
params: BaseModel | None = None,
stream: bool = False,
) -> StreamingResponse:
self.calls.append(f"send {path}")
return STREAMING
def probe(self, path: str, *, params: BaseModel) -> ProbeResult:
self.calls.append(f"probe {path}")
return PROBE
def upload[R: BaseModel](
self,
path: str,
*,
headers: BaseModel,
form: BaseModel,
filename: str,
content: bytes,
file_content_type: str = "application/jsonl",
file_field: str = "file",
params: BaseModel | None = None,
response_type: type[R],
) -> Result[R]:
self.calls.append(f"upload {path}")
return self._success(response_type)
def download(self, path: str, *, headers: BaseModel) -> StreamingResponse:
self.calls.append(f"download {path}")
return STREAMING
def make_recorder(root: Path) -> BundleRecorder:
recorder = prepare_bundle(root)
assert isinstance(recorder, BundleRecorder)
return recorder
def replay_source(root: Path) -> ReplaySource:
loaded = load_bundle(root)
assert isinstance(loaded, LoadedBundle)
return ReplaySource(bundle=loaded)
def this_tests_files(root: Path) -> list[Path]:
slug_dir = root / slug_for_test(current_test_key())
return sorted(slug_dir.glob("*.json")) if slug_dir.is_dir() else []
def write_manifest(root: Path, recorded_at: datetime) -> None:
root.mkdir(parents=True, exist_ok=True)
manifest = Manifest(
format_version=BUNDLE_FORMAT_VERSION, recorded_at=recorded_at, harness_version="abc1234"
)
(root / MANIFEST_FILENAME).write_text(manifest.model_dump_json(), encoding="utf-8")
class TestParseFixtureMode:
@pytest.mark.parametrize(
("raw", "expected"),
[("live", "live"), ("record", "record"), ("replay", "replay"), ("", "live"), (" REPLAY ", "replay")],
)
def test_known_values_normalize(self, raw: str, expected: str) -> None:
assert parse_fixture_mode(raw) == expected
def test_unknown_value_is_invalid_with_the_original_spelling(self) -> None:
assert parse_fixture_mode("cached") == InvalidFixtureMode(value="cached")
class TestDeterministicMarker:
def test_sequence_is_a_pure_function_of_test_and_ordinal(self) -> None:
"""A replay process must regenerate exactly the markers the record
process generated, so the Nth marker of a test is pinned to a pure
function of the node id and N."""
key = current_test_key()
assert deterministic_marker() == hashlib.sha1(f"{key}#0".encode()).hexdigest()[:12]
assert deterministic_marker() == hashlib.sha1(f"{key}#1".encode()).hexdigest()[:12]
class TestCurrentTestKey:
def test_names_this_test_and_strips_the_phase(self) -> None:
key = current_test_key()
assert key.endswith("TestCurrentTestKey::test_names_this_test_and_strips_the_phase")
assert "(call)" not in key
class TestRecordingTransport:
def test_passes_the_result_through_and_writes_one_file_per_call(self, tmp_path: Path) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
recording: Transport = RecordingTransport(inner=fake, recorder=make_recorder(root))
result = recording.post(
"/model/new", headers=fake.master, json=Body(prompt="x"), response_type=Payload
)
assert result == Success(status_code=200, data=Payload(value="live"))
assert fake.calls == ["post /model/new"]
files = this_tests_files(root)
assert [file.name for file in files] == ["0000-post-model-new.json"]
interaction = Interaction.model_validate_json(files[0].read_text(encoding="utf-8"))
assert interaction.request.method == "post"
assert interaction.request.path == "/model/new"
def test_redacts_auth_header_values_in_the_recorded_request(self, tmp_path: Path) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
recording: Transport = RecordingTransport(inner=fake, recorder=make_recorder(root))
headers = AuthHeaders.model_validate(
{"authorization": "Bearer sk-secret", "x-litellm-api-key": "sk-other"}
)
recording.post("/key/generate", headers=headers, json=Body(prompt="x"), response_type=Payload)
interaction = Interaction.model_validate_json(
this_tests_files(root)[0].read_text(encoding="utf-8")
)
assert interaction.request.headers == {
"authorization": "<redacted>",
"x-litellm-api-key": "<redacted>",
}
assert "sk-secret" not in this_tests_files(root)[0].read_text(encoding="utf-8")
def test_upload_records_a_content_digest_not_the_bytes(self, tmp_path: Path) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
recording: Transport = RecordingTransport(inner=fake, recorder=make_recorder(root))
recording.upload(
"/v1/files",
headers=fake.master,
form=Query(q="batch"),
filename="batch.jsonl",
content=b'{"custom_id": "1"}',
response_type=Payload,
)
interaction = Interaction.model_validate_json(
this_tests_files(root)[0].read_text(encoding="utf-8")
)
assert interaction.request.file_name == "batch.jsonl"
assert interaction.request.file_bytes == len(b'{"custom_id": "1"}')
assert interaction.request.file_sha256 is not None
assert "custom_id" not in interaction.request.model_dump_json()
class TestReplayTransport:
def test_serves_recorded_values_without_touching_the_inner_transport(
self, tmp_path: Path
) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
recording: Transport = RecordingTransport(inner=fake, recorder=make_recorder(root))
recorded_post = recording.post(
"/model/new", headers=fake.master, json=Body(prompt="x"), response_type=Payload
)
recorded_get = recording.get(
"/v1/models", headers=fake.master, params=Query(q="all"), response_type=Payload
)
recorded_stream = recording.stream(
"/chat/completions", headers=fake.master, json=Body(prompt="hi")
)
recorded_probe = recording.probe("/health/liveliness", params=Query(q="1"))
recorded_binary = recording.stream_binary(
"/v1/audio/speech", headers=fake.master, json=Body(prompt="say")
)
calls_after_record = list(fake.calls)
replay: Transport = ReplayTransport(source=replay_source(root), master_key="sk-1234")
assert (
replay.post("/model/new", headers=replay.master, json=Body(prompt="x"), response_type=Payload)
== recorded_post
)
assert (
replay.get("/v1/models", headers=replay.master, params=Query(q="all"), response_type=Payload)
== recorded_get
)
assert (
replay.stream("/chat/completions", headers=replay.master, json=Body(prompt="hi"))
== recorded_stream
)
assert replay.probe("/health/liveliness", params=Query(q="1")) == recorded_probe
assert (
replay.stream_binary("/v1/audio/speech", headers=replay.master, json=Body(prompt="say"))
== recorded_binary
)
assert fake.calls == calls_after_record
def test_mismatched_call_names_recorded_and_actual(self, tmp_path: Path) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
recording: Transport = RecordingTransport(inner=fake, recorder=make_recorder(root))
recording.post("/model/new", headers=fake.master, json=Body(prompt="x"), response_type=Payload)
replay: Transport = ReplayTransport(source=replay_source(root), master_key="sk-1234")
with pytest.raises(ReplayMiss, match=r"recorded post /model/new, test made get /v1/models"):
replay.get("/v1/models", headers=replay.master, params=Query(q="all"), response_type=Payload)
def test_exhausted_recording_names_the_call_count(self, tmp_path: Path) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
recording: Transport = RecordingTransport(inner=fake, recorder=make_recorder(root))
recording.post("/model/new", headers=fake.master, json=Body(prompt="x"), response_type=Payload)
replay: Transport = ReplayTransport(source=replay_source(root), master_key="sk-1234")
replay.post("/model/new", headers=replay.master, json=Body(prompt="x"), response_type=Payload)
with pytest.raises(ReplayMiss, match=r"call #2 \(post /model/new\) has no recorded interaction \(1 recorded"):
replay.post("/model/new", headers=replay.master, json=Body(prompt="x"), response_type=Payload)
class TestReplayLeftover:
def test_fully_consumed_recording_leaves_nothing(self, tmp_path: Path) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
recording: Transport = RecordingTransport(inner=fake, recorder=make_recorder(root))
recording.post("/model/new", headers=fake.master, json=Body(prompt="x"), response_type=Payload)
source = replay_source(root)
replay: Transport = ReplayTransport(source=source, master_key="sk-1234")
replay.post("/model/new", headers=replay.master, json=Body(prompt="x"), response_type=Payload)
assert source.leftover_error(current_test_key()) is None
def test_unconsumed_trailing_interactions_name_the_next_call(self, tmp_path: Path) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
recording: Transport = RecordingTransport(inner=fake, recorder=make_recorder(root))
recording.post("/model/new", headers=fake.master, json=Body(prompt="x"), response_type=Payload)
recording.probe("/health/liveliness", params=Query(q="1"))
source = replay_source(root)
replay: Transport = ReplayTransport(source=source, master_key="sk-1234")
replay.post("/model/new", headers=replay.master, json=Body(prompt="x"), response_type=Payload)
error = source.leftover_error(current_test_key())
assert error is not None
assert "1 of 2 recorded interactions never consumed" in error
assert "next is probe /health/liveliness" in error
assert "re-record with E2E_FIXTURE_MODE=record" in error
def test_test_without_recordings_has_no_leftover(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
make_recorder(root)
assert replay_source(root).leftover_error("suite.py::test_never_recorded") is None
def test_inert_outside_replay_mode(self, tmp_path: Path) -> None:
missing = tmp_path / "missing"
assert replay_leftover_error(mode_raw="", bundle_dir=missing, test_key="k") is None
assert replay_leftover_error(mode_raw="record", bundle_dir=missing, test_key="k") is None
def test_replay_mode_reads_the_shared_bundle(self, tmp_path: Path) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
recording: Transport = RecordingTransport(inner=fake, recorder=make_recorder(root))
recording.post("/model/new", headers=fake.master, json=Body(prompt="x"), response_type=Payload)
error = replay_leftover_error(mode_raw="replay", bundle_dir=root, test_key=current_test_key())
assert error is not None
assert "1 of 1 recorded interactions never consumed" in error
class TestSelectTransport:
def test_live_returns_the_live_transport_untouched(self, tmp_path: Path) -> None:
fake = FakeTransport()
for mode_raw in ("live", ""):
assert (
select_transport(fake, mode_raw=mode_raw, bundle_dir=tmp_path / "b", master_key="sk")
is fake
)
def test_record_wraps_live_and_starts_a_fresh_bundle(self, tmp_path: Path) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
write_manifest(root, NOW - timedelta(days=30))
(root / "old-test-slug").mkdir()
(root / "old-test-slug" / "0000-post-old.json").write_text("{}", encoding="utf-8")
selected = select_transport(fake, mode_raw="record", bundle_dir=root, master_key="sk")
assert isinstance(selected, RecordingTransport)
assert selected.inner is fake
assert {entry.name for entry in root.iterdir()} == {MANIFEST_FILENAME}
def test_replay_builds_a_transport_from_the_bundle_alone(self, tmp_path: Path) -> None:
fake = FakeTransport()
root = tmp_path / "bundle"
make_recorder(root)
selected = select_transport(fake, mode_raw="replay", bundle_dir=root, master_key="sk-master")
assert isinstance(selected, ReplayTransport)
assert selected.master == AuthHeaders(authorization="Bearer sk-master")
def test_invalid_mode_raises_naming_the_value(self, tmp_path: Path) -> None:
with pytest.raises(ValueError, match="cached"):
select_transport(
FakeTransport(), mode_raw="cached", bundle_dir=tmp_path / "b", master_key="sk"
)
class TestCollectionGate:
def test_invalid_mode_names_the_value_and_the_choices(self, tmp_path: Path) -> None:
assert (
fixture_mode_collection_error("cached", tmp_path, now=NOW)
== "E2E_FIXTURE_MODE='cached' is not one of live, record, replay"
)
@pytest.mark.parametrize("mode_raw", ["live", "", "record"])
def test_live_and_record_never_block_collection(self, mode_raw: str, tmp_path: Path) -> None:
assert fixture_mode_collection_error(mode_raw, tmp_path / "missing", now=NOW) is None
def test_replay_with_no_bundle_says_how_to_record_one(self, tmp_path: Path) -> None:
reason = fixture_mode_collection_error("replay", tmp_path / "missing", now=NOW)
assert reason is not None
assert f"no {MANIFEST_FILENAME}" in reason
assert "E2E_FIXTURE_MODE=record" in reason
def test_stale_replay_bundle_fails_naming_its_age(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
write_manifest(root, NOW - timedelta(days=9, hours=5))
reason = fixture_mode_collection_error("replay", root, now=NOW)
assert reason is not None
assert "age 9d5h exceeds the 7-day limit" in reason
assert "re-record with E2E_FIXTURE_MODE=record" in reason
def test_fresh_replay_bundle_collects(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
write_manifest(root, NOW - timedelta(days=2))
assert fixture_mode_collection_error("replay", root, now=NOW) is None
class TestReportHeader:
def test_live_mode_prints_nothing(self, tmp_path: Path) -> None:
assert fixture_report_lines("live", tmp_path, now=NOW) == []
assert fixture_report_lines("", tmp_path, now=NOW) == []
def test_record_and_replay_name_the_bundle(self, tmp_path: Path) -> None:
root = tmp_path / "bundle"
recorded_at = NOW - timedelta(days=1)
write_manifest(root, recorded_at)
assert fixture_report_lines("record", root, now=NOW) == [
f"e2e fixture mode: record -> {root}"
]
replay_lines = fixture_report_lines("replay", root, now=NOW)
assert len(replay_lines) == 1
assert "replay" in replay_lines[0]
assert recorded_at.isoformat() in replay_lines[0]