headroom/tests/test_proxy_gemini_native_integration.py
Ben Younes 9fde127534
fix(proxy): relocate stray system-role messages to the top-level system param (#765) (#1357)
## Description

On requests large enough to trigger compression, the proxy emitted an
upstream Anthropic request whose `messages[0]` had `role: "system"`.
Anthropic's Messages API rejects any `system` role inside `messages[]`:

```
400 invalid_request_error: "messages.0: use the top-level 'system' parameter for the initial system prompt"
```

The original request correctly carries its system prompt in the
top-level `system` parameter; a compression/transform/pipeline step
relocates the harness system block into `messages[0]`, so the request
fails outright (intermittent only because it requires a context large
enough to compress).

This adds a wire-contract guard in the Anthropic forwarder: as the
**last** step before sending upstream (after every transform, memory
injection, tool sort, and pipeline extension, covering both the Bedrock
and direct paths), any stray `role="system"` message is relocated out of
`messages[]` and merged back into the top-level `system` parameter.
Content order is preserved (existing system first, relocated content
after) and block-level `cache_control` survives. The guard is a no-op on
the common path (no system-role entry → inputs pass through unchanged).

Closes #765

## Type of Change

- [x] Bug fix (non-breaking change that fixes an issue)
- [ ] New feature (non-breaking change that adds functionality)
- [ ] Breaking change (fix or feature that would cause existing
functionality to change)
- [ ] Documentation update
- [ ] Performance improvement
- [ ] Code refactoring (no functional changes)

## Changes Made

- `headroom/proxy/helpers.py`: new pure helper
`relocate_system_messages_to_top_level(messages, system) ->
(clean_messages, new_system, changed)` plus `_system_message_to_blocks`.
Handles `system` being `None`/`str`/`list`, never drops content,
preserves order and content blocks.
- `headroom/proxy/handlers/anthropic.py`: invoke the guard just before
the byte-faithful forward block; on relocation, update
`body["messages"]`/`body["system"]`, mark the body mutated
(`system_role_relocated`) so the byte-faithful forwarder re-serializes,
and log a warning.
- `tests/test_proxy_handler_helpers.py`: 3 unit tests (relocate stray
system into top-level, append-to-existing-system order, no-op without a
system entry).
- `CHANGELOG.md`: Bug Fixes entry.

## Testing

- [x] Unit tests pass (`pytest`)
- [x] Linting passes (`ruff check .`)
- [x] Type checking passes (`mypy headroom`)
- [x] New tests added for new functionality
- [x] Manual testing performed

### Test Output

```text
$ uv run --extra dev python -m pytest tests/test_proxy_handler_helpers.py -q
29 passed in 4.95s

# Regression on the forward path (byte-faithful forwarding, system-prompt immutability, cache stability):
$ uv run --extra dev python -m pytest tests/test_proxy_handler_helpers.py tests/test_proxy_byte_faithful_forwarding.py tests/test_proxy_system_prompt_immutable.py tests/test_proxy_anthropic_cache_stability.py -q
90 passed, 15 warnings in 29.72s

$ uv run ruff check headroom/proxy/helpers.py headroom/proxy/handlers/anthropic.py tests/test_proxy_handler_helpers.py
All checks passed!

$ uv run ruff format --check ...   # 3 files already formatted
$ uv run mypy headroom/proxy/helpers.py headroom/proxy/handlers/anthropic.py
Success: no issues found in 2 source files
```

## Test verification (RED → GREEN)

The new tests exercise the guard directly and import the new helper at
module top, so reverting the production fix makes them fail at
collection.

**RED — production fix reverted (helper removed):**
```text
ImportError while importing test module 'tests/test_proxy_handler_helpers.py'.
E   ImportError: cannot import name 'relocate_system_messages_to_top_level' from 'headroom.proxy.helpers'
=========================== 1 error in 0.41s ===============================
```

**GREEN — production fix applied:**
```text
tests/test_proxy_handler_helpers.py ...                                  [100%]
======================= 3 passed, 26 deselected in 1.50s =======================
```

## Real Behavior Proof

- Environment: Python 3.13, `uv run` in this repo, branch
`fix/issue-765`.
- Exact command / steps: ran the guard on a body in the exact #765
failure shape — `system: None` and a `role="system"` harness block at
`messages[0]`:
- Observed result:
  ```text
  BEFORE: messages[0].role = system (Anthropic 400 trigger)
  changed       = True
  AFTER roles   = ['user', 'assistant']
system param = [{"type": "text", "text": "You are Claude Code.
<system-reminder>...</system-reminder>"}]
OK: no role=system in messages[]; system content preserved in top-level
param
  ```
The illegal `role="system"` entry is removed from `messages[]` and its
content lands in the top-level `system` parameter — exactly the body
Anthropic accepts.
- Not tested: a full live 250k+-token Claude Code session against the
real Anthropic API (needs a large live context + API key); the fix is
validated at the request-shaping boundary the 400 is raised on.

## Review Readiness

- [x] I have performed a self-review
- [x] This PR is ready for human review

## Checklist

- [x] My code follows the project's style guidelines
- [x] I have performed a self-review of my code
- [x] I have commented my code, particularly in hard-to-understand areas
- [x] I have made corresponding changes to the documentation
- [x] My changes generate no new warnings
- [x] I have added tests that prove my fix is effective or that my
feature works
- [x] New and existing unit tests pass locally with my changes
- [x] I have updated the CHANGELOG.md if applicable

## Additional Notes

The guard intentionally fires at the forwarder boundary rather than in
any single transform: the issue's captures show the relocation can
originate from the compression path, and pipeline extensions / hooks can
also mutate `messages` late. Enforcing Anthropic's wire contract once,
at the point the body is serialized upstream, fixes the 400 regardless
of which step introduced the stray entry and matches the architecture
invariant "never produce a `system`-role entry within `messages[]`".

---------

Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Co-authored-by: JerrettDavis <mxjerrett@gmail.com>
2026-08-13 11:52:09 -05:00

391 lines
16 KiB
Python

"""Integration tests for Gemini native API endpoint with real API calls.
These tests require a valid GEMINI_API_KEY environment variable.
They test the /v1beta/models/{model}:generateContent endpoint with compression.
Run with:
GEMINI_API_KEY=your-key pytest tests/test_proxy_gemini_native_integration.py -v
"""
import json
import os
import pytest
# Skip entire module if no API key
pytestmark = pytest.mark.skipif(
not os.environ.get("GEMINI_API_KEY"), reason="GEMINI_API_KEY not set"
)
pytest.importorskip("fastapi")
pytest.importorskip("httpx")
from fastapi.testclient import TestClient # noqa: E402
from headroom.proxy.server import ProxyConfig, create_app # noqa: E402
from tests._gemini_live import skip_if_gemini_quota_exhausted # noqa: E402
@pytest.fixture
def gemini_native_client():
"""Create test client for Gemini native API with optimization enabled."""
config = ProxyConfig(
optimize=True, # Enable compression
cache_enabled=False,
rate_limit_enabled=False,
cost_tracking_enabled=False,
)
app = create_app(config)
with TestClient(app) as client:
yield client
@pytest.fixture
def api_key():
"""Get Gemini API key from environment."""
return os.environ.get("GEMINI_API_KEY")
class TestGeminiNativeGenerateContent:
"""Test /v1beta/models/{model}:generateContent endpoint."""
def test_basic_generation(self, gemini_native_client, api_key):
"""Basic text generation works."""
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}",
json={"contents": [{"parts": [{"text": "What is 2+2? Reply with just the number."}]}]},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
data = response.json()
# Verify Gemini native response format
assert "candidates" in data
assert len(data["candidates"]) > 0
assert "content" in data["candidates"][0]
assert "parts" in data["candidates"][0]["content"]
text = data["candidates"][0]["content"]["parts"][0]["text"]
assert "4" in text
# Verify usage metadata
assert "usageMetadata" in data
assert "promptTokenCount" in data["usageMetadata"]
def test_with_system_instruction(self, gemini_native_client, api_key):
"""System instruction works correctly."""
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}",
json={
"contents": [{"parts": [{"text": "Hello"}]}],
"systemInstruction": {"parts": [{"text": "Always respond with exactly one word."}]},
},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
data = response.json()
text = data["candidates"][0]["content"]["parts"][0]["text"]
# Should be a short response due to system instruction
assert len(text.split()) <= 3
def test_multi_turn_conversation(self, gemini_native_client, api_key):
"""Multi-turn conversations maintain context."""
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}",
json={
"contents": [
{"role": "user", "parts": [{"text": "My name is TestUser456."}]},
{"role": "model", "parts": [{"text": "Nice to meet you, TestUser456!"}]},
{"role": "user", "parts": [{"text": "What is my name?"}]},
]
},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
data = response.json()
text = data["candidates"][0]["content"]["parts"][0]["text"].lower()
assert "testuser456" in text
def test_function_calling(self, gemini_native_client, api_key):
"""Function calling / tools work correctly."""
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}",
json={
"contents": [{"parts": [{"text": "What is the weather in Tokyo?"}]}],
"tools": [
{
"functionDeclarations": [
{
"name": "get_weather",
"description": "Get current weather for a location",
"parameters": {
"type": "object",
"properties": {
"location": {"type": "string", "description": "City name"}
},
"required": ["location"],
},
}
]
}
],
},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
data = response.json()
# Verify function call response
parts = data["candidates"][0]["content"]["parts"]
function_call = None
for part in parts:
if "functionCall" in part:
function_call = part["functionCall"]
break
assert function_call is not None
assert function_call["name"] == "get_weather"
assert "tokyo" in function_call["args"]["location"].lower()
def test_generation_config(self, gemini_native_client, api_key):
"""Generation config parameters are respected."""
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}",
json={
"contents": [{"parts": [{"text": "Write a very short poem about AI."}]}],
"generationConfig": {"maxOutputTokens": 50, "temperature": 0.1},
},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
data = response.json()
# Response should be limited by maxOutputTokens
assert data["usageMetadata"]["candidatesTokenCount"] <= 60 # Some buffer
class TestGeminiNativeCompression:
"""Test that compression works with Gemini native API."""
def test_compression_on_model_message(self, gemini_native_client, api_key):
"""Large data in model message gets compressed."""
# Create large JSON data (simulating tool output)
items = [
{"id": i, "name": f"Item {i}", "desc": f"Description for item {i}"} for i in range(100)
]
tool_output = json.dumps(items)
# Send as model message (like tool returning data)
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}",
json={
"contents": [
{"role": "user", "parts": [{"text": "Get items from database"}]},
{"role": "model", "parts": [{"text": f"Here are the results:\n{tool_output}"}]},
{"role": "user", "parts": [{"text": "How many items are there?"}]},
]
},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
data = response.json()
text = data["candidates"][0]["content"]["parts"][0]["text"]
# Model should correctly count the items
assert "100" in text
# Check that compression happened via stats
stats = gemini_native_client.get("/stats").json()
# At least some tokens should have been saved
assert stats["tokens"]["saved"] >= 0 # May or may not compress depending on size
def test_user_messages_protected(self, gemini_native_client, api_key):
"""User messages are not compressed (by design)."""
# Large data in user message
items = [{"id": i} for i in range(50)]
user_data = json.dumps(items)
# First request with data in user message
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}",
json={
"contents": [
{"role": "user", "parts": [{"text": f"Analyze this data: {user_data}"}]}
]
},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
# The request should succeed - user messages are protected from compression
class TestGeminiNativeStats:
"""Test that proxy stats track Gemini native requests correctly."""
def test_stats_track_gemini_provider(self, gemini_native_client, api_key):
"""Stats show requests under 'gemini' provider."""
# Make a request
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}",
json={"contents": [{"parts": [{"text": "Hi"}]}]},
)
skip_if_gemini_quota_exhausted(response)
stats = gemini_native_client.get("/stats").json()
assert "gemini" in stats["requests"]["by_provider"]
assert stats["requests"]["by_provider"]["gemini"] >= 1
def test_stats_track_model(self, gemini_native_client, api_key):
"""Stats track the specific model used."""
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}",
json={"contents": [{"parts": [{"text": "Hi"}]}]},
)
skip_if_gemini_quota_exhausted(response)
stats = gemini_native_client.get("/stats").json()
assert "gemini-2.0-flash" in stats["requests"]["by_model"]
class TestGeminiNativeErrorHandling:
"""Test error handling for Gemini native API."""
def test_invalid_api_key(self, gemini_native_client):
"""Invalid API key returns appropriate error."""
response = gemini_native_client.post(
"/v1beta/models/gemini-2.0-flash:generateContent?key=invalid-key-123",
json={"contents": [{"parts": [{"text": "Hi"}]}]},
)
assert response.status_code >= 400
def test_invalid_model(self, gemini_native_client, api_key):
"""Invalid model returns appropriate error."""
response = gemini_native_client.post(
f"/v1beta/models/nonexistent-model-xyz:generateContent?key={api_key}",
json={"contents": [{"parts": [{"text": "Hi"}]}]},
)
assert response.status_code >= 400
def test_empty_contents(self, gemini_native_client, api_key):
"""Empty contents handled gracefully."""
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}", json={"contents": []}
)
skip_if_gemini_quota_exhausted(response)
# Should either return error or handle gracefully
assert response.status_code in [200, 400]
class TestGeminiNativeHeaderAuth:
"""Test authentication via x-goog-api-key header."""
def test_header_auth(self, gemini_native_client, api_key):
"""API key in header works."""
response = gemini_native_client.post(
"/v1beta/models/gemini-2.0-flash:generateContent",
headers={"x-goog-api-key": api_key},
json={"contents": [{"parts": [{"text": "Hi"}]}]},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
class TestGeminiNativeCountTokens:
"""Test /v1beta/models/{model}:countTokens endpoint with compression."""
def test_count_tokens_basic(self, gemini_native_client, api_key):
"""Basic token counting works."""
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:countTokens?key={api_key}",
json={"contents": [{"parts": [{"text": "Hello, world!"}]}]},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
data = response.json()
# Verify response format
assert "totalTokens" in data
assert isinstance(data["totalTokens"], int)
assert data["totalTokens"] > 0
def test_count_tokens_with_system_instruction(self, gemini_native_client, api_key):
"""Token counting includes system instruction."""
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:countTokens?key={api_key}",
json={
"contents": [{"parts": [{"text": "Hello"}]}],
"systemInstruction": {"parts": [{"text": "You are a helpful assistant."}]},
},
)
skip_if_gemini_quota_exhausted(response)
# Note: systemInstruction may not be supported by countTokens in all versions
assert response.status_code in [200, 400]
if response.status_code == 200:
data = response.json()
assert "totalTokens" in data
assert data["totalTokens"] > 0
def test_count_tokens_reflects_compression(self, gemini_native_client, api_key):
"""Token count reflects compressed content size."""
# Create large repetitive JSON data that should compress
items = [
{
"id": i,
"name": f"Item {i}",
"description": f"This is the description for item number {i}",
}
for i in range(100)
]
tool_output = json.dumps(items)
# Count tokens with large data in model message (which gets compressed)
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:countTokens?key={api_key}",
json={
"contents": [
{"role": "user", "parts": [{"text": "Get items from database"}]},
{"role": "model", "parts": [{"text": f"Here are the results:\n{tool_output}"}]},
{"role": "user", "parts": [{"text": "Summarize these items"}]},
]
},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
data = response.json()
# Verify we got a token count
assert "totalTokens" in data
compressed_tokens = data["totalTokens"]
assert compressed_tokens > 0
# Check stats to verify compression was applied
stats = gemini_native_client.get("/stats").json()
# The request should have been tracked
assert stats["requests"]["by_provider"].get("gemini", 0) >= 1
def test_count_tokens_multi_turn(self, gemini_native_client, api_key):
"""Token counting works for multi-turn conversations."""
response = gemini_native_client.post(
f"/v1beta/models/gemini-2.0-flash:countTokens?key={api_key}",
json={
"contents": [
{"role": "user", "parts": [{"text": "My name is Alice."}]},
{"role": "model", "parts": [{"text": "Nice to meet you, Alice!"}]},
{"role": "user", "parts": [{"text": "What is my name?"}]},
]
},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
data = response.json()
assert "totalTokens" in data
assert data["totalTokens"] > 0
def test_count_tokens_header_auth(self, gemini_native_client, api_key):
"""API key in header works for countTokens."""
response = gemini_native_client.post(
"/v1beta/models/gemini-2.0-flash:countTokens",
headers={"x-goog-api-key": api_key},
json={"contents": [{"parts": [{"text": "Hello"}]}]},
)
skip_if_gemini_quota_exhausted(response)
assert response.status_code == 200
data = response.json()
assert "totalTokens" in data