diff --git a/headroom/proxy/handlers/anthropic.py b/headroom/proxy/handlers/anthropic.py index e03ef5744..2fb1e0708 100644 --- a/headroom/proxy/handlers/anthropic.py +++ b/headroom/proxy/handlers/anthropic.py @@ -1736,6 +1736,7 @@ class AnthropicHandlerMixin: tags, optimization_latency, pipeline_timing=pipeline_timing, + original_messages=original_client_messages, ) else: async with stage_timer.measure("upstream_connect"): @@ -1837,7 +1838,15 @@ class AnthropicHandlerMixin: turn_id=compute_turn_id( model, body.get("system"), body.get("messages") ), - request_messages=body.get("messages") + # `original_client_messages` is the deep-copied + # pre-compression snapshot; `body["messages"]` + # is the compressed list sent upstream. Both + # share the `log_full_messages` gate so the two + # sides stay symmetric. + request_messages=original_client_messages + if self.config.log_full_messages + else None, + compressed_messages=body.get("messages") if self.config.log_full_messages else None, ) @@ -2358,7 +2367,16 @@ class AnthropicHandlerMixin: turn_id=compute_turn_id( model, body.get("system"), body.get("messages") ), - request_messages=messages if self.config.log_full_messages else None, + # `original_client_messages` is the deep-copied + # pre-compression snapshot; `body["messages"]` is the + # compressed list sent upstream. Both gated by + # `log_full_messages`. + request_messages=original_client_messages + if self.config.log_full_messages + else None, + compressed_messages=body.get("messages") + if self.config.log_full_messages + else None, ) ) diff --git a/headroom/proxy/handlers/streaming.py b/headroom/proxy/handlers/streaming.py index 09b081052..f5896a350 100644 --- a/headroom/proxy/handlers/streaming.py +++ b/headroom/proxy/handlers/streaming.py @@ -785,6 +785,7 @@ class StreamingMixin: uncached_input_tokens=uncached_input_tokens, ttfb_ms=stream_state["ttfb_ms"] or total_latency, pipeline_timing=pipeline_timing, + original_messages=original_messages, ) await self._record_request_outcome(outcome) @@ -1349,6 +1350,7 @@ class StreamingMixin: tags: dict[str, str], optimization_latency: float, pipeline_timing: dict[str, float] | None = None, + original_messages: list[dict] | None = None, ) -> StreamingResponse: """Stream response from Bedrock backend with metrics tracking. @@ -1455,6 +1457,7 @@ class StreamingMixin: cache_write_1h_tokens=stream_state["cache_creation_ephemeral_1h_input_tokens"], ttfb_ms=stream_state["ttfb_ms"] or 0, pipeline_timing=pipeline_timing, + original_messages=original_messages, ) await self._record_request_outcome(outcome) diff --git a/headroom/proxy/models.py b/headroom/proxy/models.py index 0d1a98fef..ee4352071 100644 --- a/headroom/proxy/models.py +++ b/headroom/proxy/models.py @@ -48,6 +48,11 @@ class RequestLog: # Request/Response (optional, for debugging) request_messages: list[dict] | None = None + # Messages after compression, as actually sent upstream. Paired with + # `request_messages` (the pre-compression snapshot) so consumers can diff + # the two sides of the compression. Governed by the same + # `log_full_messages` gate as `request_messages`. + compressed_messages: list[dict] | None = None response_content: str | None = None error: str | None = None diff --git a/headroom/proxy/outcome.py b/headroom/proxy/outcome.py index ebd17516d..3e5d0fc5d 100644 --- a/headroom/proxy/outcome.py +++ b/headroom/proxy/outcome.py @@ -127,6 +127,12 @@ class RequestOutcome: num_messages: int = 0 turn_id: str | None = None request_messages: list[dict[str, Any]] | None = None + # Post-compression messages actually sent upstream, paired with + # ``request_messages`` (pre-compression) so consumers can diff the two. + # Only populated when a caller threads in the pre-compression snapshot + # (``original_messages``); otherwise ``request_messages`` carries the sent + # body for backward compatibility and this stays ``None``. + compressed_messages: list[dict[str, Any]] | None = None tags: dict[str, str] = field(default_factory=dict) client: str | None = None project: str | None = None @@ -201,6 +207,7 @@ class RequestOutcome: ttfb_ms: float = 0.0, pipeline_timing: dict[str, float] | None = None, waste_signals: dict[str, int] | None = None, + original_messages: list[dict] | None = None, ) -> RequestOutcome: """Construct an outcome from the locals available at streaming finalize. Three streaming finalizers @@ -246,6 +253,25 @@ class RequestOutcome: if system is None: system = body.get("systemInstruction") + # ``request_items`` is ``body["messages"]`` (or ``body["contents"]`` + # for Gemini, falling back to ``[]``) — the post-compression list the + # caller already mutated in place before finalize. When a + # caller threads in ``original_messages`` (the pre-compression + # snapshot), log it as ``request_messages`` and the sent body as + # ``compressed_messages`` so the two sides stay diffable. Callers that + # don't thread it in (gemini ``contents``, OpenAI-via-backend) keep the + # prior behaviour: sent body under ``request_messages``, no compressed + # side. Both sides share the ``log_full_messages`` gate. + if not log_full_messages: + log_request_messages = None + log_compressed_messages = None + elif original_messages is not None: + log_request_messages = original_messages + log_compressed_messages = request_items + else: + log_request_messages = request_items + log_compressed_messages = None + return cls( request_id=request_id, provider=provider, @@ -271,7 +297,8 @@ class RequestOutcome: turn_id=compute_turn_id(model, system, turn_messages), tags=tags or {}, client=client, - request_messages=request_items if log_full_messages else None, + request_messages=log_request_messages, + compressed_messages=log_compressed_messages, ) @@ -376,6 +403,7 @@ async def emit_request_outcome(handler: Any, outcome: RequestOutcome) -> None: transforms_applied=list(outcome.transforms_applied), waste_signals=outcome.waste_signals, request_messages=outcome.request_messages, + compressed_messages=outcome.compressed_messages, turn_id=outcome.turn_id, ) ) diff --git a/headroom/proxy/request_logger.py b/headroom/proxy/request_logger.py index c367b6f4d..21ea155a8 100644 --- a/headroom/proxy/request_logger.py +++ b/headroom/proxy/request_logger.py @@ -212,9 +212,9 @@ class RequestLogger: """Log a request. Oldest entries are automatically removed when limit reached. Phase G PR-G3 (P4-45): base64-encoded image payloads in - ``request_messages`` / ``response_content`` are redacted - before write. Redaction also applies to the in-memory deque - so the ``/stats/recent_requests`` endpoint never serves a + ``request_messages`` / ``compressed_messages`` / ``response_content`` + are redacted before write. Redaction also applies to the in-memory + deque so the ``/stats/recent_requests`` endpoint never serves a multi-MB image either. """ # Redact image payloads in-place on the deque entry so memory @@ -223,6 +223,8 @@ class RequestLogger: # ``get_recent_with_messages`` unchanged. if entry.request_messages is not None: entry.request_messages = redact_image_base64(entry.request_messages) + if entry.compressed_messages is not None: + entry.compressed_messages = redact_image_base64(entry.compressed_messages) if entry.response_content is not None: entry.response_content = redact_image_base64(entry.response_content) @@ -234,20 +236,21 @@ class RequestLogger: log_dict = asdict(entry) if not self.log_full_messages: log_dict.pop("request_messages", None) + log_dict.pop("compressed_messages", None) log_dict.pop("response_content", None) f.write(json.dumps(log_dict) + "\n") except OSError: pass # Graceful degradation: memory-only logging continues def get_recent(self, n: int = 100) -> list[dict]: - """Get recent log entries (without request_messages and response_content).""" + """Get recent log entries (without request/compressed messages and response_content).""" # Convert deque to list for slicing (deque doesn't support slicing) entries = list(self._logs)[-n:] return [ { k: v for k, v in asdict(e).items() - if k not in ("request_messages", "response_content") + if k not in ("request_messages", "compressed_messages", "response_content") } for e in entries ] @@ -289,6 +292,8 @@ class RequestLogger: # Messages and response can be large if log_entry.request_messages: size_bytes += sys.getsizeof(log_entry.request_messages) + if log_entry.compressed_messages: + size_bytes += sys.getsizeof(log_entry.compressed_messages) if log_entry.response_content: size_bytes += len(log_entry.response_content) diff --git a/headroom/proxy/server.py b/headroom/proxy/server.py index 8a6867950..f198af98e 100644 --- a/headroom/proxy/server.py +++ b/headroom/proxy/server.py @@ -2533,6 +2533,7 @@ def create_app(config: ProxyConfig | None = None) -> FastAPI: "savings_percent": log.get("savings_percent"), "transforms_applied": log.get("transforms_applied", []), "request_messages": log.get("request_messages"), + "compressed_messages": log.get("compressed_messages"), "response_content": log.get("response_content"), "turn_id": log.get("turn_id"), } diff --git a/tests/test_proxy/test_request_logger.py b/tests/test_proxy/test_request_logger.py new file mode 100644 index 000000000..b9d140baf --- /dev/null +++ b/tests/test_proxy/test_request_logger.py @@ -0,0 +1,100 @@ +"""Tests for the in-memory request logger. + +Covers the `log_full_messages` gate, which controls whether the +pre-compression (`request_messages`) and post-compression +(`compressed_messages`) payloads persist past the in-memory entry onto disk. +Both sides are governed by the same flag so the two sides of the compression +stay in sync - it's pointless to store one without the other. +""" + +from __future__ import annotations + +from headroom.proxy.models import RequestLog +from headroom.proxy.request_logger import RequestLogger + + +def _entry(**overrides) -> RequestLog: + base: dict = { + "request_id": "r1", + "timestamp": "2026-04-24T10:00:00Z", + "provider": "anthropic", + "model": "claude-sonnet-4-6", + "input_tokens_original": 100, + "input_tokens_optimized": 40, + "output_tokens": 10, + "tokens_saved": 60, + "savings_percent": 60.0, + "optimization_latency_ms": 1.0, + "total_latency_ms": 20.0, + "tags": {}, + "cache_hit": False, + "transforms_applied": ["kompress:user:0.4"], + } + base.update(overrides) + return RequestLog(**base) + + +def test_get_recent_strips_compressed_messages_alongside_request_and_response(): + logger = RequestLogger(log_file=None, log_full_messages=True) + logger.log( + _entry( + request_messages=[{"role": "user", "content": "pre"}], + compressed_messages=[{"role": "user", "content": "post"}], + response_content="ok", + ) + ) + + recent = logger.get_recent(10) + assert len(recent) == 1 + assert "request_messages" not in recent[0] + assert "compressed_messages" not in recent[0] + assert "response_content" not in recent[0] + + +def test_get_recent_with_messages_returns_compressed_messages(): + logger = RequestLogger(log_file=None, log_full_messages=True) + logger.log( + _entry( + request_messages=[{"role": "user", "content": "pre"}], + compressed_messages=[{"role": "user", "content": "post"}], + ) + ) + + recent = logger.get_recent_with_messages(10) + assert len(recent) == 1 + assert recent[0]["request_messages"] == [{"role": "user", "content": "pre"}] + assert recent[0]["compressed_messages"] == [{"role": "user", "content": "post"}] + + +def test_jsonl_file_strips_both_sides_when_log_full_messages_disabled(tmp_path): + log_file = tmp_path / "requests.jsonl" + logger = RequestLogger(log_file=str(log_file), log_full_messages=False) + logger.log( + _entry( + request_messages=[{"role": "user", "content": "pre"}], + compressed_messages=[{"role": "user", "content": "post"}], + response_content="ok", + ) + ) + + import json + + lines = log_file.read_text().strip().splitlines() + assert len(lines) == 1 + obj = json.loads(lines[0]) + assert "request_messages" not in obj + assert "compressed_messages" not in obj + assert "response_content" not in obj + + +def test_get_memory_stats_accounts_for_compressed_messages(): + logger = RequestLogger(log_file=None) + logger.log( + _entry( + compressed_messages=[{"role": "user", "content": "post"}], + ) + ) + + stats = logger.get_memory_stats() + assert stats.entry_count == 1 + assert stats.size_bytes > 0 diff --git a/tests/test_proxy/test_transformations_feed.py b/tests/test_proxy/test_transformations_feed.py index 84e3dab0a..c7615a47f 100644 --- a/tests/test_proxy/test_transformations_feed.py +++ b/tests/test_proxy/test_transformations_feed.py @@ -30,15 +30,24 @@ async def test_transformations_feed_endpoint_returns_list(app): @pytest.mark.asyncio async def test_transformations_feed_returns_messages(app): - """Each transformation should include request_messages and response_content.""" + """Each transformation exposes both the original request and the + post-compression form that was actually sent upstream, plus the response. + + The pre/post pair is what makes compression legible: consumers can diff + the two to see what the pipeline stripped, replaced, or kept. + """ async with AsyncClient(transport=ASGITransport(app=app), base_url="http://test") as client: response = await client.get("/transformations/feed") data = response.json() transformations = data["transformations"] for t in transformations: - assert "request_messages" in t or t.get("request_messages") is None - assert "response_content" in t or t.get("response_content") is None + assert "request_messages" in t + assert t["request_messages"] is None or isinstance(t["request_messages"], list) + assert "compressed_messages" in t + assert t["compressed_messages"] is None or isinstance(t["compressed_messages"], list) + assert "response_content" in t + assert t["response_content"] is None or isinstance(t["response_content"], str) @pytest.mark.asyncio diff --git a/tests/test_proxy_streaming_request_logger.py b/tests/test_proxy_streaming_request_logger.py index 035d9570c..1aed8bf38 100644 --- a/tests/test_proxy_streaming_request_logger.py +++ b/tests/test_proxy_streaming_request_logger.py @@ -114,9 +114,17 @@ async def test_finalize_stream_response_logs_request_for_feed(): @pytest.mark.asyncio -async def test_finalize_stream_response_includes_messages_when_log_full_messages_enabled(): +async def test_finalize_stream_response_logs_original_and_compressed_messages(): + """With log_full_messages enabled, both sides of the compression are + recorded: `request_messages` is the pre-compression snapshot the caller + threads in via `original_messages`, `compressed_messages` is what was + actually sent upstream (i.e. `body["messages"]` after in-place mutation).""" proxy = _build_proxy_with_real_logger(log_full_messages=True) - body = {"messages": [{"role": "user", "content": "hello"}]} + # `body["messages"]` models the post-compression list - the proxy mutates + # `body` in place before calling `_finalize_stream_response`, so this is + # already what was shipped over the wire. + body = {"messages": [{"role": "user", "content": "[compressed]"}]} + original = [{"role": "user", "content": "[original, pre-compression]"}] await proxy._finalize_stream_response( body=body, @@ -130,11 +138,13 @@ async def test_finalize_stream_response_includes_messages_when_log_full_messages optimization_latency=1.0, stream_state=_stream_state(output_tokens=5), start_time=0.0, + original_messages=original, ) entries = proxy.logger.get_recent_with_messages(10) assert len(entries) == 1 - assert entries[0]["request_messages"] == body["messages"] + assert entries[0]["request_messages"] == original + assert entries[0]["compressed_messages"] == body["messages"] @pytest.mark.asyncio @@ -153,11 +163,15 @@ async def test_finalize_stream_response_omits_messages_when_log_full_messages_di optimization_latency=1.0, stream_state=_stream_state(output_tokens=5), start_time=0.0, + original_messages=[{"role": "user", "content": "dropped"}], ) entries = proxy.logger.get_recent_with_messages(10) assert len(entries) == 1 + # Both sides share the same gate - neither leaks when log_full_messages + # is off. assert entries[0]["request_messages"] is None + assert entries[0]["compressed_messages"] is None @pytest.mark.asyncio