From 87f2926db257a31cf62a36e0a49db7d71efb2a97 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Mon, 22 Dec 2025 10:26:37 +0100 Subject: [PATCH] feat: integrate tracing throughout request pipeline MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- src/agents/delegation.py | 229 ++++++++++++-------- src/agents/steward/service.py | 25 +-- src/agents/tatlock.py | 58 ++++- src/core/preprocessing.py | 33 ++- src/core/tool_tracking.py | 47 +--- src/main.py | 10 +- src/responses/router.py | 78 +------ src/responses/service.py | 391 ++++++++++++++++++++++------------ 8 files changed, 502 insertions(+), 369 deletions(-) diff --git a/src/agents/delegation.py b/src/agents/delegation.py index 7be59cb..08d977d 100644 --- a/src/agents/delegation.py +++ b/src/agents/delegation.py @@ -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), + ) # ============================================================================= diff --git a/src/agents/steward/service.py b/src/agents/steward/service.py index b152b92..87601d0 100644 --- a/src/agents/steward/service.py +++ b/src/agents/steward/service.py @@ -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: diff --git a/src/agents/tatlock.py b/src/agents/tatlock.py index 4bcdbe5..9f4c8d7 100644 --- a/src/agents/tatlock.py +++ b/src/agents/tatlock.py @@ -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: diff --git a/src/core/preprocessing.py b/src/core/preprocessing.py index b91253a..8a37e17 100644 --- a/src/core/preprocessing.py +++ b/src/core/preprocessing.py @@ -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) diff --git a/src/core/tool_tracking.py b/src/core/tool_tracking.py index beb42c2..b42410d 100644 --- a/src/core/tool_tracking.py +++ b/src/core/tool_tracking.py @@ -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( diff --git a/src/main.py b/src/main.py index 0f7a337..d3028af 100644 --- a/src/main.py +++ b/src/main.py @@ -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 diff --git a/src/responses/router.py b/src/responses/router.py index 2d23f36..54731b1 100644 --- a/src/responses/router.py +++ b/src/responses/router.py @@ -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) diff --git a/src/responses/service.py b/src/responses/service.py index b73878f..9e96874 100644 --- a/src/responses/service.py +++ b/src/responses/service.py @@ -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(