Build and Push / build (release) Successful in 52s
### Added - Environment-aware configuration: - Auto-selected logging (DEBUG for dev, WARNING for prod) - Auto-selected default user (llm_tester for dev isolation) - User context logging at request entry - Direct delegation bypass: - Pure memory/librarian requests skip Tatlock LLM - Reduces latency for memory-only requests - Text-based delegation fallback: - Parse [DELEGATE:agent] patterns from LLM output - Sequential and parallel execution support - Comprehensive E2E test suite: - 22 orchestration tests with QdrantVerifier - assert_llm_behavior() for flexible pattern matching - Tests for memory, delegation, isolation, scenarios ### Fixed - Unit test mocks for streaming (async generator) - Temporal context handling in tests - LLM non-determinism with pytest.xfail() - Streaming test timeouts increased 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
632 lines
19 KiB
Python
632 lines
19 KiB
Python
"""
|
|
Response service for creating responses.
|
|
|
|
Handles both streaming and non-streaming response generation.
|
|
Tracks conversation history for analytics and future vector memory.
|
|
Integrates with Steward preprocessing for Phase 2 two-tier architecture.
|
|
"""
|
|
|
|
import time
|
|
import secrets
|
|
from typing import AsyncGenerator
|
|
|
|
from src.agents.registry import ModelRegistry
|
|
from src.responses.schemas import (
|
|
Response,
|
|
ResponseRequest,
|
|
ResponseUsage,
|
|
MessageOutputItem,
|
|
ReasoningOutputItem,
|
|
FunctionCallOutputItem,
|
|
OutputTextContent,
|
|
)
|
|
from src.responses.streaming import StreamingCoordinator
|
|
from src.responses.history import ConversationHistory
|
|
from src.responses.context import ContextWindow
|
|
from src.core.preprocessing import preprocess_request
|
|
from src.core.tool_tracking import ToolCallTracker
|
|
from src.core.logging_config import get_logger
|
|
from src.agents.steward.schemas import StewardRecommendation
|
|
|
|
import re
|
|
import asyncio
|
|
|
|
logger = get_logger(__name__)
|
|
|
|
|
|
async def _execute_single_delegation(
|
|
agent_name: str,
|
|
task: str,
|
|
tracker: "ToolCallTracker",
|
|
) -> tuple[str, str]:
|
|
"""
|
|
Execute a single delegation to an agent.
|
|
|
|
Args:
|
|
agent_name: Name of agent (biographer, librarian)
|
|
task: Task description
|
|
tracker: Tool call tracker
|
|
|
|
Returns:
|
|
tuple: (agent_name, result_summary)
|
|
"""
|
|
import time
|
|
start_time = time.time()
|
|
|
|
if agent_name == "biographer":
|
|
from src.agents.delegation import delegate_to_biographer
|
|
result = await delegate_to_biographer(task=task)
|
|
duration = time.time() - start_time
|
|
await tracker.track_call("delegate_to_biographer", duration)
|
|
return (agent_name, result.output)
|
|
|
|
elif agent_name == "librarian":
|
|
from src.agents.delegation import delegate_to_librarian
|
|
result = await delegate_to_librarian(task=task)
|
|
duration = time.time() - start_time
|
|
await tracker.track_call("delegate_to_librarian", duration)
|
|
return (agent_name, result.output)
|
|
|
|
else:
|
|
return (agent_name, f"Unknown agent: {agent_name}")
|
|
|
|
|
|
async def _handle_text_delegation(
|
|
response: str,
|
|
tracker: "ToolCallTracker",
|
|
conversation_id: str
|
|
) -> str:
|
|
"""
|
|
Handle text-based delegation fallback.
|
|
|
|
When Tatlock outputs [DELEGATE:agent] task="..." instead of calling
|
|
the actual function, we parse and execute it here.
|
|
|
|
Supports multiple delegations in the same response:
|
|
- Sequential: Run one after another in order
|
|
- Parallel: Run all at once if [PARALLEL] prefix is present
|
|
|
|
Patterns:
|
|
[DELEGATE:biographer] task="Remember something"
|
|
[DELEGATE:librarian] task="Search for something"
|
|
[PARALLEL][DELEGATE:biographer] task="..." [DELEGATE:librarian] task="..."
|
|
|
|
Args:
|
|
response: Tatlock's response text
|
|
tracker: Tool call tracker for metrics
|
|
conversation_id: Current conversation ID
|
|
|
|
Returns:
|
|
str: Either the original response or the delegation result(s)
|
|
"""
|
|
# Pattern 1: [DELEGATE:agent_name] task="task description"
|
|
# Pattern 2: Delegate:"agent_name", "task":"task description" (LLM variant)
|
|
# Pattern 3: delegate_to_agent(task="...") (function-like text)
|
|
patterns = [
|
|
r'\[DELEGATE:(\w+)\]\s*task=["\']([^"\']+)["\']',
|
|
r'[Dd]elegate[:\s]*["\']?(\w+)["\']?,?\s*["\']?task["\']?[:\s]*["\']([^"\']+)["\']',
|
|
r'delegate_to_(\w+)\s*\(\s*task\s*=\s*["\']([^"\']+)["\']',
|
|
]
|
|
|
|
matches = []
|
|
for pattern in patterns:
|
|
found = re.findall(pattern, response)
|
|
if found:
|
|
matches.extend(found)
|
|
break # Use first matching pattern
|
|
|
|
if not matches:
|
|
# No text delegation found, return original response
|
|
return response
|
|
|
|
logger.info(
|
|
"text_delegation_detected",
|
|
delegation_count=len(matches),
|
|
agents=[m[0] for m in matches],
|
|
conversation_id=conversation_id,
|
|
)
|
|
|
|
# Check if parallel execution is requested
|
|
is_parallel = "[PARALLEL]" in response.upper()
|
|
|
|
try:
|
|
if is_parallel and len(matches) > 1:
|
|
# Execute all delegations in parallel
|
|
logger.info(
|
|
"executing_parallel_delegations",
|
|
count=len(matches),
|
|
conversation_id=conversation_id,
|
|
)
|
|
tasks = [
|
|
_execute_single_delegation(agent.lower(), task, tracker)
|
|
for agent, task in matches
|
|
]
|
|
results = await asyncio.gather(*tasks, return_exceptions=True)
|
|
|
|
# Combine results
|
|
summaries = []
|
|
for agent_name, result in results:
|
|
if isinstance(result, Exception):
|
|
summaries.append(f"**{agent_name}**: Error - {result}")
|
|
else:
|
|
summaries.append(f"**{agent_name}**: {result}")
|
|
|
|
return "\n\n".join(summaries)
|
|
|
|
else:
|
|
# Execute sequentially
|
|
summaries = []
|
|
for agent_name, task in matches:
|
|
agent_name = agent_name.lower()
|
|
logger.info(
|
|
"executing_sequential_delegation",
|
|
agent=agent_name,
|
|
task_preview=task[:50],
|
|
conversation_id=conversation_id,
|
|
)
|
|
try:
|
|
_, result = await _execute_single_delegation(
|
|
agent_name, task, tracker
|
|
)
|
|
summaries.append(result)
|
|
except Exception as e:
|
|
logger.error(
|
|
"delegation_failed",
|
|
agent=agent_name,
|
|
error=str(e),
|
|
conversation_id=conversation_id,
|
|
)
|
|
summaries.append(
|
|
f"I apologize, sir. Delegation to {agent_name} failed: {e}"
|
|
)
|
|
|
|
return "\n\n".join(summaries)
|
|
|
|
except Exception as e:
|
|
logger.error(
|
|
"text_delegation_failed",
|
|
error=str(e),
|
|
conversation_id=conversation_id,
|
|
)
|
|
return f"I apologize, sir. I encountered an error processing delegations: {e}"
|
|
|
|
|
|
async def _direct_delegation(
|
|
user_message: str,
|
|
recommendation: "StewardRecommendation",
|
|
tracker: "ToolCallTracker",
|
|
conversation_id: str,
|
|
) -> str:
|
|
"""
|
|
Directly delegate to expert agents, bypassing Tatlock.
|
|
|
|
When Steward recommends ONLY delegation agents (biographer/librarian),
|
|
we skip Tatlock's LLM call and delegate directly. This works around
|
|
models that don't reliably call tools.
|
|
|
|
Args:
|
|
user_message: User's request
|
|
recommendation: Steward's recommendation
|
|
tracker: Tool call tracker
|
|
conversation_id: Conversation ID
|
|
|
|
Returns:
|
|
str: Combined results from delegations
|
|
"""
|
|
logger.info(
|
|
"direct_delegation_triggered",
|
|
agents=recommendation.recommended_capabilities,
|
|
conversation_id=conversation_id,
|
|
)
|
|
|
|
results = []
|
|
for agent in recommendation.recommended_capabilities:
|
|
try:
|
|
agent_name, result = await _execute_single_delegation(
|
|
agent, user_message, tracker
|
|
)
|
|
results.append(result)
|
|
logger.info(
|
|
"direct_delegation_complete",
|
|
agent=agent_name,
|
|
result_preview=result[:100] if result else "empty",
|
|
conversation_id=conversation_id,
|
|
)
|
|
except Exception as e:
|
|
logger.error(
|
|
"direct_delegation_failed",
|
|
agent=agent,
|
|
error=str(e),
|
|
conversation_id=conversation_id,
|
|
)
|
|
results.append(f"I apologize, sir. Delegation to {agent} failed: {e}")
|
|
|
|
return "\n\n".join(results) if results else "I apologize, sir. No delegation results available."
|
|
|
|
|
|
# Global conversation history tracker
|
|
# In production, this would be backed by a database or Redis
|
|
_conversation_history = ConversationHistory(max_turns=20)
|
|
|
|
# Global context window manager (4096 token default)
|
|
_context_window = ContextWindow(max_tokens=4096)
|
|
|
|
|
|
def generate_id() -> str:
|
|
"""Generate unique ID for responses."""
|
|
return secrets.token_hex(16)
|
|
|
|
|
|
def _calculate_usage(input_messages: list[dict], output_items: list) -> ResponseUsage:
|
|
"""
|
|
Calculate token usage for the response.
|
|
|
|
For now: approximate token counting.
|
|
Future: Use tiktoken or similar for accurate counting.
|
|
|
|
Args:
|
|
input_messages: Input messages
|
|
output_items: Output items generated
|
|
|
|
Returns:
|
|
ResponseUsage: Token usage statistics
|
|
"""
|
|
# Approximate input tokens (chars / 4)
|
|
input_text = " ".join(str(msg) for msg in input_messages)
|
|
input_tokens = len(input_text) // 4
|
|
|
|
# Approximate output tokens
|
|
output_tokens = 0
|
|
reasoning_tokens = 0
|
|
|
|
for item in output_items:
|
|
# Check if it's a schema object (has summary/content attributes directly)
|
|
if isinstance(item, ReasoningOutputItem):
|
|
reasoning_text = " ".join(item.summary)
|
|
reasoning_tokens += len(reasoning_text) // 4
|
|
elif isinstance(item, MessageOutputItem):
|
|
message_text = item.content[0].text
|
|
output_tokens += len(message_text) // 4
|
|
elif isinstance(item, FunctionCallOutputItem):
|
|
func_text = item.arguments
|
|
output_tokens += len(func_text) // 4
|
|
elif hasattr(item, 'type'):
|
|
# Agent OutputItem objects (backward compatibility)
|
|
if item.type == "reasoning":
|
|
reasoning_text = " ".join(item.data.get("summary", []))
|
|
reasoning_tokens += len(reasoning_text) // 4
|
|
elif item.type == "message":
|
|
message_text = item.data["content"][0]["text"]
|
|
output_tokens += len(message_text) // 4
|
|
elif item.type == "function_call":
|
|
func_text = item.data["arguments"]
|
|
output_tokens += len(func_text) // 4
|
|
|
|
total_tokens = input_tokens + output_tokens + reasoning_tokens
|
|
|
|
return ResponseUsage(
|
|
input_tokens=input_tokens,
|
|
output_tokens=output_tokens,
|
|
reasoning_tokens=reasoning_tokens,
|
|
total_tokens=total_tokens
|
|
)
|
|
|
|
|
|
async def create_response(request: ResponseRequest) -> Response:
|
|
"""
|
|
Create non-streaming response.
|
|
|
|
Tracks conversation history if conversation_id is provided in metadata.
|
|
|
|
Args:
|
|
request: Response request
|
|
|
|
Returns:
|
|
Response: Complete response object
|
|
|
|
Example:
|
|
request = ResponseRequest(
|
|
model="lorem-tester",
|
|
input=[{"role": "user", "content": "Hello"}],
|
|
reasoning={"effort": "medium", "summary": "auto"},
|
|
metadata={"conversation_id": "conv_abc123"} # Optional
|
|
)
|
|
response = await create_response(request)
|
|
"""
|
|
# Get or generate conversation ID
|
|
conversation_id = await _conversation_history.get_conversation_id(request)
|
|
|
|
# Strip pipeline prefix if present (e.g., "pipeline.model" -> "model")
|
|
model_id = request.model
|
|
if "." in model_id:
|
|
model_id = model_id.split(".", 1)[1]
|
|
|
|
# Get agent for model
|
|
agent = ModelRegistry.get_agent(model_id)
|
|
|
|
# Collect all output items from agent
|
|
output_items = []
|
|
async for item in agent.generate_response(
|
|
messages=request.input,
|
|
reasoning=request.reasoning,
|
|
tools=request.tools,
|
|
temperature=request.temperature,
|
|
max_tokens=request.max_output_tokens,
|
|
stop=request.stop,
|
|
):
|
|
output_items.append(item)
|
|
|
|
# Convert agent OutputItems to schema OutputItems
|
|
converted_items = _convert_output_items(output_items)
|
|
|
|
# Calculate token usage
|
|
usage = _calculate_usage(request.input, output_items)
|
|
|
|
response = Response(
|
|
id=f"resp_{generate_id()}",
|
|
created_at=int(time.time()),
|
|
model=request.model,
|
|
status="completed",
|
|
output=converted_items,
|
|
usage=usage
|
|
)
|
|
|
|
# Track conversation history (for analytics and future vector memory)
|
|
await _conversation_history.add_response(conversation_id, response)
|
|
|
|
return response
|
|
|
|
|
|
async def create_response_with_steward(request: ResponseRequest) -> Response:
|
|
"""
|
|
Create response using Steward preprocessing (Phase 2 flow).
|
|
|
|
This is the two-tier architecture where:
|
|
1. Steward analyzes the request and recommends capabilities
|
|
2. Tatlock runs with scoped tools based on recommendations
|
|
3. Tool usage is tracked for benchmarking
|
|
|
|
Args:
|
|
request: Response request
|
|
|
|
Returns:
|
|
Response: Complete response object with Steward analysis included
|
|
|
|
Example:
|
|
request = ResponseRequest(
|
|
model="tatlock",
|
|
input=[{"role": "user", "content": "What's sqrt(144)?"}],
|
|
metadata={"conversation_id": "conv_abc123"}
|
|
)
|
|
response = await create_response_with_steward(request)
|
|
"""
|
|
# Get or generate conversation ID
|
|
conversation_id = await _conversation_history.get_conversation_id(request)
|
|
|
|
# Extract user message and conversation history
|
|
user_message = ""
|
|
for msg in reversed(request.input):
|
|
if msg.get("role") == "user":
|
|
user_message = msg.get("content", "")
|
|
break
|
|
|
|
# Conversation history is all messages except the current one
|
|
conversation_history = request.input[:-1] if len(request.input) > 1 else []
|
|
|
|
logger.info(
|
|
"creating_response_with_steward",
|
|
user_message_preview=user_message[:100],
|
|
history_length=len(conversation_history),
|
|
conversation_id=conversation_id,
|
|
)
|
|
|
|
# Phase 1: Steward preprocessing
|
|
enriched = await preprocess_request(
|
|
user_message,
|
|
conversation_history=conversation_history,
|
|
conversation_id=conversation_id,
|
|
)
|
|
|
|
# Phase 2: Initialize tool tracker
|
|
tracker = ToolCallTracker(
|
|
recommended_capabilities=enriched.recommendation.recommended_capabilities,
|
|
conversation_id=conversation_id,
|
|
)
|
|
|
|
# Phase 3: Check if direct delegation is recommended
|
|
# If Steward recommends ONLY delegation agents (biographer/librarian),
|
|
# skip Tatlock and delegate directly
|
|
delegation_only = all(
|
|
cap in ("biographer", "librarian")
|
|
for cap in enriched.recommendation.recommended_capabilities
|
|
) and enriched.recommendation.recommended_capabilities
|
|
|
|
if delegation_only:
|
|
tatlock_response = await _direct_delegation(
|
|
user_message, enriched.recommendation, tracker, conversation_id
|
|
)
|
|
else:
|
|
# Phase 3a: Run Tatlock with scoped tools
|
|
from src.agents.tatlock import TatlockAgent
|
|
tatlock = TatlockAgent()
|
|
|
|
tatlock_response = await tatlock.run_with_scoped_tools(
|
|
user_message=user_message,
|
|
steward_note=enriched.steward_note,
|
|
scoped_tools=enriched.scoped_tools,
|
|
message_history=conversation_history,
|
|
tool_tracker=tracker,
|
|
)
|
|
|
|
# Phase 3b: Check for text-based delegation fallback
|
|
# If Tatlock outputs [DELEGATE:...] instead of calling the function,
|
|
# we parse and execute it here
|
|
tatlock_response = await _handle_text_delegation(
|
|
tatlock_response, tracker, conversation_id
|
|
)
|
|
|
|
# Phase 4: Finalize tool tracking
|
|
await tracker.finalize()
|
|
|
|
# Build response output items
|
|
output_items = []
|
|
|
|
# Add Steward reasoning as a reasoning output item
|
|
output_items.append(ReasoningOutputItem(
|
|
id=f"reasoning_{generate_id()}",
|
|
summary=[
|
|
"🎩 Steward's Analysis:",
|
|
enriched.steward_reasoning,
|
|
],
|
|
status="completed"
|
|
))
|
|
|
|
# Add Tatlock's message
|
|
output_items.append(MessageOutputItem(
|
|
id=f"msg_{generate_id()}",
|
|
role="assistant",
|
|
content=[OutputTextContent(
|
|
type="output_text",
|
|
text=tatlock_response,
|
|
annotations=[]
|
|
)],
|
|
status="completed"
|
|
))
|
|
|
|
# Calculate usage (approximate)
|
|
usage = _calculate_usage(request.input, output_items)
|
|
|
|
response = Response(
|
|
id=f"resp_{generate_id()}",
|
|
created_at=int(time.time()),
|
|
model=request.model,
|
|
status="completed",
|
|
output=output_items,
|
|
usage=usage
|
|
)
|
|
|
|
# Track conversation history
|
|
await _conversation_history.add_response(conversation_id, response)
|
|
|
|
logger.info(
|
|
"response_with_steward_complete",
|
|
response_id=response.id,
|
|
recommended_capabilities=enriched.recommendation.recommended_capabilities,
|
|
tool_summary=tracker.get_summary(),
|
|
)
|
|
|
|
return response
|
|
|
|
|
|
async def create_response_stream(
|
|
request: ResponseRequest
|
|
) -> AsyncGenerator[dict, None]:
|
|
"""
|
|
Create streaming response.
|
|
|
|
Yields SSE-formatted events for streaming to client.
|
|
|
|
Args:
|
|
request: Response request
|
|
|
|
Yields:
|
|
dict: SSE event dict with 'event' and 'data' keys
|
|
|
|
Example:
|
|
async for event in create_response_stream(request):
|
|
# event = {"event": "response.output_text.delta", "data": "..."}
|
|
yield event
|
|
"""
|
|
coordinator = StreamingCoordinator()
|
|
|
|
async for event in coordinator.stream_response(request):
|
|
yield {
|
|
"event": event.event,
|
|
"data": event.model_dump_json()
|
|
}
|
|
|
|
|
|
async def get_conversation_history(conversation_id: str) -> list[Response]:
|
|
"""
|
|
Get conversation history for a given conversation ID.
|
|
|
|
Args:
|
|
conversation_id: Conversation identifier
|
|
|
|
Returns:
|
|
list: List of Response objects
|
|
"""
|
|
return await _conversation_history.get_history(conversation_id)
|
|
|
|
|
|
async def get_conversation_stats() -> dict:
|
|
"""
|
|
Get conversation history statistics.
|
|
|
|
Returns:
|
|
dict: Statistics including total conversations tracked
|
|
"""
|
|
return {
|
|
"total_conversations": await _conversation_history.get_conversation_count(),
|
|
"max_turns_per_conversation": _conversation_history._max_turns
|
|
}
|
|
|
|
|
|
async def clear_conversation(conversation_id: str) -> bool:
|
|
"""
|
|
Clear a specific conversation history.
|
|
|
|
Args:
|
|
conversation_id: Conversation identifier
|
|
|
|
Returns:
|
|
bool: True if conversation was cleared
|
|
"""
|
|
return await _conversation_history.clear_conversation(conversation_id)
|
|
|
|
|
|
def get_context_window() -> ContextWindow:
|
|
"""
|
|
Get the global context window manager.
|
|
|
|
Returns:
|
|
ContextWindow: Context window manager instance
|
|
"""
|
|
return _context_window
|
|
|
|
|
|
def _convert_output_items(items: list) -> list:
|
|
"""
|
|
Convert agent OutputItem objects to schema OutputItem objects.
|
|
|
|
Args:
|
|
items: List of agent OutputItem objects
|
|
|
|
Returns:
|
|
list: List of schema OutputItem objects
|
|
"""
|
|
converted = []
|
|
|
|
for item in items:
|
|
if item.type == "message":
|
|
converted.append(MessageOutputItem(
|
|
id=item.id,
|
|
content=[OutputTextContent(**c) for c in item.data["content"]],
|
|
status=item.data.get("status", "completed")
|
|
))
|
|
elif item.type == "reasoning":
|
|
converted.append(ReasoningOutputItem(
|
|
id=item.id,
|
|
summary=item.data["summary"],
|
|
status=item.data.get("status", "completed")
|
|
))
|
|
elif item.type == "function_call":
|
|
converted.append(FunctionCallOutputItem(
|
|
id=item.id,
|
|
name=item.data["name"],
|
|
arguments=item.data["arguments"],
|
|
status=item.data.get("status", "completed")
|
|
))
|
|
|
|
return converted
|