mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-29 01:42:19 +00:00
tests(e2e): add transport-agnostic recording proxy sidecar for e2e CI jobs
Introduces a mitmproxy-based recording HTTP/HTTPS sidecar that any CI
job can opt into to cache LLM-provider responses across runs. Unlike
the in-process VCR persister at tests/_vcr_redis_persister.py — which
can only intercept HTTP traffic from the same Python process where it
was loaded — this sidecar operates at the network layer, so it works
for any e2e job whose system-under-test runs in a Docker container
(every job under e2e_*, proxy_*, etc.).
Components
- tests/e2e_cassette_proxy/cache_key.py: pure-function cache-key
derivation. Hashes (method, scheme, host, path, sorted query,
allowlisted headers, canonical-JSON body); strips auth, tracing,
and SDK-metadata headers so equivalent requests collide regardless
of run-to-run noise.
- tests/e2e_cassette_proxy/redis_store.py: thin Redis wrapper that
stores one (request, response) pair per key as MessagePack
(JSON+base64 fallback). Caps per-key payload size, drops oversize
responses with a log line, and never blocks the request path on
Redis errors.
- tests/e2e_cassette_proxy/addon.py: mitmproxy addon that ties the
two together. Hosts on the passthrough list (localhost, the proxy
itself) are never cached; non-2xx upstream responses are not
persisted.
- tests/e2e_cassette_proxy/Dockerfile: pinned python:3.12-slim base +
pinned mitmproxy 11.0.2 + pinned redis-py + pinned msgpack.
- tests/e2e_cassette_proxy/trust_ca.sh: helper for SUT containers to
trust the proxy CA in every Python / curl / boto3 / node trust
store at once.
- tests/e2e_cassette_proxy/README.md: usage guide + opt-in checklist
for other e2e jobs.
CI integration
- New reusable command 'start_cassette_proxy' in .circleci/config.yml.
Builds the image, runs the sidecar wired to the project Redis, fetches
the proxy CA, and exports CASSETTE_PROXY_URL / CASSETTE_PROXY_CA into
$BASH_ENV for downstream steps.
- e2e_openai_endpoints is wired up as the canonical demo: two-line opt-in
pattern documented in the README.
- Job logs include a 'Cassette-proxy stats' step that dumps the
hit/miss/store summary via 'docker logs cassette-proxy | grep
[E2ECASS]'.
Tests
- 31 hermetic unit tests under tests/test_litellm/e2e_cassette_proxy:
- test_cache_key.py: 14 tests pinning equivalence-class behavior of
the key derivation (auth header, tracing header, JSON key order,
query order, host case all collapse; method/path/body/allowlisted
headers / query-param values do not).
- test_redis_store.py: 9 tests covering set/get round-trip, default
TTL, binary body round-trip, oversize-payload rejection, corrupt-
blob eviction, and graceful behavior when the Redis client raises.
- test_addon.py: 8 tests using a fake-mitmproxy flow to exercise
the addon end-to-end (passthrough miss, persist on 2xx, hit on
canonicalized-equivalent re-request, no-persist on 5xx, host
passthrough, replay-only 599-on-miss, record-only never serves
cache).
- 31/31 pass.
Co-authored-by: Mateo Wang <mateo-berri@users.noreply.github.com>
This commit is contained in:
parent
eab0075353
commit
765aef4ff8
12 changed files with 1327 additions and 0 deletions
|
|
@ -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
|
||||
|
|
|
|||
52
tests/e2e_cassette_proxy/Dockerfile
Normal file
52
tests/e2e_cassette_proxy/Dockerfile
Normal file
|
|
@ -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"]
|
||||
116
tests/e2e_cassette_proxy/README.md
Normal file
116
tests/e2e_cassette_proxy/README.md
Normal file
|
|
@ -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:<sha256>`.
|
||||
16
tests/e2e_cassette_proxy/__init__.py
Normal file
16
tests/e2e_cassette_proxy/__init__.py
Normal file
|
|
@ -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.
|
||||
"""
|
||||
193
tests/e2e_cassette_proxy/addon.py
Normal file
193
tests/e2e_cassette_proxy/addon.py
Normal file
|
|
@ -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"]
|
||||
158
tests/e2e_cassette_proxy/cache_key.py
Normal file
158
tests/e2e_cassette_proxy/cache_key.py
Normal file
|
|
@ -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()}"
|
||||
186
tests/e2e_cassette_proxy/redis_store.py
Normal file
186
tests/e2e_cassette_proxy/redis_store.py
Normal file
|
|
@ -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
|
||||
69
tests/e2e_cassette_proxy/trust_ca.sh
Executable file
69
tests/e2e_cassette_proxy/trust_ca.sh
Executable file
|
|
@ -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"
|
||||
0
tests/test_litellm/e2e_cassette_proxy/__init__.py
Normal file
0
tests/test_litellm/e2e_cassette_proxy/__init__.py
Normal file
207
tests/test_litellm/e2e_cassette_proxy/test_addon.py
Normal file
207
tests/test_litellm/e2e_cassette_proxy/test_addon.py
Normal file
|
|
@ -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
|
||||
146
tests/test_litellm/e2e_cassette_proxy/test_cache_key.py
Normal file
146
tests/test_litellm/e2e_cassette_proxy/test_cache_key.py
Normal file
|
|
@ -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
|
||||
111
tests/test_litellm/e2e_cassette_proxy/test_redis_store.py
Normal file
111
tests/test_litellm/e2e_cassette_proxy/test_redis_store.py
Normal file
|
|
@ -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
|
||||
Loading…
Add table
Reference in a new issue