headroom/examples/strands_bundle_demo.py
chopratejas 20dc1f28f3 fix(proxy): Strands MCP bundle + backend path fixes + Codex fail-closed protection
Three logically-related sets of proxy changes ship in this branch:

1. Strands integration on the Bedrock path (HeadroomBundle + 4 OpenAI
   handler fixes + LiteLLM cache stats + dep pin)
2. /stats MCP aggregation (cross-process events log → proxy summary)
3. Codex compression-failure fail-closed (WS + HTTP /v1/responses)

== 1. Strands integration on the Bedrock path ==

* HeadroomBundle (headroom/integrations/strands/bundle.py): single-helper
  MCP wiring for a Strands Agent — Headroom MCP server (headroom_compress
  / headroom_retrieve / headroom_stats) plus optional Serena MCP and
  optional in-process compression hook. Constructor builds unstarted
  MCPClient instances per server; Strands' Agent owns the subprocess
  lifecycle. Default config: MCP enabled, Serena enabled, hook OFF
  (proxy is the single source of truth for compression). User-side
  integration is two lines in any Strands app.

* headroom/proxy/handlers/openai.py — backend path now:
  - calls PrefixCacheTracker.update_from_response (was direct-OpenAI only)
  - intercepts CCR headroom_retrieve tool_calls server-side, mirroring
    the Anthropic handler pattern; NO silent fallback, re-raises on
    CCR errors (per feedback_no_silent_fallbacks)
  - works for both non-streaming and streaming paths

* headroom/proxy/handlers/streaming.py: _stream_openai_via_backend now
  accepts prefix_tracker + optimized_messages, parses cache stats from
  the SSE final-usage frame (cache_creation_input_tokens added to the
  state machine), records CCR retrieve feedback via a new
  _record_ccr_feedback_from_openai_sse helper. Streaming CCR intercept
  is intentionally out of scope (mirrors Anthropic streaming behaviour).

* headroom/backends/litellm.py: send_openai_message response usage block
  now carries cache_read_input_tokens / cache_creation_input_tokens
  (Anthropic/Bedrock dialect) and prompt_tokens_details.cached_tokens
  (OpenAI dialect). Backwards-compatible — cold-start callers see the
  same 3-key shape; cache keys appear only when the underlying provider
  returns them. Pinned by test_no_cache_fields_means_no_cache_keys.

* headroom/proxy/auth_mode.py: ("strands-agents/", "strands") added to
  CLIENT_UA_MAP. Production callers should also set X-Client: strands
  since the default openai-python UA carries no Strands signal.

* pyproject.toml: huggingface-hub>=1.5.0,<2.0 pinned in [ml] so a sibling
  install (e.g. strands-agents) can't drag the version below the floor
  transformers 5.x requires (otherwise Kompress silently goes
  "unavailable").

== 2. /stats MCP aggregation ==

* headroom/proxy/cost.py: _aggregate_mcp_events() reads the cross-process
  shared events file the Headroom MCP server already writes to and
  surfaces summary.mcp with three new keys:
    - compressions       (count of headroom_compress invocations)
    - tokens_removed     (sum of input - output across those)
    - retrievals         (count of headroom_retrieve — the load-bearing
                          over-compression alarm; if it grows linearly
                          with turn count, lossy compressors are
                          dropping info the model actually needs)
  Defensive on every axis — missing MCP SDK, missing file, malformed
  events, read errors — never blocks /stats.

* examples/strands_bundle_demo.py: stats panel prints the new fields so
  the demo shows the full proxy-HTTP + MCP-tool story in one view.

== 3. Codex compression-failure fail-closed protection ==

Reported by Camille (2026-05-21): Codex threads were locking with
"ran out of room in the model's context window" after Headroom's
compression timed out on an oversized response.create frame and
forwarded the original ~1.7 MB frame to the upstream, which then
rejected it. Codex's auto-compact heuristic gates on the upstream-
reported total_usage_tokens (which Headroom had been shrinking on
earlier turns), so its compaction never fired and the thread locked.

Validated against open Codex issues (CLI + Desktop share codex-rs/core):
* #16068 — confirms compaction gates on total_usage_tokens,
  estimated_token_count is computed but only logged
* #19806 — confirms image token estimator unbounded, contributes to
  the same ContextManager.get_total_token_usage → auto-compaction chain

* headroom/proxy/helpers.py: decide_compression_failure_action() with a
  unit-tested decision matrix:
    - asyncio.TimeoutError                              → refuse, always
    - non-timeout failure + frame > 256 KiB (configurable) → refuse
    - non-timeout failure + small frame                 → forward (legacy)
  Operator escape hatches:
    - HEADROOM_WS_FAIL_OPEN_ON_COMPRESSION_FAILURE=1 restores legacy
    - HEADROOM_WS_COMPRESSION_FAIL_THRESHOLD_BYTES tunes the threshold

* headroom/proxy/handlers/openai.py (WS /v1/responses): consults the
  helper after compression failure. On refuse: close client websocket
  code 1009 with "headroom: compression <reason> — please compact
  context and retry" reason; set termination_cause for the outer
  lifecycle finally; return.

* headroom/proxy/handlers/openai.py (HTTP /v1/responses): same helper.
  On refuse: raise HTTPException(413) with a structured error body so
  FastAPI's HTTPException handler emits a clean 413. The existing
  `except HTTPException: raise` guard in this handler already ensures
  the 413 propagates without being swallowed by the 502 catch-all.

Anthropic /v1/messages NOT changed in this branch: no equivalent bug
report on Anthropic-protocol clients, Claude Code (Anthropic-owned)
handles context overflow via its own cache_control/ephemeral
primitives, and Cursor/Aider don't maintain the local-Y estimate the
Codex bug requires. Deferred until a real report lands; the patch is
a one-liner reusing the same helper.

== Tests + verification ==

* tests/test_backends/test_litellm_cache_stats.py — 3 tests pinning
  cache-stat surfacing across Anthropic/OpenAI dialects + backwards-
  compat for no-cache responses.
* tests/test_proxy/test_openai_backend_path.py — 5 tests (Bedrock cache
  fields, OpenAI fallback shape, CCR intercept with provider="openai",
  CCR re-raise on exception, streaming signature contract).
* tests/test_proxy/test_mcp_stats_aggregation.py — 5 tests pinning the
  aggregator across compress+retrieve mixes, empty events, unknown event
  types, missing token fields, and read failures.
* tests/test_proxy/test_compression_failure_action.py — 12 tests pinning
  the fail-closed decision matrix (timeout always refuses, small
  transient passes through, oversize refuses, env override variants,
  custom threshold, invalid threshold falls back, 0/negative ignored).

* examples/strands_bedrock_demo.py — model_id bumped from deprecated
  Claude 3 Haiku to Sonnet 4.5 (the deprecated model now errors on
  account access).
* examples/strands_via_proxy_demo.py — proxy + Bedrock cache + streaming
  smoke test.
* examples/strands_mcp_dispatch_test.py — pure MCP round-trip probe.
* examples/strands_bundle_demo.py — full Strands + HeadroomBundle E2E
  demo (this is the shape a real Strands user copies into their app).

Full pytest: 5327 passed, 178 skipped. The previously-failing
test_core_operations.py::TestAddBatch::test_add_batch_basic passes now
that the huggingface-hub pin in pyproject.toml unblocks transformers
imports.

E2E verified live against AWS Bedrock (Sonnet 4.5):
* cache_write=10,438 on turn A → cache_read=10,438 on turn B
* streaming SSE final usage frame carries cache_read_input_tokens
* 78.7% reduction on a 50 KB JSON tool_result via SmartCrusher (
  dispatched per-content-type by ContentRouter)
* Strands Agent + HeadroomBundle: model autonomously called
  headroom_compress + headroom_retrieve via MCP; CompressionStore
  round-trip succeeded; final answer correct.
2026-05-21 11:00:14 -07:00

255 lines
8.9 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""Strands agent + Headroom — drop-in compression demo.
This is what a real Strands user writes. The ONLY Headroom-specific
lines are the import and the bundle construction. Everything else
is normal Strands code.
The demo:
1. Defines a normal Strands `@tool` that returns a verbose JSON blob
(mimicking a real RAG or DB tool result).
2. Builds a normal Strands `Agent` with that tool + the bundle's MCP
tools (headroom_compress / headroom_retrieve / headroom_stats).
3. Sends ONE user query.
4. Lets the agent loop autonomously.
5. Prints the answer + the proxy /stats summary so you can SEE that
Headroom compressed the verbose tool output on the way to Bedrock.
The Headroom proxy is started as a background process here so the
script is self-contained. In production, the proxy runs as a
long-lived service (ECS / k8s / EC2) and the application code looks
exactly like what's below the `=== USER CODE ===` line.
Run
---
AWS_REGION=us-west-2 python examples/strands_bundle_demo.py
Cost
----
Sonnet 4.5 via Bedrock, multi-turn agent loop. Expect ~$0.020.10
per run depending on how many tool calls the model makes.
"""
from __future__ import annotations
import json
import os
import subprocess
import sys
import time
import urllib.request
from contextlib import suppress
from pathlib import Path
# ============================================================================
# Boilerplate: start the Headroom proxy as a child process for the demo.
# In production this is a long-lived service — none of this code is in your
# Strands app.
# ============================================================================
PROXY_PORT = 8787
PROXY_URL = f"http://127.0.0.1:{PROXY_PORT}"
def _start_proxy() -> subprocess.Popen[bytes]:
cmd = [
sys.executable,
"-m",
"headroom.cli",
"proxy",
"--backend",
"bedrock",
"--region",
os.environ.get("AWS_REGION", "us-west-2"),
"--port",
str(PROXY_PORT),
]
log_path = Path("/tmp/headroom_bundle_demo_proxy.log")
log = log_path.open("wb")
proc = subprocess.Popen(cmd, stdout=log, stderr=subprocess.STDOUT) # noqa: S603
print(f" → proxy started (pid={proc.pid}); log: {log_path}")
return proc
def _wait_for_proxy(timeout_s: float = 30.0) -> None:
deadline = time.time() + timeout_s
while time.time() < deadline:
try:
with urllib.request.urlopen(f"{PROXY_URL}/readyz", timeout=1) as r: # noqa: S310
if r.status == 200:
return
except Exception: # noqa: BLE001
time.sleep(0.5)
raise RuntimeError(f"Proxy did not become ready within {timeout_s}s")
def _stop_proxy(proc: subprocess.Popen[bytes]) -> None:
with suppress(ProcessLookupError):
proc.terminate()
try:
proc.wait(timeout=5)
except subprocess.TimeoutExpired:
proc.kill()
def _print_stats_panel() -> None:
try:
with urllib.request.urlopen(f"{PROXY_URL}/stats", timeout=5) as r: # noqa: S310
stats = json.loads(r.read())
except Exception: # noqa: BLE001
return
summary = stats.get("summary", {})
comp = summary.get("compression", {})
uncomp = summary.get("uncompressed_requests", {})
mcp = summary.get("mcp", {}) or {}
print()
print(" Proxy /stats summary")
print(" --------------------")
print(f" api_requests: {summary.get('api_requests', 0)}")
print(f" requests_compressed: {comp.get('requests_compressed', 0)}")
print(f" total_tokens_removed: {comp.get('total_tokens_removed', 0)}")
if comp.get("best_compression_pct"):
print(f" best_compression_pct: {comp['best_compression_pct']:.1f}%")
print(f" uncompressed reasons: {uncomp}")
# MCP-side work (headroom_compress / headroom_retrieve called by the LLM
# via Strands' MCP dispatcher). The retrievals counter is the
# over-compression alarm — if it grows linearly with turn count, our
# lossy compressors are dropping info the model actually needs.
print(f" mcp_compressions: {mcp.get('compressions', 0)}")
print(f" mcp_tokens_removed: {mcp.get('tokens_removed', 0)}")
print(f" ccr_retrievals_count: {mcp.get('retrievals', 0)}")
# ============================================================================
# === USER CODE ===
# Everything below is what a normal Strands user writes. The only
# Headroom-specific bits are the `HeadroomBundle` import and one
# construction call. The rest is vanilla Strands.
# ============================================================================
from strands import Agent, tool # noqa: E402 (kept under USER CODE banner for readability)
from strands.models.openai import OpenAIModel # noqa: E402
from headroom.integrations.strands import HeadroomBundle # noqa: E402
@tool
def search_documentation(query: str) -> str:
"""Search the documentation for articles matching `query`. Returns up to 30 results as JSON."""
# Mock data — pretend this hit a real search API. The point is the
# response is big and repetitive, which is the kind of tool output
# Headroom is designed to shrink.
articles = [
{
"id": f"doc-{i:04d}",
"title": f"{query.title()} Guide — Part {i + 1}",
"url": f"https://docs.example.com/{query}/article-{i + 1}",
"category": "tutorial" if i % 3 == 0 else "reference",
"snippet": (
f"This article covers {query} implementation in depth. "
f"It walks through setup, configuration, and common pitfalls. "
f"Section {i + 1} of the comprehensive guide series."
)
* 4,
"metadata": {
"author": "docs-team",
"tags": ["how-to", query, "production-ready"],
"last_updated": f"2026-04-{(i % 28) + 1:02d}",
},
}
for i in range(30)
]
return json.dumps(articles)
def run_agent_demo() -> None:
"""The actual Strands agent demo — looks like any other Strands app."""
# === The only Headroom-specific code in your app ===
bundle = HeadroomBundle(
proxy_url=PROXY_URL,
enable_serena_mcp=False, # disabled here so the demo runs fast
)
# === Normal Strands agent setup ===
model = OpenAIModel(
model_id="bedrock/us.anthropic.claude-sonnet-4-5-20250929-v1:0",
client_args={
"base_url": f"{PROXY_URL}/v1",
"api_key": "dummy-bedrock-uses-aws-creds-at-proxy",
"default_headers": {
"x-headroom-session-id": "demo-1",
"X-Client": "strands",
},
},
params={"max_tokens": 800, "temperature": 0.2},
)
agent = Agent(
model=model,
# bundle.tools = [Headroom MCP client]; we also add our local tool
tools=bundle.tools + [search_documentation],
system_prompt=(
"You are a documentation assistant. When the user asks about a "
"topic, use the search_documentation tool to look it up, then "
"answer concisely."
),
)
print("\n → sending user query (agent runs autonomously) ...")
user_query = (
"Search the documentation for 'authentication'. "
"Tell me how many results you got, then summarize the top 3 in one sentence each."
)
print(f" user: {user_query!r}\n")
t0 = time.time()
response = agent(user_query)
elapsed = time.time() - t0
print(f"\n agent response (after {elapsed:.1f}s):")
print(" " + "-" * 68)
for line in str(response).splitlines():
print(f" {line}")
print(" " + "-" * 68)
# ============================================================================
# Bootstrap
# ============================================================================
def main() -> int:
print("=" * 72)
print(" Strands agent + Headroom (drop-in compression)")
print("=" * 72)
print(f" proxy: {PROXY_URL} region: {os.environ.get('AWS_REGION', 'us-west-2')}")
print()
print(" Starting Headroom proxy (in production this is a long-lived service) ...")
proxy = _start_proxy()
try:
_wait_for_proxy()
print(" → proxy ready.")
run_agent_demo()
_print_stats_panel()
print("\n" + "=" * 72)
print(" Done. If requests_compressed > 0 above, Headroom shrunk the")
print(" verbose tool output on its way to Bedrock — automatically,")
print(" with no code changes in the agent.")
print("=" * 72)
return 0
except Exception as e: # noqa: BLE001
import traceback
print(f"\n ! FAILED: {type(e).__name__}: {e}")
traceback.print_exc()
return 1
finally:
print("\n → stopping proxy ...")
_stop_proxy(proxy)
if __name__ == "__main__":
sys.exit(main())