diff --git a/.circleci/config.yml b/.circleci/config.yml index d2c4906ef6b..04da13167cc 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -111,6 +111,64 @@ commands: - wait_for_service: url: tcp://localhost:6379 timeout: "60" + start_cassette_proxy: + description: | + Start the e2e cassette proxy sidecar (mitmproxy + Redis-backed cache). + Egress HTTPS traffic from any container that points HTTPS_PROXY at + this sidecar is captured / replayed. After this command runs you'll + have: + - a container named ``cassette-proxy`` listening on host port 8080 + - the proxy CA at /tmp/cassette-proxy-ca.crt on the host + - $CASSETTE_PROXY_URL exported in $BASH_ENV + - $CASSETTE_PROXY_CA exported in $BASH_ENV + Consumers should ``-e HTTPS_PROXY=$CASSETTE_PROXY_URL`` and mount + $CASSETTE_PROXY_CA into the container's trust store. The + ``trust_ca.sh`` helper inside the image takes care of all known + Python / curl / boto3 trust stores in one shot. + steps: + - run: + name: Build cassette-proxy image + command: | + docker build -t litellm-cassette-proxy:ci \ + -f tests/e2e_cassette_proxy/Dockerfile \ + . + - run: + name: Run cassette-proxy + command: | + # Prefer the project-level REDIS_SSL_URL (already set in + # CircleCI's env), fall back to constructing one from + # REDIS_HOST/PORT/PASSWORD if it isn't. + if [ -n "${REDIS_SSL_URL:-}" ]; then + REDIS_TARGET="$REDIS_SSL_URL" + else + : "${REDIS_HOST:?REDIS_HOST or REDIS_SSL_URL must be set}" + : "${REDIS_PORT:?REDIS_PORT must be set}" + : "${REDIS_PASSWORD:?REDIS_PASSWORD must be set}" + REDIS_TARGET="rediss://default:${REDIS_PASSWORD}@${REDIS_HOST}:${REDIS_PORT}" + fi + docker run -d \ + --name cassette-proxy \ + -p 8080:8080 \ + -e LITELLM_E2E_CASS_REDIS_URL="$REDIS_TARGET" \ + litellm-cassette-proxy:ci + - wait_for_service: + url: tcp://localhost:8080 + timeout: "60" + - run: + name: Fetch CA from cassette-proxy + command: | + mkdir -p /tmp + for i in 1 2 3 4 5; do + if curl --silent --show-error --max-time 30 \ + --proxy http://localhost:8080 \ + -o /tmp/cassette-proxy-ca.crt http://mitm.it/cert/pem; then + break + fi + sleep 2 + done + test -s /tmp/cassette-proxy-ca.crt + echo "export CASSETTE_PROXY_URL=http://host.docker.internal:8080" >> "$BASH_ENV" + echo "export CASSETTE_PROXY_CA=/tmp/cassette-proxy-ca.crt" >> "$BASH_ENV" setup_litellm_enterprise_pip: steps: - run: @@ -1351,6 +1409,7 @@ jobs: command: | uv sync --frozen --all-groups --all-extras --python 3.12 - start_postgres + - start_cassette_proxy - attach_workspace: at: ~/project - run: @@ -1389,9 +1448,18 @@ jobs: -e LANGFUSE_PROJECT2_PUBLIC=$LANGFUSE_PROJECT2_PUBLIC \ -e LANGFUSE_PROJECT1_SECRET=$LANGFUSE_PROJECT1_SECRET \ -e LANGFUSE_PROJECT2_SECRET=$LANGFUSE_PROJECT2_SECRET \ + -e HTTP_PROXY="$CASSETTE_PROXY_URL" \ + -e HTTPS_PROXY="$CASSETTE_PROXY_URL" \ + -e NO_PROXY="localhost,127.0.0.1,host.docker.internal,postgres-db" \ + -e SSL_CERT_FILE=/etc/litellm-cassette-proxy-ca.crt \ + -e REQUESTS_CA_BUNDLE=/etc/litellm-cassette-proxy-ca.crt \ + -e CURL_CA_BUNDLE=/etc/litellm-cassette-proxy-ca.crt \ + -e AWS_CA_BUNDLE=/etc/litellm-cassette-proxy-ca.crt \ + -e NODE_EXTRA_CA_CERTS=/etc/litellm-cassette-proxy-ca.crt \ --add-host host.docker.internal:host-gateway \ --name my-app \ -v $(pwd)/litellm/proxy/example_config_yaml/oai_misc_config.yaml:/app/config.yaml \ + -v "$CASSETTE_PROXY_CA":/etc/litellm-cassette-proxy-ca.crt:ro \ litellm-docker-database:ci \ --config /app/config.yaml \ --port 4000 \ @@ -1409,6 +1477,11 @@ jobs: uv run --no-sync python -m pytest -s -vv tests/openai_endpoints_tests --junitxml=test-results/junit.xml --durations=5 no_output_timeout: 15m + - run: + name: Cassette-proxy stats + command: docker logs cassette-proxy 2>&1 | grep -E '\[E2ECASS\]' | tail -200 || true + when: always + # Store test results - store_test_results: path: test-results diff --git a/tests/e2e_cassette_proxy/Dockerfile b/tests/e2e_cassette_proxy/Dockerfile new file mode 100644 index 00000000000..d6edb666aa5 --- /dev/null +++ b/tests/e2e_cassette_proxy/Dockerfile @@ -0,0 +1,52 @@ +# Recording HTTP proxy used by e2e CI jobs. +# +# Pinned to ``python:3.12-slim`` rather than the official mitmproxy image +# because we need: +# - msgpack (binary wheels not always present in the upstream image) +# - redis-py +# - the addon source from this repo +# and pinning + verifying the upstream mitmproxy image would force us to +# carry yet another sha256 in CI. Building from a known python base +# keeps the supply chain to one tag we already control. +# +# Image is intentionally fat; it only runs in CI sidecars. + +FROM python:3.12-slim@sha256:46cb7cc2877e60fbd5e21a9ae6115c30ace7a077b9f8772da879e4590c18c2e3 + +ENV PYTHONDONTWRITEBYTECODE=1 \ + PYTHONUNBUFFERED=1 \ + PIP_DISABLE_PIP_VERSION_CHECK=1 \ + PIP_NO_CACHE_DIR=1 + +RUN apt-get update \ + && apt-get install -y --no-install-recommends curl ca-certificates \ + && rm -rf /var/lib/apt/lists/* + +# Pin every dependency. +RUN pip install \ + "mitmproxy==11.0.2" \ + "redis==5.2.0" \ + "msgpack==1.1.0" + +WORKDIR /app + +COPY tests/e2e_cassette_proxy /app/tests/e2e_cassette_proxy + +ENV PYTHONPATH=/app + +VOLUME ["/root/.mitmproxy"] + +EXPOSE 8080 + +# Args explained in tests/e2e_cassette_proxy/README.md. +# +# We run as root inside the container so VOLUME ["/root/.mitmproxy"] works +# without uid juggling — this image only runs in CI sidecars, not in +# production. +ENTRYPOINT ["mitmdump", \ + "--listen-host", "0.0.0.0", \ + "--listen-port", "8080", \ + "--set", "block_global=false", \ + "--set", "ssl_insecure=true", \ + "--set", "termlog_verbosity=info", \ + "-s", "/app/tests/e2e_cassette_proxy/addon.py"] diff --git a/tests/e2e_cassette_proxy/README.md b/tests/e2e_cassette_proxy/README.md new file mode 100644 index 00000000000..ada17f49504 --- /dev/null +++ b/tests/e2e_cassette_proxy/README.md @@ -0,0 +1,116 @@ +# e2e cassette proxy + +A sidecar HTTP/HTTPS proxy that records and replays upstream responses +in Redis. Designed for CircleCI e2e jobs whose system-under-test runs +inside Docker — the in-process VCR persister at +`tests/_vcr_redis_persister.py` can't see those requests, this can. + +## What it caches + +Every HTTPS egress that flows through the proxy and: + +- uses `GET` / `POST` / `PUT` / `PATCH` / `DELETE` +- is not destined for `localhost`, `127.0.0.1`, `host.docker.internal`, + or any host listed in `LITELLM_E2E_CASS_PASSTHROUGH_HOSTS` +- received a `2xx` response from the upstream + +is keyed on a canonical hash of `(method, scheme, host, path, sorted +query, allowlisted headers, canonical body)` and stored in Redis. On +subsequent runs, matching requests are served straight from Redis +without ever hitting the upstream. + +What's intentionally *not* keyed on (so cache hits survive normal +churn): + +- `Authorization`, `x-api-key`, `anthropic-api-key`, `openai-api-key`, + `azure-api-key`, `cookie`, AWS sigv4 headers, `x-goog-api-key`, + `x-goog-user-project` — auth rotates every run +- `User-Agent`, `x-stainless-*`, `traceparent`, `tracestate`, + `x-request-id`, `request-id` — tracing / SDK metadata +- JSON key order or whitespace inside the request body — bodies are + re-serialized in canonical form before hashing + +## How to opt a CI job in + +Two changes to the job in `.circleci/config.yml`: + +1. Add the `start_cassette_proxy` reusable command after your other + sidecars (postgres, redis, etc.) and *before* you start the + container under test: + + ```yaml + - start_postgres + - start_cassette_proxy + ``` + +2. When you `docker run` the container under test, route its egress + through the sidecar and trust its CA: + + ```yaml + docker run -d \ + ...your existing env... + -e HTTP_PROXY="$CASSETTE_PROXY_URL" \ + -e HTTPS_PROXY="$CASSETTE_PROXY_URL" \ + -e NO_PROXY="localhost,127.0.0.1,host.docker.internal" \ + -e SSL_CERT_FILE=/etc/litellm-cassette-proxy-ca.crt \ + -e REQUESTS_CA_BUNDLE=/etc/litellm-cassette-proxy-ca.crt \ + -e CURL_CA_BUNDLE=/etc/litellm-cassette-proxy-ca.crt \ + -e AWS_CA_BUNDLE=/etc/litellm-cassette-proxy-ca.crt \ + -e NODE_EXTRA_CA_CERTS=/etc/litellm-cassette-proxy-ca.crt \ + -v "$CASSETTE_PROXY_CA":/etc/litellm-cassette-proxy-ca.crt:ro \ + ...your image and command... + ``` + +`e2e_openai_endpoints` is the canonical example in this PR. To opt the +others (`proxy_e2e_anthropic_messages_tests`, +`proxy_pass_through_endpoint_tests`, `e2e_ui_testing`, +`google_generate_content_endpoint_testing`, +`proxy_logging_guardrails_model_info_tests`, +`proxy_multi_instance_tests`, `proxy_spend_accuracy_tests`, +`proxy_store_model_in_db_tests`) in, copy the same two changes. + +## Knobs + +| Env var | Where set | Effect | +|---|---|---| +| `LITELLM_E2E_CASS_REDIS_URL` | sidecar container | Override the Redis URL the sidecar uses to store cassettes. Falls back to `REDIS_URL` / `REDIS_SSL_URL` / `REDIS_HOST + REDIS_PORT + REDIS_PASSWORD`. | +| `LITELLM_E2E_CASS_PASSTHROUGH_HOSTS` | sidecar container | Extra hosts (comma-separated) to never cache. | +| `LITELLM_E2E_CASS_RECORD_ONLY` | sidecar container | When `1`, never serve from cache; always forward + persist. Use during cassette refresh runs. | +| `LITELLM_E2E_CASS_REPLAY_ONLY` | sidecar container | When `1`, never forward to upstream; serve `599` on miss. Use to *prove* a job is fully cached. | + +## Why mitmproxy and not vcrpy + +vcrpy is a Python in-process monkey-patch on `httpx`/`aiohttp`/`httpcore`. +It can only intercept HTTP traffic in the *same process* where it was +installed. Every CI job that runs the LiteLLM proxy in a Docker container +issues its upstream traffic from that container's process — not the +pytest process — so vcrpy literally has no hook to attach to. + +A network-level recording proxy is language-agnostic, in-process-agnostic, +and transport-agnostic. It works for the LiteLLM proxy (Python aiohttp), +for the websocket realtime tests (a separate server), for the OpenAI SDK +(node fetch in some paths), and for any future containerized e2e job +without any changes to the SUT. + +## Why one Redis key per request, not vcrpy-style cassettes + +vcrpy stores an ordered list of `(request, response)` episodes per test +file. That model is the source of every footgun the +`tests/llm_translation/` recorder hit (unbounded per-key growth in +`new_episodes` mode, ordering brittleness, OOM under `noeviction`). Here +each Redis key holds exactly one `(request_summary, response)` pair, +so: + +- Per-key size is bounded by one response (with an explicit + `max_payload_bytes` ceiling on top). +- Two tests issuing the same request share the cache entry for free. +- "Order" is no longer a thing. +- "Record mode" is no longer a thing — if the entry is present, replay; + else record. + +## Refreshing cassettes + +Add the `LITELLM_E2E_CASS_RECORD_ONLY=1` env var to the +`start_cassette_proxy` `docker run` flags for one CI run; every cassette +the job exercises will be re-recorded. Or wipe specific keys with +`redis-cli del litellm:e2ecass:`. diff --git a/tests/e2e_cassette_proxy/__init__.py b/tests/e2e_cassette_proxy/__init__.py new file mode 100644 index 00000000000..cd8f8ee81eb --- /dev/null +++ b/tests/e2e_cassette_proxy/__init__.py @@ -0,0 +1,16 @@ +"""HTTP recording proxy used by e2e CI jobs. + +This is intentionally separate from the in-process VCR persister under +``tests/_vcr_redis_persister.py``: that one only sees HTTP traffic from +the pytest process, and so cannot record requests originating from the +LiteLLM proxy when it runs in a Docker container (the case in every CI +job under ``e2e_*`` and ``proxy_*``). This package runs as a sidecar +mitmproxy that any container can route egress through via ``HTTPS_PROXY``, +making the recording layer transport-agnostic and language-agnostic. + +Public surface for unit tests: + +- ``cache_key.derive_cache_key`` — pure function. +- ``redis_store.RedisCassetteStore`` — thin wrapper over ``redis.Redis``. +- ``addon.CassetteAddon`` — mitmproxy addon class. +""" diff --git a/tests/e2e_cassette_proxy/addon.py b/tests/e2e_cassette_proxy/addon.py new file mode 100644 index 00000000000..4a131a0eb4a --- /dev/null +++ b/tests/e2e_cassette_proxy/addon.py @@ -0,0 +1,193 @@ +"""mitmproxy addon: cache HTTPS request/response pairs in Redis. + +Loaded by ``mitmdump -s addon.py``. Two hooks: + +- ``request(flow)``: derive cache key, look up in Redis, short-circuit + the response if we have one. +- ``response(flow)``: persist the upstream's response under that same + cache key for the next run. + +This intentionally caches *every* upstream call that flows through the +proxy, regardless of host. The expensive surface (LLM provider APIs) +is what we care about, but capturing everything also dedupes the +boring stuff (token endpoints, model-list calls, control-plane pings) +for free. +""" + +from __future__ import annotations + +import logging +import os +from typing import Optional + +from mitmproxy import ctx, http # type: ignore[import-not-found] + +from tests.e2e_cassette_proxy.cache_key import ( + DEFAULT_HEADER_ALLOWLIST, + DEFAULT_HEADER_BLOCKLIST, + derive_cache_key, +) +from tests.e2e_cassette_proxy.redis_store import ( + CachedResponse, + RedisCassetteStore, +) + +_log = logging.getLogger("litellm.e2e_cassette_proxy.addon") +if not _log.handlers: + _h = logging.StreamHandler() + _h.setFormatter(logging.Formatter("[E2ECASS] %(message)s")) + _log.addHandler(_h) + _log.setLevel(logging.INFO) + _log.propagate = False + + +# Hosts we should *never* cache — Redis itself, the proxy admin UI, +# anything pointed at localhost. Extended via env var +# ``LITELLM_E2E_CASS_PASSTHROUGH_HOSTS`` (comma-separated). +_PASSTHROUGH_HOSTS = { + "localhost", + "127.0.0.1", + "0.0.0.0", + "host.docker.internal", +} + + +def _passthrough_hosts_from_env() -> set[str]: + raw = os.environ.get("LITELLM_E2E_CASS_PASSTHROUGH_HOSTS", "") + return {h.strip().lower() for h in raw.split(",") if h.strip()} + + +def _record_only() -> bool: + return os.environ.get("LITELLM_E2E_CASS_RECORD_ONLY", "").lower() in ("1", "true") + + +def _replay_only() -> bool: + return os.environ.get("LITELLM_E2E_CASS_REPLAY_ONLY", "").lower() in ("1", "true") + + +def _shape_summary(flow: http.HTTPFlow) -> str: + return f"{flow.request.method} {flow.request.pretty_host}{flow.request.path}" + + +def _is_2xx(status_code: int) -> bool: + return 200 <= status_code < 300 + + +class CassetteAddon: + """mitmproxy addon class. Created once at startup; ``request`` and + ``response`` are invoked per flow.""" + + def __init__(self, store: Optional[RedisCassetteStore] = None) -> None: + self._store: Optional[RedisCassetteStore] = store + self._passthrough = _PASSTHROUGH_HOSTS | _passthrough_hosts_from_env() + self._record_only = _record_only() + self._replay_only = _replay_only() + self._stats = {"hit": 0, "miss": 0, "stored": 0, "skipped": 0} + + def load(self, loader) -> None: # noqa: ARG002 - mitmproxy hook signature + if self._store is None: + try: + self._store = RedisCassetteStore() + except Exception as exc: + _log.info( + f"redis-init-failed err={type(exc).__name__}: {exc}; " + f"running in passthrough mode" + ) + self._store = None + _log.info( + f"loaded record_only={self._record_only} replay_only={self._replay_only} " + f"passthrough_hosts={sorted(self._passthrough)}" + ) + + def done(self) -> None: + _log.info( + f"summary hits={self._stats['hit']} misses={self._stats['miss']} " + f"stored={self._stats['stored']} skipped={self._stats['skipped']}" + ) + + def _should_skip(self, flow: http.HTTPFlow) -> bool: + host = flow.request.pretty_host.lower() + if host in self._passthrough: + return True + # CONNECT tunnels are handled implicitly by mitmproxy; only filter + # the inner HTTP request/response. Method allowlist keeps the keys + # readable in Redis. + if flow.request.method.upper() not in ("GET", "POST", "PUT", "PATCH", "DELETE"): + return True + return False + + def _key_for(self, flow: http.HTTPFlow) -> str: + return derive_cache_key( + method=flow.request.method, + url=flow.request.pretty_url, + body=flow.request.raw_content or b"", + headers=dict(flow.request.headers), + allowlist=DEFAULT_HEADER_ALLOWLIST, + blocklist=DEFAULT_HEADER_BLOCKLIST, + ) + + def request(self, flow: http.HTTPFlow) -> None: + if self._store is None or self._should_skip(flow): + self._stats["skipped"] += 1 + return + if self._record_only: + return + key = self._key_for(flow) + cached = self._store.get(key) + if cached is None: + self._stats["miss"] += 1 + _log.info(f"miss key={key} {_shape_summary(flow)}") + if self._replay_only: + # Don't fall through to the upstream; serve a stable 599 so + # the test surfaces the missing recording loudly. + flow.response = http.Response.make( + 599, + b"e2e-cassette-proxy: replay-only and no recording for this request", + {"content-type": "text/plain"}, + ) + return + _log.info( + f"hit key={key} {_shape_summary(flow)} " + f"status={cached.status_code} bytes={len(cached.body)}" + ) + self._stats["hit"] += 1 + flow.response = http.Response.make( + cached.status_code, + cached.body, + list(cached.headers), + ) + + def response(self, flow: http.HTTPFlow) -> None: + if self._store is None or self._should_skip(flow): + return + if self._replay_only: + return + # If we already served from cache, don't re-store. + if flow.response is None or flow.response.status_code == 599: + return + # Don't cache redirects, errors, or rate-limit responses — they're + # not useful for replay and would just churn keys. + if not _is_2xx(flow.response.status_code): + return + key = self._key_for(flow) + cached = CachedResponse( + status_code=flow.response.status_code, + headers=tuple((k, v) for k, v in flow.response.headers.items()), + body=flow.response.raw_content or b"", + reason=flow.response.reason or "", + ) + if self._store.set(key, cached): + self._stats["stored"] += 1 + _log.info( + f"store key={key} {_shape_summary(flow)} " + f"status={cached.status_code} bytes={len(cached.body)}" + ) + + +# mitmproxy entrypoint. +addons = [CassetteAddon()] + + +# Re-exported for parity with mitmproxy quirks (some versions look for +# ``ctx`` symbol availability at import time). +__all__ = ["CassetteAddon", "addons", "ctx"] diff --git a/tests/e2e_cassette_proxy/cache_key.py b/tests/e2e_cassette_proxy/cache_key.py new file mode 100644 index 00000000000..1b2d3b5f320 --- /dev/null +++ b/tests/e2e_cassette_proxy/cache_key.py @@ -0,0 +1,158 @@ +"""Cache key derivation for the e2e recording proxy. + +The proxy is intentionally promiscuous about *what* it caches: any HTTPS +egress that flows through it is a candidate. The cache key has to be +stable across runs that send equivalent requests, while staying robust +to per-call noise (auth headers, tracing IDs, dates). + +We hash a canonical tuple of: + +- HTTP method +- scheme + host + port +- URL path +- query string (sorted) +- request body (raw bytes; JSON bodies are re-serialized in canonical + form so that semantically-equal payloads collide) +- a small allowlist of headers the upstream actually keys on (e.g. + ``content-type``, ``accept``) + +Anything not on that allowlist is dropped before hashing. This is the +same trade-off vcrpy makes when configuring ``filter_headers`` — +strict enough to dedupe equivalent requests, loose enough to ignore +the auth/tracing churn that varies run-to-run. +""" + +from __future__ import annotations + +import hashlib +import json +from typing import Iterable, Mapping, Optional, Sequence, Tuple +from urllib.parse import parse_qsl, urlsplit + +CACHE_KEY_PREFIX = "litellm:e2ecass:" + +# Headers that materially change what an upstream returns. +DEFAULT_HEADER_ALLOWLIST: Tuple[str, ...] = ( + "accept", + "accept-encoding", + "content-type", + "openai-beta", + "anthropic-beta", + "anthropic-version", + "x-stainless-lang", +) + +# Headers we *never* want in the key (auth, tracing, per-request noise). +DEFAULT_HEADER_BLOCKLIST: Tuple[str, ...] = ( + "authorization", + "x-api-key", + "anthropic-api-key", + "openai-api-key", + "azure-api-key", + "api-key", + "cookie", + "user-agent", + "x-amz-security-token", + "x-amz-date", + "x-amz-content-sha256", + "amz-sdk-invocation-id", + "amz-sdk-request", + "x-goog-api-key", + "x-goog-user-project", + "x-request-id", + "request-id", + "traceparent", + "tracestate", + "x-stainless-arch", + "x-stainless-os", + "x-stainless-runtime", + "x-stainless-runtime-version", + "x-stainless-package-version", + "host", + "content-length", + "connection", +) + + +def _canonical_body(body: bytes) -> bytes: + """JSON bodies often differ only in key order or whitespace; collapse + those into a single canonical form so the cache hits across sessions. + Non-JSON bodies are passed through unchanged.""" + if not body: + return b"" + try: + decoded = json.loads(body) + except (ValueError, UnicodeDecodeError): + return body + return json.dumps(decoded, sort_keys=True, separators=(",", ":")).encode("utf-8") + + +def _canonical_query(raw_query: str) -> str: + if not raw_query: + return "" + pairs = sorted(parse_qsl(raw_query, keep_blank_values=True)) + return "&".join(f"{k}={v}" for k, v in pairs) + + +def _normalized_headers( + headers: Mapping[str, str], + allowlist: Sequence[str], + blocklist: Sequence[str], +) -> Tuple[Tuple[str, str], ...]: + allow = {h.lower() for h in allowlist} + block = {h.lower() for h in blocklist} + out: list[Tuple[str, str]] = [] + for key, value in headers.items(): + lk = key.lower() + if lk in block: + continue + if allow and lk not in allow: + continue + out.append((lk, value)) + out.sort() + return tuple(out) + + +def derive_cache_key( + method: str, + url: str, + body: bytes, + headers: Mapping[str, str], + *, + allowlist: Optional[Iterable[str]] = None, + blocklist: Optional[Iterable[str]] = None, +) -> str: + """Hash the canonicalized (method, url, body, allowlisted-headers) tuple + into a stable Redis key. Equal inputs (modulo header / JSON noise) always + produce the same key. + """ + parts = urlsplit(url) + canonical_headers = _normalized_headers( + headers, + allowlist=( + tuple(allowlist) if allowlist is not None else DEFAULT_HEADER_ALLOWLIST + ), + blocklist=( + tuple(blocklist) if blocklist is not None else DEFAULT_HEADER_BLOCKLIST + ), + ) + canonical_body = _canonical_body(body) + digest = hashlib.sha256() + digest.update(method.upper().encode("ascii")) + digest.update(b"\x1f") + digest.update(parts.scheme.lower().encode("ascii")) + digest.update(b"://") + digest.update(parts.netloc.lower().encode("ascii")) + digest.update(b"\x1f") + digest.update(parts.path.encode("utf-8")) + digest.update(b"\x1f") + digest.update(_canonical_query(parts.query).encode("utf-8")) + digest.update(b"\x1f") + for k, v in canonical_headers: + digest.update(k.encode("ascii")) + digest.update(b"=") + digest.update(v.encode("utf-8", errors="replace")) + digest.update(b"\x1e") + digest.update(b"\x1f") + digest.update(canonical_body) + return f"{CACHE_KEY_PREFIX}{digest.hexdigest()}" diff --git a/tests/e2e_cassette_proxy/redis_store.py b/tests/e2e_cassette_proxy/redis_store.py new file mode 100644 index 00000000000..4f4fd8c11d2 --- /dev/null +++ b/tests/e2e_cassette_proxy/redis_store.py @@ -0,0 +1,186 @@ +"""Redis-backed key-value store for cached HTTP responses. + +The on-the-wire format for a cached entry is a single MessagePack blob +(or, if msgpack isn't available, length-prefixed JSON). We pick this +shape over vcrpy's YAML cassettes because: + +- Each Redis key holds exactly one (request_summary, response) pair. + No "growing list of episodes," no `record_mode` semantics, no + ordering brittleness — see PR #26967 thread for context. +- Binary response bodies (gzip, audio, image) round-trip without + base64 expansion on the JSON path *because* they're stored as raw + bytes inside MessagePack. +- We can bound per-key size at the storage layer (oversize responses + are dropped with a log line and never block the request). +""" + +from __future__ import annotations + +import base64 +import json +import logging +import os +from dataclasses import dataclass +from typing import Mapping, Optional, Sequence, Tuple + +DEFAULT_TTL_SECONDS = 7 * 24 * 60 * 60 # 7 days +DEFAULT_MAX_PAYLOAD_BYTES = 4 * 1024 * 1024 # 4 MiB per entry + +_log = logging.getLogger("litellm.e2e_cassette_proxy.store") +if not _log.handlers: + _h = logging.StreamHandler() + _h.setFormatter(logging.Formatter("[E2ECASS] %(message)s")) + _log.addHandler(_h) + _log.setLevel(logging.INFO) + _log.propagate = False + + +try: + import msgpack # type: ignore + + _HAVE_MSGPACK = True +except ImportError: # pragma: no cover + _HAVE_MSGPACK = False + + +@dataclass(frozen=True) +class CachedResponse: + """Wire-format-agnostic response cached for a single request.""" + + status_code: int + headers: Tuple[Tuple[str, str], ...] + body: bytes + reason: str = "" + + def to_blob(self) -> bytes: + if _HAVE_MSGPACK: + return msgpack.packb( # type: ignore[no-any-return] + { + "v": 1, + "status": self.status_code, + "reason": self.reason, + "headers": [[k, v] for k, v in self.headers], + "body": self.body, + }, + use_bin_type=True, + ) + # JSON fallback: body must be base64 because JSON can't carry raw bytes. + return json.dumps( + { + "v": 1, + "status": self.status_code, + "reason": self.reason, + "headers": [[k, v] for k, v in self.headers], + "body_b64": base64.b64encode(self.body).decode("ascii"), + } + ).encode("utf-8") + + @classmethod + def from_blob(cls, blob: bytes) -> "CachedResponse": + if _HAVE_MSGPACK: + try: + payload = msgpack.unpackb(blob, raw=False) + except Exception: # pragma: no cover - corrupt blob + payload = json.loads(blob.decode("utf-8")) + else: + payload = json.loads(blob.decode("utf-8")) + body = payload.get("body") + if body is None and "body_b64" in payload: + body = base64.b64decode(payload["body_b64"]) + headers = tuple((str(k), str(v)) for k, v in payload.get("headers", [])) + return cls( + status_code=int(payload["status"]), + headers=headers, + body=body or b"", + reason=str(payload.get("reason", "")), + ) + + +def _redis_url_from_env() -> Optional[str]: + for var in ("LITELLM_E2E_CASS_REDIS_URL", "REDIS_URL", "REDIS_SSL_URL"): + url = os.environ.get(var) + if url: + return url + host = os.environ.get("REDIS_HOST") + if not host: + return None + scheme = "rediss" if os.environ.get("REDIS_SSL", "").lower() == "true" else "redis" + auth = "" + if os.environ.get("REDIS_PASSWORD"): + user = os.environ.get("REDIS_USERNAME", "") + auth = f"{user}:{os.environ['REDIS_PASSWORD']}@" + port = os.environ.get("REDIS_PORT", "6379") + return f"{scheme}://{auth}{host}:{port}" + + +class RedisCassetteStore: + """Thin GET/SET wrapper around ``redis.Redis``. + + Operations either succeed or are dropped with a log line — the + sidecar must never block the request path because the cache is + misbehaving. ``get`` returning ``None`` is the "miss" signal; + ``set`` returning ``False`` means "we tried but couldn't persist." + """ + + def __init__( + self, + client=None, + ttl_seconds: int = DEFAULT_TTL_SECONDS, + max_payload_bytes: int = DEFAULT_MAX_PAYLOAD_BYTES, + ) -> None: + self._ttl = ttl_seconds + self._max = max_payload_bytes + self._client = client if client is not None else self._build_default_client() + + @staticmethod + def _build_default_client(): + import redis + + url = _redis_url_from_env() + if not url: + raise RuntimeError( + "Set LITELLM_E2E_CASS_REDIS_URL / REDIS_URL / REDIS_SSL_URL / REDIS_HOST" + " to enable the e2e cassette proxy" + ) + return redis.Redis.from_url( + url, + socket_timeout=5, + socket_connect_timeout=5, + decode_responses=False, + ) + + def get(self, key: str) -> Optional[CachedResponse]: + try: + blob = self._client.get(key) + except Exception as exc: + _log.info(f"get-failed key={key} err={type(exc).__name__}: {exc}") + return None + if blob is None: + return None + try: + return CachedResponse.from_blob(blob) + except Exception as exc: + _log.info(f"corrupt-blob key={key} err={type(exc).__name__}: {exc}") + try: + self._client.delete(key) + except Exception: + pass + return None + + def set(self, key: str, response: CachedResponse) -> bool: + blob = response.to_blob() + if len(blob) > self._max: + _log.info( + f"persist-skipped-oversize key={key} bytes={len(blob)} " + f"max={self._max}" + ) + return False + try: + self._client.set(key, blob, ex=self._ttl) + return True + except Exception as exc: + _log.info( + f"persist-failed key={key} bytes={len(blob)} " + f"err={type(exc).__name__}: {exc}" + ) + return False diff --git a/tests/e2e_cassette_proxy/trust_ca.sh b/tests/e2e_cassette_proxy/trust_ca.sh new file mode 100755 index 00000000000..cc8daa9dcac --- /dev/null +++ b/tests/e2e_cassette_proxy/trust_ca.sh @@ -0,0 +1,69 @@ +#!/usr/bin/env sh +# Fetch the recording proxy's CA cert and add it to the local trust store +# *and* every Python TLS client we know about. Idempotent. +# +# Why so many trust-store paths? Different clients in litellm read different +# CA bundles: +# +# - openssl / curl / requests / aiohttp default: /etc/ssl/certs/ca-certificates.crt +# (Debian/Wolfi) or the OpenSSL default (Alpine) +# - boto3 / botocore: $AWS_CA_BUNDLE if set, else certifi +# - httpx: certifi by default unless SSL_CERT_FILE points elsewhere +# - openai-python: certifi (httpx) +# - google.auth: certifi (httpx) or system roots depending on transport +# +# We update certifi in-place inside the litellm venv and additionally export +# SSL_CERT_FILE / REQUESTS_CA_BUNDLE / CURL_CA_BUNDLE so anything that +# honors the environment finds it. +# +# Caller is expected to: +# 1) export PROXY_HOST / PROXY_PORT (default cassette-proxy:8080) +# 2) source this script (so the env vars persist), or eval its output + +set -eu + +PROXY_HOST="${PROXY_HOST:-cassette-proxy}" +PROXY_PORT="${PROXY_PORT:-8080}" +CA_PATH="${CA_PATH:-/usr/local/share/ca-certificates/litellm-cassette-proxy.crt}" +CA_FETCH_ENDPOINT="http://${PROXY_HOST}:${PROXY_PORT}/cert/pem" + +mkdir -p "$(dirname "$CA_PATH")" + +# mitmproxy serves its CA at /cert/pem when it sees a plain HTTP request to +# any host (the magic bypasses the CONNECT path). +if ! curl --silent --show-error --max-time 30 \ + --proxy "http://${PROXY_HOST}:${PROXY_PORT}" \ + -o "$CA_PATH" "$CA_FETCH_ENDPOINT"; then + echo "trust_ca.sh: could not fetch CA from $CA_FETCH_ENDPOINT" >&2 + exit 1 +fi + +# Update system trust store. Best-effort across distros. +if command -v update-ca-certificates >/dev/null 2>&1; then + update-ca-certificates --fresh >/dev/null 2>&1 || true +elif command -v update-ca-trust >/dev/null 2>&1; then + cp "$CA_PATH" /etc/pki/ca-trust/source/anchors/ 2>/dev/null || true + update-ca-trust 2>/dev/null || true +fi + +# Update certifi in *every* Python venv we can find (litellm runs out of +# /app/.venv in the database image, but be permissive). +for cacert in $(find / -name cacert.pem 2>/dev/null); do + # Only append if our CA isn't already in the bundle. + if ! grep -q -F "$(head -n 2 "$CA_PATH")" "$cacert" 2>/dev/null; then + cat "$CA_PATH" >> "$cacert" + fi +done + +# Export for downstream processes. The caller should source this script +# (or otherwise propagate these env vars) for them to take effect. +export HTTP_PROXY="http://${PROXY_HOST}:${PROXY_PORT}" +export HTTPS_PROXY="http://${PROXY_HOST}:${PROXY_PORT}" +export NO_PROXY="${NO_PROXY:-localhost,127.0.0.1,host.docker.internal}" +export SSL_CERT_FILE="$CA_PATH" +export REQUESTS_CA_BUNDLE="$CA_PATH" +export CURL_CA_BUNDLE="$CA_PATH" +export AWS_CA_BUNDLE="$CA_PATH" +export NODE_EXTRA_CA_CERTS="$CA_PATH" + +echo "trust_ca.sh: trusted CA at $CA_PATH; routing egress via $HTTPS_PROXY" diff --git a/tests/test_litellm/e2e_cassette_proxy/__init__.py b/tests/test_litellm/e2e_cassette_proxy/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/tests/test_litellm/e2e_cassette_proxy/test_addon.py b/tests/test_litellm/e2e_cassette_proxy/test_addon.py new file mode 100644 index 00000000000..9c251e42224 --- /dev/null +++ b/tests/test_litellm/e2e_cassette_proxy/test_addon.py @@ -0,0 +1,207 @@ +"""Addon-level tests using a minimal fake-mitmproxy flow. + +Mitmproxy itself is heavy and fiddly to install in CI for unit tests. +The addon only touches a small surface of the ``http.HTTPFlow`` API +(``request.method``, ``request.pretty_host``, ``request.pretty_url``, +``request.path``, ``request.headers``, ``request.raw_content``, +``response.status_code``, ``response.headers``, ``response.raw_content``, +``response.reason``, plus ``http.Response.make``), so we stub exactly +that subset and exercise the addon directly. +""" + +from __future__ import annotations + +import os +import sys +import types + +import fakeredis + +sys.path.insert( + 0, os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", "..")) +) + +# Install a fake `mitmproxy` package *before* importing the addon, so the +# real mitmproxy doesn't need to be present in the unit-test env. +_mitmproxy_pkg = types.ModuleType("mitmproxy") +_http_mod = types.ModuleType("mitmproxy.http") + + +class _FakeHeaders(dict): + def items(self): + return list(super().items()) + + +class _FakeResponse: + def __init__(self, status_code, body=b"", headers=None, reason=""): + self.status_code = status_code + self.raw_content = body + self.headers = _FakeHeaders(headers or {}) + self.reason = reason + + @classmethod + def make(cls, status_code, body=b"", headers=None): + if isinstance(headers, list): + headers = dict(headers) + return cls(status_code, body=body, headers=headers) + + +class _FakeRequest: + def __init__( + self, + method, + url, + body=b"", + headers=None, + host="api.openai.com", + path="/v1/chat/completions", + ): + self.method = method + self.pretty_url = url + self.pretty_host = host + self.path = path + self.raw_content = body + self.headers = _FakeHeaders(headers or {}) + + +class _FakeFlow: + def __init__(self, request, response=None): + self.request = request + self.response = response + + +_http_mod.Response = _FakeResponse # type: ignore[attr-defined] +_http_mod.HTTPFlow = _FakeFlow # type: ignore[attr-defined] +_mitmproxy_pkg.http = _http_mod # type: ignore[attr-defined] +_mitmproxy_pkg.ctx = types.SimpleNamespace(log=types.SimpleNamespace(info=lambda *_a, **_kw: None)) # type: ignore[attr-defined] +sys.modules.setdefault("mitmproxy", _mitmproxy_pkg) +sys.modules.setdefault("mitmproxy.http", _http_mod) + +from tests.e2e_cassette_proxy.addon import CassetteAddon # noqa: E402 +from tests.e2e_cassette_proxy.redis_store import ( # noqa: E402 + CachedResponse, + RedisCassetteStore, +) + + +def _addon_with_fake_redis(): + fake = fakeredis.FakeStrictRedis() + store = RedisCassetteStore(client=fake) + return fake, CassetteAddon(store=store) + + +def _make_request(body=b'{"model":"gpt-4o","messages":[]}'): + return _FakeRequest( + method="POST", + url="https://api.openai.com/v1/chat/completions", + body=body, + headers={"content-type": "application/json", "authorization": "Bearer sk-1"}, + ) + + +def test_should_pass_through_when_no_cache_entry_exists(): + _, addon = _addon_with_fake_redis() + flow = _FakeFlow(_make_request()) + addon.request(flow) + assert flow.response is None # mitmproxy will then hit the real upstream + assert addon._stats["miss"] == 1 + + +def test_should_persist_response_on_response_hook_when_2xx(): + fake, addon = _addon_with_fake_redis() + flow = _FakeFlow(_make_request()) + addon.request(flow) # miss + flow.response = _FakeResponse( + 200, + body=b'{"id":"chatcmpl-1"}', + headers={"content-type": "application/json"}, + reason="OK", + ) + addon.response(flow) + assert addon._stats["stored"] == 1 + assert any(k.startswith(b"litellm:e2ecass:") for k in fake.keys("*")) + + +def test_should_short_circuit_on_subsequent_request_with_cache_hit(): + _, addon = _addon_with_fake_redis() + first = _FakeFlow(_make_request()) + addon.request(first) + first.response = _FakeResponse( + 200, + body=b'{"id":"chatcmpl-1"}', + headers={"content-type": "application/json"}, + ) + addon.response(first) + + # Second request: same canonical shape, different auth header. + second_req = _make_request() + second_req.headers["authorization"] = "Bearer sk-different" + second = _FakeFlow(second_req) + addon.request(second) + assert second.response is not None + assert second.response.status_code == 200 + assert second.response.raw_content == b'{"id":"chatcmpl-1"}' + assert addon._stats["hit"] == 1 + + +def test_should_not_persist_non_2xx_response(): + fake, addon = _addon_with_fake_redis() + flow = _FakeFlow(_make_request()) + addon.request(flow) + flow.response = _FakeResponse( + 500, + body=b'{"error":"internal"}', + headers={"content-type": "application/json"}, + ) + addon.response(flow) + assert addon._stats["stored"] == 0 + assert len(fake.keys("*")) == 0 + + +def test_should_skip_passthrough_hosts_completely(): + _, addon = _addon_with_fake_redis() + req = _FakeRequest( + method="POST", + url="http://localhost:4000/key/generate", + body=b"{}", + headers={"content-type": "application/json"}, + host="localhost", + path="/key/generate", + ) + flow = _FakeFlow(req) + addon.request(flow) + assert flow.response is None + assert addon._stats["skipped"] == 1 + + +def test_replay_only_should_serve_599_on_miss(): + fake, _ = _addon_with_fake_redis() + addon = CassetteAddon(store=RedisCassetteStore(client=fake)) + addon._replay_only = True + + flow = _FakeFlow(_make_request()) + addon.request(flow) + assert flow.response is not None + assert flow.response.status_code == 599 + + +def test_record_only_should_not_short_circuit_even_on_hit(): + fake = fakeredis.FakeStrictRedis() + store = RedisCassetteStore(client=fake) + pre_addon = CassetteAddon(store=store) + flow = _FakeFlow(_make_request()) + pre_addon.request(flow) + flow.response = _FakeResponse( + 200, + body=b'{"id":"chatcmpl-1"}', + headers={"content-type": "application/json"}, + ) + pre_addon.response(flow) + + # Now flip a new addon into record-only mode and verify it never serves + # from cache. + record_only_addon = CassetteAddon(store=store) + record_only_addon._record_only = True + second = _FakeFlow(_make_request()) + record_only_addon.request(second) + assert second.response is None diff --git a/tests/test_litellm/e2e_cassette_proxy/test_cache_key.py b/tests/test_litellm/e2e_cassette_proxy/test_cache_key.py new file mode 100644 index 00000000000..e184198704a --- /dev/null +++ b/tests/test_litellm/e2e_cassette_proxy/test_cache_key.py @@ -0,0 +1,146 @@ +"""Unit tests for ``tests.e2e_cassette_proxy.cache_key.derive_cache_key``. + +The cache key is the only thing standing between "cache hit" and "cache +miss," so its invariants need to be pinned down precisely. Every test +below is a single equivalence claim: "these two requests should hash +to the same key" or "these two should not." +""" + +from __future__ import annotations + +import os +import sys + +sys.path.insert( + 0, os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", "..")) +) + +from tests.e2e_cassette_proxy.cache_key import ( # noqa: E402 + CACHE_KEY_PREFIX, + derive_cache_key, +) + + +def _key( + method="POST", + url="https://api.openai.com/v1/chat/completions", + body=b"", + headers=None, +): + return derive_cache_key(method, url, body, headers or {}) + + +def test_should_produce_stable_key_with_prefix(): + key = _key() + assert key.startswith(CACHE_KEY_PREFIX) + # 64 hex chars after the prefix. + assert len(key) == len(CACHE_KEY_PREFIX) + 64 + + +def test_should_match_when_inputs_are_byte_for_byte_identical(): + a = _key(body=b"hello", headers={"content-type": "application/json"}) + b = _key(body=b"hello", headers={"content-type": "application/json"}) + assert a == b + + +def test_should_match_when_only_auth_headers_differ(): + a = _key( + headers={"authorization": "Bearer one", "content-type": "application/json"} + ) + b = _key( + headers={"authorization": "Bearer two", "content-type": "application/json"} + ) + assert a == b + + +def test_should_match_when_only_tracing_headers_differ(): + a = _key(headers={"x-request-id": "abc", "content-type": "application/json"}) + b = _key(headers={"x-request-id": "xyz", "content-type": "application/json"}) + assert a == b + + +def test_should_match_when_user_agent_differs(): + a = _key( + headers={"user-agent": "openai-python/1.0", "content-type": "application/json"} + ) + b = _key(headers={"user-agent": "curl/8", "content-type": "application/json"}) + assert a == b + + +def test_should_match_when_json_body_differs_only_in_key_order(): + a = _key( + body=b'{"model":"gpt-4o","messages":[]}', + headers={"content-type": "application/json"}, + ) + b = _key( + body=b'{"messages":[],"model":"gpt-4o"}', + headers={"content-type": "application/json"}, + ) + assert a == b + + +def test_should_match_when_query_params_only_differ_in_order(): + a = _key(url="https://api.openai.com/v1/x?b=2&a=1") + b = _key(url="https://api.openai.com/v1/x?a=1&b=2") + assert a == b + + +def test_should_match_when_host_capitalization_differs(): + a = _key(url="https://api.openai.com/v1/chat/completions") + b = _key(url="https://API.OPENAI.COM/v1/chat/completions") + assert a == b + + +def test_should_differ_when_method_changes(): + assert _key(method="GET") != _key(method="POST") + + +def test_should_differ_when_path_changes(): + a = _key(url="https://api.openai.com/v1/chat/completions") + b = _key(url="https://api.openai.com/v1/responses") + assert a != b + + +def test_should_differ_when_body_changes_meaningfully(): + a = _key( + body=b'{"model":"gpt-4o","messages":[{"role":"user","content":"hi"}]}', + headers={"content-type": "application/json"}, + ) + b = _key( + body=b'{"model":"gpt-4o","messages":[{"role":"user","content":"bye"}]}', + headers={"content-type": "application/json"}, + ) + assert a != b + + +def test_should_differ_when_allowlisted_header_changes(): + # ``anthropic-version`` is on the allowlist — changing it must miss. + a = _key(headers={"anthropic-version": "2023-06-01"}) + b = _key(headers={"anthropic-version": "2024-01-01"}) + assert a != b + + +def test_should_match_when_blocklisted_header_added_or_removed(): + a = _key(headers={"content-type": "application/json"}) + b = _key( + headers={"content-type": "application/json", "x-amz-date": "20260501T000000Z"} + ) + assert a == b + + +def test_should_handle_non_json_body_passthrough(): + a = _key( + body=b"\x00\x01\x02binary", headers={"content-type": "application/octet-stream"} + ) + b = _key( + body=b"\x00\x01\x02binary", headers={"content-type": "application/octet-stream"} + ) + c = _key(body=b"different", headers={"content-type": "application/octet-stream"}) + assert a == b + assert a != c + + +def test_should_treat_query_string_difference_as_a_miss(): + a = _key(url="https://api.openai.com/v1/files?purpose=batch") + b = _key(url="https://api.openai.com/v1/files?purpose=fine-tune") + assert a != b diff --git a/tests/test_litellm/e2e_cassette_proxy/test_redis_store.py b/tests/test_litellm/e2e_cassette_proxy/test_redis_store.py new file mode 100644 index 00000000000..41638776eb7 --- /dev/null +++ b/tests/test_litellm/e2e_cassette_proxy/test_redis_store.py @@ -0,0 +1,111 @@ +"""Unit tests for the Redis cassette store. + +We test against ``fakeredis`` so the suite stays hermetic. +""" + +from __future__ import annotations + +import os +import sys + +import fakeredis +import pytest + +sys.path.insert( + 0, os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", "..")) +) + +from tests.e2e_cassette_proxy.redis_store import ( # noqa: E402 + DEFAULT_TTL_SECONDS, + CachedResponse, + RedisCassetteStore, +) + + +def _store(max_payload_bytes=None): + fake = fakeredis.FakeStrictRedis() + if max_payload_bytes is None: + return fake, RedisCassetteStore(client=fake) + return fake, RedisCassetteStore(client=fake, max_payload_bytes=max_payload_bytes) + + +def _sample_response(body: bytes = b'{"id":"msg_1"}'): + return CachedResponse( + status_code=200, + headers=(("content-type", "application/json"),), + body=body, + reason="OK", + ) + + +def test_should_round_trip_a_set_then_get(): + _, store = _store() + store.set("k", _sample_response()) + got = store.get("k") + assert got is not None + assert got.status_code == 200 + assert got.body == b'{"id":"msg_1"}' + assert ("content-type", "application/json") in got.headers + + +def test_should_return_none_on_miss(): + _, store = _store() + assert store.get("never-set") is None + + +def test_should_apply_default_ttl_on_set(): + fake, store = _store() + store.set("k", _sample_response()) + ttl = fake.ttl("k") + # fakeredis returns a float; allow a tiny window for clock skew. + assert DEFAULT_TTL_SECONDS - 5 <= ttl <= DEFAULT_TTL_SECONDS + + +def test_should_round_trip_binary_response_bodies(): + _, store = _store() + payload = bytes(range(256)) + store.set("k", _sample_response(body=payload)) + got = store.get("k") + assert got is not None + assert got.body == payload + + +def test_should_skip_oversize_payloads_silently(): + fake, store = _store(max_payload_bytes=128) + huge = b"x" * 10_000 + persisted = store.set("k", _sample_response(body=huge)) + assert persisted is False + assert fake.exists("k") == 0 + + +def test_should_persist_payloads_at_the_size_threshold(): + fake, store = _store(max_payload_bytes=10_000) + fits = b"x" * 100 + persisted = store.set("k", _sample_response(body=fits)) + assert persisted is True + assert fake.exists("k") == 1 + + +def test_should_treat_corrupt_blob_as_miss_and_evict_it(): + fake, store = _store() + fake.set("corrupt-key", b"\x00\x01not a valid blob\xff") + assert store.get("corrupt-key") is None + assert fake.exists("corrupt-key") == 0 + + +def test_should_return_none_when_get_raises(): + class _ExplodingClient: + def get(self, key): + raise ConnectionError("redis is down") + + store = RedisCassetteStore(client=_ExplodingClient()) + assert store.get("k") is None + + +def test_should_return_false_when_set_raises(): + class _ExplodingClient: + def set(self, key, value, ex=None): # noqa: ARG002 + raise ConnectionError("redis is down") + + store = RedisCassetteStore(client=_ExplodingClient()) + assert store.set("k", _sample_response()) is False