feat: integrate tracing throughout request pipeline
Instrument the full request flow with trace spans for debugging: - Wrap expert delegations (librarian/biographer/housekeeper) in spans - Add orchestrate and synthesize spans to TatlockAgent - Trace Steward analysis in preprocessing - Start/end traces in response service with context management - Simplify router by moving context handling to service layer - Include tracing router in debug mode - Remove benchmark recording from tool_tracking and steward service 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
+145
-84
@@ -13,6 +13,7 @@ from enum import Enum
|
||||
from typing import AsyncGenerator, Callable, Optional, Any
|
||||
|
||||
from src.core.logging_config import get_logger
|
||||
from src.core.tracing import trace_span, SpanType
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
@@ -239,38 +240,58 @@ async def delegate_to_librarian(
|
||||
has_context=bool(context),
|
||||
)
|
||||
|
||||
try:
|
||||
# Use run() not run_stream() - avoids Ollama bug
|
||||
output = await run_librarian(task=task, context=context)
|
||||
async with trace_span(
|
||||
"delegate_to_librarian",
|
||||
SpanType.EXPERT,
|
||||
metadata={
|
||||
"expert": "librarian",
|
||||
"task_preview": task[:100],
|
||||
"has_context": bool(context),
|
||||
},
|
||||
) as span:
|
||||
try:
|
||||
# Use run() not run_stream() - avoids Ollama bug
|
||||
output = await run_librarian(task=task, context=context)
|
||||
|
||||
logger.info(
|
||||
"delegation_to_librarian_completed",
|
||||
task=task[:50],
|
||||
output_length=len(output),
|
||||
)
|
||||
logger.info(
|
||||
"delegation_to_librarian_completed",
|
||||
task=task[:50],
|
||||
output_length=len(output),
|
||||
)
|
||||
|
||||
return DelegationResult(
|
||||
expert_name="librarian",
|
||||
task=task,
|
||||
success=True,
|
||||
output=output,
|
||||
)
|
||||
if span:
|
||||
span.metadata["success"] = True
|
||||
span.metadata["output_length"] = len(output)
|
||||
span.details["task"] = task
|
||||
span.details["context"] = context[:500] if context else None
|
||||
span.details["result_preview"] = output[:1000]
|
||||
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"delegation_to_librarian_error",
|
||||
task=task[:50],
|
||||
error=str(e),
|
||||
exc_info=True,
|
||||
)
|
||||
return DelegationResult(
|
||||
expert_name="librarian",
|
||||
task=task,
|
||||
success=True,
|
||||
output=output,
|
||||
)
|
||||
|
||||
return DelegationResult(
|
||||
expert_name="librarian",
|
||||
task=task,
|
||||
success=False,
|
||||
output="",
|
||||
error=str(e),
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"delegation_to_librarian_error",
|
||||
task=task[:50],
|
||||
error=str(e),
|
||||
exc_info=True,
|
||||
)
|
||||
|
||||
if span:
|
||||
span.metadata["success"] = False
|
||||
span.details["error"] = str(e)
|
||||
|
||||
return DelegationResult(
|
||||
expert_name="librarian",
|
||||
task=task,
|
||||
success=False,
|
||||
output="",
|
||||
error=str(e),
|
||||
)
|
||||
|
||||
|
||||
async def delegate_to_biographer(
|
||||
@@ -317,38 +338,58 @@ async def delegate_to_biographer(
|
||||
has_context=bool(context),
|
||||
)
|
||||
|
||||
try:
|
||||
# Use run() not run_stream() - avoids Ollama bug
|
||||
output = await run_biographer(task=task, context=context)
|
||||
async with trace_span(
|
||||
"delegate_to_biographer",
|
||||
SpanType.EXPERT,
|
||||
metadata={
|
||||
"expert": "biographer",
|
||||
"task_preview": task[:100],
|
||||
"has_context": bool(context),
|
||||
},
|
||||
) as span:
|
||||
try:
|
||||
# Use run() not run_stream() - avoids Ollama bug
|
||||
output = await run_biographer(task=task, context=context)
|
||||
|
||||
logger.info(
|
||||
"delegation_to_biographer_completed",
|
||||
task=task[:50],
|
||||
output_length=len(output),
|
||||
)
|
||||
logger.info(
|
||||
"delegation_to_biographer_completed",
|
||||
task=task[:50],
|
||||
output_length=len(output),
|
||||
)
|
||||
|
||||
return DelegationResult(
|
||||
expert_name="biographer",
|
||||
task=task,
|
||||
success=True,
|
||||
output=output,
|
||||
)
|
||||
if span:
|
||||
span.metadata["success"] = True
|
||||
span.metadata["output_length"] = len(output)
|
||||
span.details["task"] = task
|
||||
span.details["context"] = context[:500] if context else None
|
||||
span.details["result_preview"] = output[:1000]
|
||||
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"delegation_to_biographer_error",
|
||||
task=task[:50],
|
||||
error=str(e),
|
||||
exc_info=True,
|
||||
)
|
||||
return DelegationResult(
|
||||
expert_name="biographer",
|
||||
task=task,
|
||||
success=True,
|
||||
output=output,
|
||||
)
|
||||
|
||||
return DelegationResult(
|
||||
expert_name="biographer",
|
||||
task=task,
|
||||
success=False,
|
||||
output="",
|
||||
error=str(e),
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"delegation_to_biographer_error",
|
||||
task=task[:50],
|
||||
error=str(e),
|
||||
exc_info=True,
|
||||
)
|
||||
|
||||
if span:
|
||||
span.metadata["success"] = False
|
||||
span.details["error"] = str(e)
|
||||
|
||||
return DelegationResult(
|
||||
expert_name="biographer",
|
||||
task=task,
|
||||
success=False,
|
||||
output="",
|
||||
error=str(e),
|
||||
)
|
||||
|
||||
|
||||
async def delegate_to_housekeeper(
|
||||
@@ -394,38 +435,58 @@ async def delegate_to_housekeeper(
|
||||
has_context=bool(context),
|
||||
)
|
||||
|
||||
try:
|
||||
# Use run() not run_stream() - avoids Ollama bug
|
||||
output = await run_housekeeper(task=task, context=context)
|
||||
async with trace_span(
|
||||
"delegate_to_housekeeper",
|
||||
SpanType.EXPERT,
|
||||
metadata={
|
||||
"expert": "housekeeper",
|
||||
"task_preview": task[:100],
|
||||
"has_context": bool(context),
|
||||
},
|
||||
) as span:
|
||||
try:
|
||||
# Use run() not run_stream() - avoids Ollama bug
|
||||
output = await run_housekeeper(task=task, context=context)
|
||||
|
||||
logger.info(
|
||||
"delegation_to_housekeeper_completed",
|
||||
task=task[:50],
|
||||
output_length=len(output),
|
||||
)
|
||||
logger.info(
|
||||
"delegation_to_housekeeper_completed",
|
||||
task=task[:50],
|
||||
output_length=len(output),
|
||||
)
|
||||
|
||||
return DelegationResult(
|
||||
expert_name="housekeeper",
|
||||
task=task,
|
||||
success=True,
|
||||
output=output,
|
||||
)
|
||||
if span:
|
||||
span.metadata["success"] = True
|
||||
span.metadata["output_length"] = len(output)
|
||||
span.details["task"] = task
|
||||
span.details["context"] = context[:500] if context else None
|
||||
span.details["result_preview"] = output[:1000]
|
||||
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"delegation_to_housekeeper_error",
|
||||
task=task[:50],
|
||||
error=str(e),
|
||||
exc_info=True,
|
||||
)
|
||||
return DelegationResult(
|
||||
expert_name="housekeeper",
|
||||
task=task,
|
||||
success=True,
|
||||
output=output,
|
||||
)
|
||||
|
||||
return DelegationResult(
|
||||
expert_name="housekeeper",
|
||||
task=task,
|
||||
success=False,
|
||||
output="",
|
||||
error=str(e),
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"delegation_to_housekeeper_error",
|
||||
task=task[:50],
|
||||
error=str(e),
|
||||
exc_info=True,
|
||||
)
|
||||
|
||||
if span:
|
||||
span.metadata["success"] = False
|
||||
span.details["error"] = str(e)
|
||||
|
||||
return DelegationResult(
|
||||
expert_name="housekeeper",
|
||||
task=task,
|
||||
success=False,
|
||||
output="",
|
||||
error=str(e),
|
||||
)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
"""
|
||||
Steward service layer.
|
||||
|
||||
Provides high-level interface for request analysis with logging,
|
||||
benchmarking, and error handling.
|
||||
Provides high-level interface for request analysis with logging
|
||||
and error handling.
|
||||
|
||||
Parses plain text recommendations into structured data.
|
||||
Includes memory pre-fetch for user context injection.
|
||||
@@ -10,7 +10,6 @@ Includes memory pre-fetch for user context injection.
|
||||
import re
|
||||
from typing import Any, Optional
|
||||
|
||||
from src.core.benchmarks import PerformanceBenchmark, get_benchmark_store
|
||||
from src.core.household_registry import get_household_registry
|
||||
from src.core.logging_config import get_logger, log_operation
|
||||
from src.core.memory_service import memory_service
|
||||
@@ -285,8 +284,7 @@ async def analyze_request(
|
||||
This is the main entry point for Steward analysis. It:
|
||||
1. Calls the Steward agent with full conversation history
|
||||
2. Logs the operation with timing
|
||||
3. Records performance benchmarks to Redis
|
||||
4. Returns structured recommendations
|
||||
3. Returns structured recommendations
|
||||
|
||||
Args:
|
||||
user_request: The current user message to analyze
|
||||
@@ -365,23 +363,6 @@ async def analyze_request(
|
||||
reasoning=analysis_text[:200], # First 200 chars
|
||||
)
|
||||
|
||||
# Record performance benchmark
|
||||
if log_ctx.get("duration_seconds"):
|
||||
benchmark = PerformanceBenchmark(
|
||||
operation="steward_analysis",
|
||||
duration_seconds=log_ctx["duration_seconds"],
|
||||
success=True,
|
||||
recommendation_count=len(recommendation.recommended_capabilities),
|
||||
confidence=None, # Could add confidence scoring in future
|
||||
conversation_id=conversation_id,
|
||||
metadata={
|
||||
"complexity": recommendation.estimated_complexity,
|
||||
"has_context": recommendation.conversation_context.has_previous_context,
|
||||
"missing_capabilities": recommendation.missing_capabilities is not None,
|
||||
},
|
||||
)
|
||||
await get_benchmark_store().record(benchmark)
|
||||
|
||||
return recommendation
|
||||
|
||||
except Exception as e:
|
||||
|
||||
+56
-2
@@ -20,6 +20,11 @@ from src.agents.tatlock_core.tools import (
|
||||
)
|
||||
from src.core.config import config
|
||||
from src.core.logging_config import get_logger
|
||||
from src.core.tracing import (
|
||||
start_span, end_span, get_current_span,
|
||||
add_tool_spans_from_messages,
|
||||
SpanType, SpanStatus,
|
||||
)
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
@@ -433,7 +438,7 @@ class TatlockAgent(AgentInterface):
|
||||
steward_note: Note from Steward (prepended to request, invisible to user)
|
||||
scoped_tools: List of tool definitions from household registry
|
||||
message_history: Conversation history in PydanticAI format
|
||||
tool_tracker: Optional tool call tracker for benchmarking
|
||||
tool_tracker: Optional tool call tracker for analysis
|
||||
|
||||
Returns:
|
||||
str: Tatlock's response text
|
||||
@@ -627,7 +632,7 @@ class TatlockAgent(AgentInterface):
|
||||
steward_note: Note from Steward (invisible to user)
|
||||
scoped_tools: List of tool definitions from household registry
|
||||
message_history: Conversation history
|
||||
tool_tracker: Optional tool call tracker for benchmarking
|
||||
tool_tracker: Optional tool call tracker for analysis
|
||||
|
||||
Returns:
|
||||
dict with:
|
||||
@@ -655,6 +660,16 @@ class TatlockAgent(AgentInterface):
|
||||
history_length=len(message_history),
|
||||
)
|
||||
|
||||
# Start tracing span for orchestration phase
|
||||
orchestrate_span = start_span(
|
||||
"tatlock_orchestrate",
|
||||
SpanType.TATLOCK,
|
||||
metadata={
|
||||
"scoped_tool_count": len(scoped_tools),
|
||||
"tool_names": [getattr(t, '__name__', str(t)) for t in scoped_tools[:5]],
|
||||
},
|
||||
)
|
||||
|
||||
# Create a fresh agent instance with scoped tools only
|
||||
clean_host = self.ollama_host.rstrip('/')
|
||||
base_url = f"{clean_host}/v1"
|
||||
@@ -731,6 +746,23 @@ class TatlockAgent(AgentInterface):
|
||||
tool_output_count=len(tool_outputs),
|
||||
)
|
||||
|
||||
# Add tool-level spans from result messages
|
||||
if orchestrate_span:
|
||||
add_tool_spans_from_messages(result.new_messages(), orchestrate_span)
|
||||
|
||||
# End orchestration span with results
|
||||
end_span(
|
||||
orchestrate_span,
|
||||
metadata_update={
|
||||
"tools_called": tools_called,
|
||||
"expert_count": len(expert_results),
|
||||
"tool_output_count": len(tool_outputs),
|
||||
},
|
||||
details_update={
|
||||
"steward_note_preview": steward_note[:500] if steward_note else None,
|
||||
},
|
||||
)
|
||||
|
||||
return {
|
||||
"tools_called": tools_called,
|
||||
"expert_results": expert_results,
|
||||
@@ -769,6 +801,16 @@ class TatlockAgent(AgentInterface):
|
||||
tool_count=len(orchestration_results.get("tool_outputs", {})),
|
||||
)
|
||||
|
||||
# Start tracing span for synthesis phase
|
||||
synthesize_span = start_span(
|
||||
"tatlock_synthesize",
|
||||
SpanType.TATLOCK,
|
||||
metadata={
|
||||
"expert_count": len(orchestration_results.get("expert_results", {})),
|
||||
"tool_output_count": len(orchestration_results.get("tool_outputs", {})),
|
||||
},
|
||||
)
|
||||
|
||||
# Build synthesis prompt with all available information
|
||||
synthesis_parts = []
|
||||
synthesis_parts.append(f"The user asked: {user_message}")
|
||||
@@ -841,6 +883,18 @@ class TatlockAgent(AgentInterface):
|
||||
response_preview=result.output[:100],
|
||||
)
|
||||
|
||||
# End synthesis span with result
|
||||
end_span(
|
||||
synthesize_span,
|
||||
metadata_update={
|
||||
"response_length": len(result.output),
|
||||
},
|
||||
details_update={
|
||||
"synthesis_prompt": synthesis_prompt[:1000],
|
||||
"response_preview": result.output[:500],
|
||||
},
|
||||
)
|
||||
|
||||
return result.output
|
||||
|
||||
async def get_capabilities(self) -> dict:
|
||||
|
||||
@@ -11,6 +11,7 @@ from src.agents.steward import analyze_request, format_steward_note
|
||||
from src.agents.steward.schemas import StewardRecommendation
|
||||
from src.core.household_registry import get_household_registry
|
||||
from src.core.logging_config import get_logger
|
||||
from src.core.tracing import trace_span, SpanType
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
@@ -93,12 +94,32 @@ async def preprocess_request(
|
||||
conversation_id=conversation_id,
|
||||
)
|
||||
|
||||
# Call Steward with full conversation history
|
||||
recommendation = await analyze_request(
|
||||
enriched_request,
|
||||
conversation_history=conversation_history,
|
||||
conversation_id=conversation_id,
|
||||
)
|
||||
# Call Steward with full conversation history (traced)
|
||||
async with trace_span(
|
||||
"steward_analysis",
|
||||
SpanType.STEWARD,
|
||||
metadata={
|
||||
"request_preview": user_request[:100],
|
||||
"history_length": len(conversation_history),
|
||||
},
|
||||
) as span:
|
||||
recommendation = await analyze_request(
|
||||
enriched_request,
|
||||
conversation_history=conversation_history,
|
||||
conversation_id=conversation_id,
|
||||
)
|
||||
|
||||
# Update span with results
|
||||
if span:
|
||||
span.metadata.update({
|
||||
"recommended_capabilities": recommendation.recommended_capabilities,
|
||||
"complexity": recommendation.estimated_complexity,
|
||||
"has_memory_context": bool(recommendation.memory_context),
|
||||
"has_conversation_context": recommendation.conversation_context.has_previous_context,
|
||||
})
|
||||
span.details["reasoning"] = recommendation.reasoning
|
||||
if recommendation.enriched_query:
|
||||
span.details["enriched_query"] = recommendation.enriched_query
|
||||
|
||||
# Format note for Tatlock (includes conversation context)
|
||||
steward_note = await format_steward_note(recommendation)
|
||||
|
||||
@@ -1,13 +1,11 @@
|
||||
"""
|
||||
Tool call tracking and benchmarking.
|
||||
Tool call tracking.
|
||||
|
||||
Tracks which tools are recommended by the Steward versus which tools
|
||||
are actually used by Tatlock, recording benchmarks for analysis.
|
||||
are actually used by Tatlock for debugging and analysis.
|
||||
"""
|
||||
from datetime import datetime, timezone
|
||||
from typing import Optional
|
||||
|
||||
from src.core.benchmarks import PerformanceBenchmark, get_benchmark_store
|
||||
from src.core.logging_config import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
@@ -15,7 +13,7 @@ logger = get_logger(__name__)
|
||||
|
||||
class ToolCallTracker:
|
||||
"""
|
||||
Tracks tool calls for benchmarking and accuracy analysis.
|
||||
Tracks tool calls for accuracy analysis.
|
||||
|
||||
Compares Steward's recommendations with Tatlock's actual tool usage
|
||||
to measure recommendation accuracy.
|
||||
@@ -53,6 +51,10 @@ class ToolCallTracker:
|
||||
return tool_name.replace("delegate_to_", "")
|
||||
return tool_name
|
||||
|
||||
def log_call(self, message: str):
|
||||
"""Log a tool call message (for UI display)."""
|
||||
logger.debug("tool_call_message", message=message)
|
||||
|
||||
async def track_call(self, tool_name: str, duration: float):
|
||||
"""
|
||||
Record a tool call with timing.
|
||||
@@ -78,23 +80,6 @@ class ToolCallTracker:
|
||||
recommended=list(self.recommended_capabilities),
|
||||
)
|
||||
|
||||
# Record benchmark to Redis
|
||||
benchmark = PerformanceBenchmark(
|
||||
timestamp=datetime.now(timezone.utc),
|
||||
operation="tool_call",
|
||||
duration_seconds=duration,
|
||||
success=True, # If we got here, the call succeeded
|
||||
tool_name=tool_name,
|
||||
was_recommended=was_recommended,
|
||||
was_actually_used=True,
|
||||
conversation_id=self.conversation_id,
|
||||
metadata={
|
||||
"recommended_capabilities": list(self.recommended_capabilities),
|
||||
},
|
||||
)
|
||||
|
||||
await get_benchmark_store().record(benchmark)
|
||||
|
||||
logger.debug(
|
||||
"tool_call_tracked",
|
||||
tool_name=tool_name,
|
||||
@@ -124,24 +109,6 @@ class ToolCallTracker:
|
||||
conversation_id=self.conversation_id,
|
||||
)
|
||||
|
||||
# Record benchmarks for unused recommendations
|
||||
for tool_name in unused_tools:
|
||||
benchmark = PerformanceBenchmark(
|
||||
timestamp=datetime.now(timezone.utc),
|
||||
operation="tool_call",
|
||||
duration_seconds=0.0, # Not used
|
||||
success=True,
|
||||
tool_name=tool_name,
|
||||
was_recommended=True,
|
||||
was_actually_used=False,
|
||||
conversation_id=self.conversation_id,
|
||||
metadata={
|
||||
"recommended_capabilities": list(self.recommended_capabilities),
|
||||
"reason": "recommended_but_unused",
|
||||
},
|
||||
)
|
||||
await get_benchmark_store().record(benchmark)
|
||||
|
||||
# Log summary
|
||||
total_calls = sum(len(durations) for durations in self.actual_calls.values())
|
||||
logger.info(
|
||||
|
||||
+8
-2
@@ -23,6 +23,7 @@ from src.core.exceptions import AppException
|
||||
from src.core.logging_config import get_logger
|
||||
from src.core.router import router as core_router
|
||||
from src.core.startup import initialize_application
|
||||
from src.core.tracing_router import router as tracing_router
|
||||
from src.models.router import router as models_router
|
||||
from src.responses.router import router as responses_router
|
||||
|
||||
@@ -45,7 +46,7 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
|
||||
environment=config.ENVIRONMENT.value,
|
||||
ollama_host=str(config.OLLAMA_HOST),
|
||||
ollama_model=config.OLLAMA_DEFAULT_MODEL,
|
||||
redis_url=config.redis_url,
|
||||
redis_url=config.redis_memory_url,
|
||||
log_format=config.log_format,
|
||||
)
|
||||
|
||||
@@ -90,7 +91,12 @@ def create_application() -> FastAPI:
|
||||
application.include_router(chat_router, prefix=config.API_PREFIX)
|
||||
application.include_router(models_router, prefix=config.API_PREFIX)
|
||||
application.include_router(responses_router, prefix=config.API_PREFIX) # Responses API
|
||||
|
||||
|
||||
# Conditionally include tracing router (only in debug mode)
|
||||
if config.DEBUG:
|
||||
application.include_router(tracing_router)
|
||||
logger.info("tracing_router_enabled")
|
||||
|
||||
return application
|
||||
|
||||
|
||||
|
||||
+5
-73
@@ -10,7 +10,6 @@ from sse_starlette.sse import EventSourceResponse
|
||||
from src.responses import service
|
||||
from src.responses.schemas import ResponseRequest, Response
|
||||
from src.core.exceptions import ModelNotFoundError, AppException
|
||||
from src.core.context import current_user, current_conversation, get_default_user
|
||||
from src.core.logging_config import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
@@ -37,77 +36,16 @@ async def create_response(
|
||||
|
||||
Returns:
|
||||
Response object or SSE stream
|
||||
|
||||
Example non-streaming request:
|
||||
POST /v1/responses
|
||||
{
|
||||
"model": "lorem-tester",
|
||||
"input": [{"role": "user", "content": "Hello"}],
|
||||
"reasoning": {"effort": "medium", "summary": "auto"},
|
||||
"stream": false
|
||||
}
|
||||
|
||||
Example streaming request:
|
||||
POST /v1/responses
|
||||
{
|
||||
"model": "lorem-tester",
|
||||
"input": [{"role": "user", "content": "Hello"}],
|
||||
"stream": true
|
||||
}
|
||||
|
||||
Response format (non-streaming):
|
||||
{
|
||||
"id": "resp_...",
|
||||
"object": "response",
|
||||
"created_at": 1733529600,
|
||||
"model": "lorem-tester",
|
||||
"status": "completed",
|
||||
"output": [
|
||||
{
|
||||
"type": "reasoning",
|
||||
"id": "rs_...",
|
||||
"summary": ["Analyzing...", "Considering..."]
|
||||
},
|
||||
{
|
||||
"type": "message",
|
||||
"id": "msg_...",
|
||||
"role": "assistant",
|
||||
"content": [{"type": "output_text", "text": "Lorem ipsum..."}]
|
||||
}
|
||||
],
|
||||
"usage": {
|
||||
"input_tokens": 10,
|
||||
"output_tokens": 50,
|
||||
"reasoning_tokens": 20,
|
||||
"total_tokens": 80
|
||||
}
|
||||
}
|
||||
|
||||
Streaming format (SSE):
|
||||
event: response.reasoning_summary_text.delta
|
||||
data: {"delta": "Analyzing..."}
|
||||
|
||||
event: response.output_text.delta
|
||||
data: {"delta": "Lorem"}
|
||||
|
||||
event: response.done
|
||||
data: {"response": {...}}
|
||||
"""
|
||||
# Set request context (propagates through all async calls)
|
||||
effective_user = request.user or get_default_user()
|
||||
user_token = current_user.set(effective_user)
|
||||
conv_id = request.metadata.get("conversation_id") if request.metadata else None
|
||||
conv_token = current_conversation.set(conv_id)
|
||||
|
||||
logger.info(
|
||||
"response_request_received",
|
||||
model=request.model,
|
||||
user=effective_user,
|
||||
conversation_id=conv_id,
|
||||
user=request.user,
|
||||
streaming=request.stream,
|
||||
)
|
||||
|
||||
try:
|
||||
# Check if this is a Tatlock request - use Steward preprocessing (Phase 2)
|
||||
# Check if this is a Tatlock request - use Steward preprocessing
|
||||
model_id = request.model
|
||||
if "." in model_id:
|
||||
model_id = model_id.split(".", 1)[1]
|
||||
@@ -116,21 +54,20 @@ async def create_response(
|
||||
|
||||
if request.stream:
|
||||
logger.info("Streaming response requested")
|
||||
|
||||
if use_steward:
|
||||
logger.info("Streaming with Steward preprocessing for Tatlock request")
|
||||
# Use Steward + Tatlock streaming (Milestone 3.5)
|
||||
from src.responses.streaming import StreamingCoordinator
|
||||
coordinator = StreamingCoordinator()
|
||||
return EventSourceResponse(
|
||||
coordinator.stream_response_with_steward(request)
|
||||
)
|
||||
else:
|
||||
# Regular streaming for non-Tatlock models
|
||||
return EventSourceResponse(
|
||||
service.create_response_stream(request)
|
||||
)
|
||||
|
||||
# Use appropriate service method
|
||||
# Non-streaming response
|
||||
if use_steward:
|
||||
logger.info("Using Steward preprocessing for Tatlock request")
|
||||
return await service.create_response_with_steward(request)
|
||||
@@ -148,8 +85,3 @@ async def create_response(
|
||||
except Exception as e:
|
||||
logger.error(f"Unexpected error: {e}", exc_info=True)
|
||||
raise HTTPException(status_code=500, detail="Internal server error")
|
||||
|
||||
finally:
|
||||
# Reset context (important for connection reuse)
|
||||
current_user.reset(user_token)
|
||||
current_conversation.reset(conv_token)
|
||||
|
||||
+251
-140
@@ -26,6 +26,8 @@ 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.core.tracing import start_trace, end_trace, start_span, SpanType
|
||||
from src.core.context import current_user, current_conversation, get_default_user
|
||||
from src.agents.steward.schemas import StewardRecommendation
|
||||
|
||||
import re
|
||||
@@ -34,6 +36,29 @@ import asyncio
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
def _extract_user_input(input_data) -> str:
|
||||
"""Extract user input text from request input for tracing."""
|
||||
if isinstance(input_data, str):
|
||||
return input_data
|
||||
elif isinstance(input_data, list) and input_data:
|
||||
last_msg = input_data[-1]
|
||||
if isinstance(last_msg, dict):
|
||||
return last_msg.get("content", str(last_msg))
|
||||
return str(last_msg)
|
||||
return ""
|
||||
|
||||
|
||||
def _extract_response_preview(response: Response) -> str:
|
||||
"""Extract response preview text for tracing."""
|
||||
if response.output:
|
||||
for item in response.output:
|
||||
if hasattr(item, 'content'):
|
||||
for content in item.content:
|
||||
if hasattr(content, 'text'):
|
||||
return content.text[:200]
|
||||
return ""
|
||||
|
||||
|
||||
async def _execute_single_delegation(
|
||||
agent_name: str,
|
||||
task: str,
|
||||
@@ -405,45 +430,88 @@ async def create_response(request: ResponseRequest) -> Response:
|
||||
# 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]
|
||||
# Set context for tracing
|
||||
effective_user = request.user or get_default_user()
|
||||
current_user.set(effective_user)
|
||||
current_conversation.set(conversation_id)
|
||||
|
||||
# Get agent for model
|
||||
agent = ModelRegistry.get_agent(model_id)
|
||||
# Extract user input for tracing
|
||||
user_input = _extract_user_input(request.input)
|
||||
|
||||
# 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
|
||||
# Start trace
|
||||
trace = start_trace(
|
||||
conversation_id=conversation_id,
|
||||
user=effective_user,
|
||||
request={
|
||||
"model": request.model,
|
||||
"input_preview": user_input[:200] if user_input else "",
|
||||
"full_input": request.input,
|
||||
"streaming": False,
|
||||
},
|
||||
)
|
||||
|
||||
# Track conversation history (for analytics and future vector memory)
|
||||
await _conversation_history.add_response(conversation_id, response)
|
||||
# Start service span
|
||||
service_span = start_span(
|
||||
"create_response",
|
||||
SpanType.ROUTER,
|
||||
metadata={"model": request.model, "user": effective_user},
|
||||
)
|
||||
|
||||
return response
|
||||
try:
|
||||
# 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)
|
||||
|
||||
# End trace with response info
|
||||
response_preview = _extract_response_preview(response)
|
||||
end_trace(
|
||||
response={
|
||||
"output_preview": response_preview,
|
||||
"output_count": len(response.output) if response.output else 0,
|
||||
"status": response.status,
|
||||
},
|
||||
status="completed",
|
||||
)
|
||||
|
||||
return response
|
||||
|
||||
except Exception as e:
|
||||
end_trace(status="error")
|
||||
raise
|
||||
|
||||
|
||||
async def create_response_with_steward(request: ResponseRequest) -> Response:
|
||||
@@ -454,7 +522,7 @@ async def create_response_with_steward(request: ResponseRequest) -> Response:
|
||||
1. Steward analyzes the request and recommends capabilities
|
||||
2. Phase 1: Tatlock orchestrates tool calls and expert delegations
|
||||
3. Phase 2: Tatlock synthesizes butler-toned response from results
|
||||
4. Tool usage is tracked for benchmarking
|
||||
4. Tool usage is tracked for analysis
|
||||
|
||||
Args:
|
||||
request: Response request
|
||||
@@ -473,133 +541,176 @@ async def create_response_with_steward(request: ResponseRequest) -> Response:
|
||||
# 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
|
||||
# Set context for tracing
|
||||
effective_user = request.user or get_default_user()
|
||||
current_user.set(effective_user)
|
||||
current_conversation.set(conversation_id)
|
||||
|
||||
# Conversation history is all messages except the current one
|
||||
conversation_history = request.input[:-1] if len(request.input) > 1 else []
|
||||
# Extract user input for tracing
|
||||
user_input = _extract_user_input(request.input)
|
||||
|
||||
logger.info(
|
||||
"creating_response_with_steward",
|
||||
user_message_preview=user_message[:100],
|
||||
history_length=len(conversation_history),
|
||||
# Start trace
|
||||
trace = start_trace(
|
||||
conversation_id=conversation_id,
|
||||
user=effective_user,
|
||||
request={
|
||||
"model": request.model,
|
||||
"input_preview": user_input[:200] if user_input else "",
|
||||
"full_input": request.input,
|
||||
"streaming": False,
|
||||
},
|
||||
)
|
||||
|
||||
# Steward preprocessing
|
||||
enriched = await preprocess_request(
|
||||
user_message,
|
||||
conversation_history=conversation_history,
|
||||
conversation_id=conversation_id,
|
||||
# Start service span
|
||||
service_span = start_span(
|
||||
"create_response_with_steward",
|
||||
SpanType.ROUTER,
|
||||
metadata={"model": request.model, "user": effective_user},
|
||||
)
|
||||
|
||||
# Initialize tool tracker
|
||||
tracker = ToolCallTracker(
|
||||
recommended_capabilities=enriched.recommendation.recommended_capabilities,
|
||||
conversation_id=conversation_id,
|
||||
)
|
||||
try:
|
||||
# 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
|
||||
|
||||
# Check if direct delegation is recommended
|
||||
# If Steward recommends ONLY delegation agents (biographer/librarian/housekeeper),
|
||||
# we still use two-phase but delegate directly in Phase 1
|
||||
delegation_agents = {"biographer", "librarian", "housekeeper"}
|
||||
delegation_only = all(
|
||||
cap in delegation_agents
|
||||
for cap in enriched.recommendation.recommended_capabilities
|
||||
) and enriched.recommendation.recommended_capabilities
|
||||
# Conversation history is all messages except the current one
|
||||
conversation_history = request.input[:-1] if len(request.input) > 1 else []
|
||||
|
||||
from src.agents.tatlock import TatlockAgent
|
||||
tatlock = TatlockAgent()
|
||||
|
||||
# Use enriched query (with location/timezone context) if available
|
||||
effective_query = enriched.recommendation.enriched_query or user_message
|
||||
|
||||
if delegation_only:
|
||||
# Direct delegation path - collect results then synthesize
|
||||
orchestration_results = await _direct_delegation_with_results(
|
||||
effective_query, enriched.recommendation, tracker, conversation_id
|
||||
)
|
||||
else:
|
||||
# Phase 1: Orchestrate tool calls
|
||||
orchestration_results = await tatlock.orchestrate_tool_calls(
|
||||
user_message=effective_query,
|
||||
steward_note=enriched.steward_note,
|
||||
scoped_tools=enriched.scoped_tools,
|
||||
message_history=conversation_history,
|
||||
tool_tracker=tracker,
|
||||
logger.info(
|
||||
"creating_response_with_steward",
|
||||
user_message_preview=user_message[:100],
|
||||
history_length=len(conversation_history),
|
||||
conversation_id=conversation_id,
|
||||
)
|
||||
|
||||
# Handle text-based delegation fallback if present
|
||||
if "[DELEGATE:" in orchestration_results.get("raw_output", ""):
|
||||
text_delegation_results = await _handle_text_delegation(
|
||||
orchestration_results["raw_output"], tracker, conversation_id
|
||||
# Steward preprocessing
|
||||
enriched = await preprocess_request(
|
||||
user_message,
|
||||
conversation_history=conversation_history,
|
||||
conversation_id=conversation_id,
|
||||
)
|
||||
|
||||
# Initialize tool tracker
|
||||
tracker = ToolCallTracker(
|
||||
recommended_capabilities=enriched.recommendation.recommended_capabilities,
|
||||
conversation_id=conversation_id,
|
||||
)
|
||||
|
||||
# Check if direct delegation is recommended
|
||||
# If Steward recommends ONLY delegation agents (biographer/librarian/housekeeper),
|
||||
# we still use two-phase but delegate directly in Phase 1
|
||||
delegation_agents = {"biographer", "librarian", "housekeeper"}
|
||||
delegation_only = all(
|
||||
cap in delegation_agents
|
||||
for cap in enriched.recommendation.recommended_capabilities
|
||||
) and enriched.recommendation.recommended_capabilities
|
||||
|
||||
from src.agents.tatlock import TatlockAgent
|
||||
tatlock = TatlockAgent()
|
||||
|
||||
# Use enriched query (with location/timezone context) if available
|
||||
effective_query = enriched.recommendation.enriched_query or user_message
|
||||
|
||||
if delegation_only:
|
||||
# Direct delegation path - collect results then synthesize
|
||||
orchestration_results = await _direct_delegation_with_results(
|
||||
effective_query, enriched.recommendation, tracker, conversation_id
|
||||
)
|
||||
else:
|
||||
# Phase 1: Orchestrate tool calls
|
||||
orchestration_results = await tatlock.orchestrate_tool_calls(
|
||||
user_message=effective_query,
|
||||
steward_note=enriched.steward_note,
|
||||
scoped_tools=enriched.scoped_tools,
|
||||
message_history=conversation_history,
|
||||
tool_tracker=tracker,
|
||||
)
|
||||
# Add text delegation results to expert_results
|
||||
if text_delegation_results != orchestration_results["raw_output"]:
|
||||
orchestration_results["expert_results"]["text_delegation"] = text_delegation_results
|
||||
|
||||
# Phase 2: Synthesize butler-toned response from all results
|
||||
tatlock_response = await tatlock.synthesize_from_results(
|
||||
user_message=user_message,
|
||||
orchestration_results=orchestration_results,
|
||||
message_history=conversation_history,
|
||||
)
|
||||
# Handle text-based delegation fallback if present
|
||||
if "[DELEGATE:" in orchestration_results.get("raw_output", ""):
|
||||
text_delegation_results = await _handle_text_delegation(
|
||||
orchestration_results["raw_output"], tracker, conversation_id
|
||||
)
|
||||
# Add text delegation results to expert_results
|
||||
if text_delegation_results != orchestration_results["raw_output"]:
|
||||
orchestration_results["expert_results"]["text_delegation"] = text_delegation_results
|
||||
|
||||
# Finalize tool tracking
|
||||
await tracker.finalize()
|
||||
# Phase 2: Synthesize butler-toned response from all results
|
||||
tatlock_response = await tatlock.synthesize_from_results(
|
||||
user_message=user_message,
|
||||
orchestration_results=orchestration_results,
|
||||
message_history=conversation_history,
|
||||
)
|
||||
|
||||
# Build response output items
|
||||
output_items = []
|
||||
# Finalize tool tracking
|
||||
await tracker.finalize()
|
||||
|
||||
# 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"
|
||||
))
|
||||
# Build response output items
|
||||
output_items = []
|
||||
|
||||
# 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"
|
||||
))
|
||||
# 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"
|
||||
))
|
||||
|
||||
# Calculate usage (approximate)
|
||||
usage = _calculate_usage(request.input, output_items)
|
||||
# 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"
|
||||
))
|
||||
|
||||
response = Response(
|
||||
id=f"resp_{generate_id()}",
|
||||
created_at=int(time.time()),
|
||||
model=request.model,
|
||||
status="completed",
|
||||
output=output_items,
|
||||
usage=usage
|
||||
)
|
||||
# Calculate usage (approximate)
|
||||
usage = _calculate_usage(request.input, output_items)
|
||||
|
||||
# Track conversation history
|
||||
await _conversation_history.add_response(conversation_id, response)
|
||||
response = Response(
|
||||
id=f"resp_{generate_id()}",
|
||||
created_at=int(time.time()),
|
||||
model=request.model,
|
||||
status="completed",
|
||||
output=output_items,
|
||||
usage=usage
|
||||
)
|
||||
|
||||
logger.info(
|
||||
"response_with_steward_complete",
|
||||
response_id=response.id,
|
||||
recommended_capabilities=enriched.recommendation.recommended_capabilities,
|
||||
tool_summary=tracker.get_summary(),
|
||||
)
|
||||
# Track conversation history
|
||||
await _conversation_history.add_response(conversation_id, response)
|
||||
|
||||
return response
|
||||
logger.info(
|
||||
"response_with_steward_complete",
|
||||
response_id=response.id,
|
||||
recommended_capabilities=enriched.recommendation.recommended_capabilities,
|
||||
tool_summary=tracker.get_summary(),
|
||||
)
|
||||
|
||||
# End trace with response info
|
||||
response_preview = _extract_response_preview(response)
|
||||
end_trace(
|
||||
response={
|
||||
"output_preview": response_preview,
|
||||
"output_count": len(response.output) if response.output else 0,
|
||||
"status": response.status,
|
||||
},
|
||||
status="completed",
|
||||
)
|
||||
|
||||
return response
|
||||
|
||||
except Exception as e:
|
||||
end_trace(status="error")
|
||||
raise
|
||||
|
||||
|
||||
async def create_response_stream(
|
||||
|
||||
Reference in New Issue
Block a user