Mechanical only, and separated from the judgment calls that follow so the reviewable changes are not buried in a 98-file whitespace diff. 227 automatic fixes: 60 blank lines carrying whitespace, 60 unsorted import blocks, 34 Optional[X] to X | None, 28 unused imports, 16 deprecated typing imports, 12 datetime.timezone.utc to datetime.UTC, and assorted smaller modernisations. Then `ruff format` over src and tests: 98 files reformatted, 35 already conforming. No file among the unused-import findings defines __all__ or is an __init__.py, so nothing here removes a re-export. `make test`: 658 passed, unchanged from HEAD. Two things observed while verifying, neither addressed here: `pytest tests/` cannot collect — tests/e2e/test_orchestration_e2e.py uses an `e2e` marker that is not registered, and the config is strict about markers. This fails identically at HEAD, so it predates this change; `make test` passes because it ignores tests/e2e, tests/integration and tests/contracts. test_tatlock_tool_call_logging_calculator is flaky. It failed once in a full run with these changes and passed on the next, passes in isolation with them, and fails in isolation at HEAD. It is order- or timing-dependent, not a regression from this commit — established by running the full suite both ways rather than by reasoning about which change could have caused it. Co-Authored-By: Claude <noreply@anthropic.com>
306 lines
9.5 KiB
Python
306 lines
9.5 KiB
Python
"""
|
|
Tests for chat completions streaming wrapper.
|
|
|
|
Tests that the wrapper correctly:
|
|
- Wraps Responses API
|
|
- Enables reasoning automatically
|
|
- Streams reasoning via reasoning_content field (DeepSeek R1 format)
|
|
- Streams both reasoning and content
|
|
"""
|
|
|
|
import json
|
|
|
|
import pytest
|
|
from httpx import AsyncClient
|
|
|
|
from src.chat import constants
|
|
|
|
|
|
@pytest.mark.unit
|
|
@pytest.mark.asyncio
|
|
async def test_streaming_wrapper_enables_reasoning(async_client: AsyncClient):
|
|
"""Test that streaming wrapper automatically enables reasoning via reasoning_content."""
|
|
request_data = {
|
|
"model": "lorem-tester",
|
|
"messages": [{"role": "user", "content": "Test message"}],
|
|
"stream": True,
|
|
}
|
|
|
|
chunks_received = []
|
|
reasoning_content_found = False
|
|
|
|
async with async_client.stream(
|
|
"POST",
|
|
"/v1/chat/completions",
|
|
json=request_data,
|
|
timeout=20.0,
|
|
) as response:
|
|
assert response.status_code == 200
|
|
assert response.headers["content-type"] == "text/event-stream; charset=utf-8"
|
|
|
|
async for line in response.aiter_lines():
|
|
if not line.strip():
|
|
continue
|
|
|
|
if line.startswith("data: "):
|
|
data_str = line[6:].strip()
|
|
if data_str == "[DONE]":
|
|
break
|
|
|
|
try:
|
|
chunk = json.loads(data_str)
|
|
chunks_received.append(chunk)
|
|
|
|
# Check for reasoning_content in delta (DeepSeek R1 format)
|
|
if "choices" in chunk and len(chunk["choices"]) > 0:
|
|
delta = chunk["choices"][0].get("delta", {})
|
|
reasoning = delta.get("reasoning_content")
|
|
if reasoning:
|
|
reasoning_content_found = True
|
|
|
|
except json.JSONDecodeError:
|
|
pass
|
|
|
|
# Should have received chunks
|
|
assert len(chunks_received) > 0
|
|
|
|
# Should have found reasoning_content (reasoning enabled automatically)
|
|
assert reasoning_content_found, "Expected reasoning_content in streaming output"
|
|
|
|
|
|
@pytest.mark.unit
|
|
@pytest.mark.asyncio
|
|
async def test_streaming_wrapper_reasoning_before_content(async_client: AsyncClient):
|
|
"""Test that reasoning_content comes before regular content."""
|
|
request_data = {
|
|
"model": "lorem-tester",
|
|
"messages": [{"role": "user", "content": "Explain something"}],
|
|
"stream": True,
|
|
}
|
|
|
|
chunk_types = [] # Track order: 'reasoning' or 'content'
|
|
|
|
async with async_client.stream(
|
|
"POST",
|
|
"/v1/chat/completions",
|
|
json=request_data,
|
|
timeout=20.0,
|
|
) as response:
|
|
assert response.status_code == 200
|
|
|
|
async for line in response.aiter_lines():
|
|
if not line.strip():
|
|
continue
|
|
|
|
if line.startswith("data: "):
|
|
data_str = line[6:].strip()
|
|
if data_str == "[DONE]":
|
|
break
|
|
|
|
try:
|
|
chunk = json.loads(data_str)
|
|
if "choices" in chunk and len(chunk["choices"]) > 0:
|
|
delta = chunk["choices"][0].get("delta", {})
|
|
reasoning = delta.get("reasoning_content")
|
|
content = delta.get("content")
|
|
|
|
if reasoning:
|
|
chunk_types.append("reasoning")
|
|
if content:
|
|
chunk_types.append("content")
|
|
|
|
except json.JSONDecodeError:
|
|
pass
|
|
|
|
# Verify reasoning comes before content
|
|
if "reasoning" in chunk_types and "content" in chunk_types:
|
|
first_reasoning = chunk_types.index("reasoning")
|
|
first_content = chunk_types.index("content")
|
|
assert first_reasoning < first_content, "reasoning_content should come before content"
|
|
|
|
|
|
@pytest.mark.unit
|
|
@pytest.mark.asyncio
|
|
async def test_streaming_wrapper_proper_chunk_structure(async_client: AsyncClient):
|
|
"""Test that streaming chunks have proper structure."""
|
|
request_data = {
|
|
"model": "lorem-tester",
|
|
"messages": [{"role": "user", "content": "Hello"}],
|
|
"temperature": 0.8,
|
|
"stream": True,
|
|
}
|
|
|
|
first_chunk = None
|
|
last_chunk = None
|
|
chunk_count = 0
|
|
|
|
async with async_client.stream(
|
|
"POST",
|
|
"/v1/chat/completions",
|
|
json=request_data,
|
|
timeout=20.0,
|
|
) as response:
|
|
assert response.status_code == 200
|
|
|
|
async for line in response.aiter_lines():
|
|
if not line.strip():
|
|
continue
|
|
|
|
if line.startswith("data: "):
|
|
data_str = line[6:].strip()
|
|
if data_str == "[DONE]":
|
|
break
|
|
|
|
try:
|
|
chunk = json.loads(data_str)
|
|
chunk_count += 1
|
|
|
|
# Verify chunk structure
|
|
assert "id" in chunk
|
|
assert "object" in chunk
|
|
assert chunk["object"] == constants.CHAT_COMPLETION_CHUNK_OBJECT
|
|
assert "created" in chunk
|
|
assert "model" in chunk
|
|
assert chunk["model"] == "lorem-tester"
|
|
assert "choices" in chunk
|
|
assert len(chunk["choices"]) == 1
|
|
|
|
choice = chunk["choices"][0]
|
|
assert "index" in choice
|
|
assert choice["index"] == 0
|
|
assert "delta" in choice
|
|
|
|
if first_chunk is None:
|
|
first_chunk = chunk
|
|
last_chunk = chunk
|
|
|
|
except json.JSONDecodeError:
|
|
pass
|
|
|
|
# Verify we got chunks
|
|
assert chunk_count > 0
|
|
assert first_chunk is not None
|
|
assert last_chunk is not None
|
|
|
|
# First chunk should have role
|
|
assert first_chunk["choices"][0]["delta"].get("role") == constants.ROLE_ASSISTANT
|
|
|
|
# Last chunk should have finish_reason
|
|
assert last_chunk["choices"][0].get("finish_reason") == constants.FINISH_REASON_STOP
|
|
|
|
|
|
@pytest.mark.unit
|
|
@pytest.mark.asyncio
|
|
async def test_streaming_wrapper_with_system_message(async_client: AsyncClient):
|
|
"""Test streaming with system message."""
|
|
request_data = {
|
|
"model": "lorem-tester",
|
|
"messages": [
|
|
{"role": "system", "content": "You are a helpful assistant."},
|
|
{"role": "user", "content": "Hello"},
|
|
],
|
|
"stream": True,
|
|
}
|
|
|
|
chunks_received = []
|
|
|
|
async with async_client.stream(
|
|
"POST",
|
|
"/v1/chat/completions",
|
|
json=request_data,
|
|
timeout=20.0,
|
|
) as response:
|
|
assert response.status_code == 200
|
|
|
|
async for line in response.aiter_lines():
|
|
if not line.strip():
|
|
continue
|
|
|
|
if line.startswith("data: "):
|
|
data_str = line[6:].strip()
|
|
if data_str == "[DONE]":
|
|
break
|
|
|
|
try:
|
|
chunk = json.loads(data_str)
|
|
chunks_received.append(chunk)
|
|
except json.JSONDecodeError:
|
|
pass
|
|
|
|
# Should handle system message properly
|
|
assert len(chunks_received) > 0
|
|
# First chunk should still have assistant role
|
|
assert chunks_received[0]["choices"][0]["delta"].get("role") == constants.ROLE_ASSISTANT
|
|
|
|
|
|
@pytest.mark.unit
|
|
@pytest.mark.asyncio
|
|
async def test_streaming_wrapper_pipeline_prefix(async_client: AsyncClient):
|
|
"""Test streaming with pipeline prefix in model name."""
|
|
request_data = {
|
|
"model": "some_pipeline.lorem-tester",
|
|
"messages": [{"role": "user", "content": "Test"}],
|
|
"stream": True,
|
|
}
|
|
|
|
chunks_received = []
|
|
|
|
async with async_client.stream(
|
|
"POST",
|
|
"/v1/chat/completions",
|
|
json=request_data,
|
|
timeout=20.0,
|
|
) as response:
|
|
assert response.status_code == 200
|
|
|
|
async for line in response.aiter_lines():
|
|
if not line.strip():
|
|
continue
|
|
|
|
if line.startswith("data: "):
|
|
data_str = line[6:].strip()
|
|
if data_str == "[DONE]":
|
|
break
|
|
|
|
try:
|
|
chunk = json.loads(data_str)
|
|
chunks_received.append(chunk)
|
|
# Model should keep original name (with prefix)
|
|
assert chunk["model"] == "some_pipeline.lorem-tester"
|
|
except json.JSONDecodeError:
|
|
pass
|
|
|
|
assert len(chunks_received) > 0
|
|
|
|
|
|
@pytest.mark.unit
|
|
@pytest.mark.asyncio
|
|
async def test_streaming_wrapper_non_streaming_fallback(async_client: AsyncClient):
|
|
"""Test that non-streaming request works through wrapper."""
|
|
request_data = {
|
|
"model": "lorem-tester",
|
|
"messages": [{"role": "user", "content": "Hello"}],
|
|
"stream": False, # Non-streaming
|
|
}
|
|
|
|
response = await async_client.post("/v1/chat/completions", json=request_data, timeout=20.0)
|
|
|
|
assert response.status_code == 200
|
|
data = response.json()
|
|
|
|
# Verify structure
|
|
assert "id" in data
|
|
assert "object" in data
|
|
assert data["object"] == constants.CHAT_COMPLETION_OBJECT
|
|
assert "choices" in data
|
|
assert len(data["choices"]) == 1
|
|
|
|
choice = data["choices"][0]
|
|
assert "message" in choice
|
|
assert choice["message"]["role"] == constants.ROLE_ASSISTANT
|
|
assert choice["message"]["content"] # Should have content
|
|
|
|
# Should have <think> tags in content (reasoning enabled)
|
|
assert "<think>" in choice["message"]["content"]
|
|
assert "</think>" in choice["message"]["content"]
|