Twelve gain `-> None`, each confirmed by AST to contain no returning `return` and no `yield` rather than by reading the name and assuming. The three context-manager exits gain the canonical type[BaseException]/BaseException/TracebackType argument triple. Both files taking TracebackType needed the import, and inserting it before the first import broke ruff's I001 — lint was exit 0 at the baseline commit, verified by stashing this work and re-running, so that breakage was mine. Fixed with `ruff check --fix` on the two files, which placed the import in sorted position. 86 errors -> 75; no-untyped-def 29 -> 14. Suite: 658 passed. The baseline was 657 passed with one failure in test_tatlock_tool_call_logging_calculator, which asserts on the content of a live model's reply. It passing here is nondeterminism, NOT evidence this commit fixed anything, and it may fail again on the next run. Co-Authored-By: Claude <noreply@anthropic.com>
580 lines
20 KiB
Python
580 lines
20 KiB
Python
"""
|
|
Delegation infrastructure for expert agent calls.
|
|
|
|
Provides delegation wrappers that Tatlock uses to call expert agents.
|
|
Each wrapper encapsulates the complexity of calling an expert and
|
|
returns a structured result for synthesis.
|
|
|
|
This implements the agent-as-tool pattern recommended by PydanticAI:
|
|
agents call other agents via tool wrappers, keeping each agent focused.
|
|
"""
|
|
|
|
import asyncio
|
|
from dataclasses import dataclass, field
|
|
from enum import Enum
|
|
|
|
from src.core.config import config
|
|
from src.core.logging_config import get_logger
|
|
from src.core.tracing import SpanType, trace_span
|
|
|
|
logger = get_logger(__name__)
|
|
|
|
|
|
# =============================================================================
|
|
# Action Types for Think Slug Selection
|
|
# =============================================================================
|
|
|
|
|
|
class ActionType(Enum):
|
|
"""
|
|
Categories of actions for selecting appropriate think messages.
|
|
|
|
Each expert has different action types that warrant different
|
|
butler-perspective messages to the user.
|
|
"""
|
|
|
|
RETRIEVE = "retrieve" # Looking up existing information
|
|
RESEARCH = "research" # Conducting new research (web search, etc.)
|
|
CREATE = "create" # Creating new content (pages, notes)
|
|
CONTROL = "control" # Controlling devices/automations
|
|
RECORD = "record" # Recording memories/notes
|
|
|
|
|
|
# =============================================================================
|
|
# Household Think Messages (Butler's Perspective)
|
|
# =============================================================================
|
|
|
|
HOUSEHOLD_THINK_MESSAGES: dict[str, dict[ActionType, dict[str, str]]] = {
|
|
# Note: No <think> wrappers needed - these go to reasoning_content field
|
|
"librarian": {
|
|
ActionType.RETRIEVE: {
|
|
"start": "Allow me to consult the archives, sir.",
|
|
"success": "The Librarian has compiled the relevant findings.",
|
|
"error": "I'm afraid the archives proved difficult to access.",
|
|
},
|
|
ActionType.RESEARCH: {
|
|
"start": "I've dispatched the Librarian to conduct some fresh research.",
|
|
"success": "The Librarian has returned with findings, sir.",
|
|
"error": "The research proved inconclusive, I'm afraid.",
|
|
},
|
|
ActionType.CREATE: {
|
|
"start": "I'm having the Librarian prepare a new entry.",
|
|
"success": "The new material has been properly catalogued, sir.",
|
|
"error": "I'm afraid there was difficulty filing the entry.",
|
|
},
|
|
},
|
|
"biographer": {
|
|
ActionType.RETRIEVE: {
|
|
"start": "Let me consult the household records.",
|
|
"success": "The Biographer has located the relevant information, sir.",
|
|
"error": "I'm unable to locate those particular records.",
|
|
},
|
|
ActionType.RECORD: {
|
|
"start": "I've asked the Biographer to take note of this, sir.",
|
|
"success": "The household records have been updated accordingly.",
|
|
"error": "I'm afraid there was difficulty recording the entry.",
|
|
},
|
|
},
|
|
"housekeeper": {
|
|
ActionType.RETRIEVE: {
|
|
"start": "Allow me to inquire with the household staff.",
|
|
"success": "The staff reports the current status, sir.",
|
|
"error": "The household staff is momentarily unavailable, I'm afraid.",
|
|
},
|
|
ActionType.CONTROL: {
|
|
"start": "I'm instructing the household staff now, sir.",
|
|
"success": "The household has been configured as requested.",
|
|
"error": "I'm afraid the staff reports an issue with that request.",
|
|
},
|
|
},
|
|
}
|
|
|
|
|
|
def _detect_action_type(expert: str, task: str) -> ActionType:
|
|
"""
|
|
Detect action type from expert name and task description.
|
|
|
|
Used to select appropriate butler-perspective think messages.
|
|
|
|
Args:
|
|
expert: Name of the expert (librarian, biographer, housekeeper)
|
|
task: Task description
|
|
|
|
Returns:
|
|
ActionType: Detected action type for message selection
|
|
"""
|
|
task_lower = task.lower()
|
|
|
|
if expert == "librarian":
|
|
# Web search, URL reading = RESEARCH (fresh external data)
|
|
if any(w in task_lower for w in ["search", "find", "look up", "research"]):
|
|
if any(w in task_lower for w in ["web", "online", "internet"]):
|
|
return ActionType.RESEARCH
|
|
return ActionType.RETRIEVE
|
|
if any(w in task_lower for w in ["read", "fetch", "url", "http"]):
|
|
return ActionType.RESEARCH # Reading URLs is research
|
|
if any(w in task_lower for w in ["create", "write", "add", "make", "new"]):
|
|
return ActionType.CREATE
|
|
return ActionType.RETRIEVE
|
|
|
|
elif expert == "biographer":
|
|
if any(w in task_lower for w in ["remember", "note", "record", "save", "store"]):
|
|
return ActionType.RECORD
|
|
return ActionType.RETRIEVE
|
|
|
|
elif expert == "housekeeper":
|
|
if any(w in task_lower for w in ["turn", "set", "activate", "enable", "disable", "toggle"]):
|
|
return ActionType.CONTROL
|
|
return ActionType.RETRIEVE
|
|
|
|
return ActionType.RETRIEVE
|
|
|
|
|
|
def build_delegation_context(
|
|
conversation_history: list[dict] | None,
|
|
max_turns: int = 6,
|
|
max_chars_per_turn: int = 500,
|
|
) -> str:
|
|
"""
|
|
Format the most recent conversation turns as delegation context.
|
|
|
|
Experts accept a context string but the live paths never passed the
|
|
in-scope conversation history; this trims it to the last few turns
|
|
so follow-up questions ("and what about X?") keep their referent.
|
|
|
|
Args:
|
|
conversation_history: Prior messages as {"role", "content"} dicts
|
|
max_turns: How many trailing turns to include
|
|
max_chars_per_turn: Truncation limit per turn
|
|
|
|
Returns:
|
|
str: Newline-joined "role: content" lines ("" when no history)
|
|
"""
|
|
if not conversation_history:
|
|
return ""
|
|
|
|
lines = []
|
|
for msg in conversation_history[-max_turns:]:
|
|
if not isinstance(msg, dict):
|
|
continue
|
|
role = msg.get("role", "user")
|
|
content = msg.get("content", "")
|
|
if isinstance(content, list):
|
|
# Tolerate structured content parts
|
|
content = " ".join(
|
|
part.get("text", "") if isinstance(part, dict) else str(part) for part in content
|
|
)
|
|
content = str(content).strip()
|
|
if content:
|
|
lines.append(f"{role}: {content[:max_chars_per_turn]}")
|
|
|
|
if not lines:
|
|
return ""
|
|
return "Recent conversation:\n" + "\n".join(lines)
|
|
|
|
|
|
def get_think_message(expert: str, task: str, phase: str) -> str:
|
|
"""
|
|
Get the appropriate think message for an expert delegation.
|
|
|
|
Args:
|
|
expert: Name of the expert
|
|
task: Task description (used to detect action type)
|
|
phase: One of "start", "success", "error"
|
|
|
|
Returns:
|
|
str: Butler-perspective think message
|
|
"""
|
|
action_type = _detect_action_type(expert, task)
|
|
expert_messages = HOUSEHOLD_THINK_MESSAGES.get(expert, {})
|
|
action_messages = expert_messages.get(action_type, expert_messages.get(ActionType.RETRIEVE, {}))
|
|
return action_messages.get(phase, f"Consulting {expert}...")
|
|
|
|
|
|
@dataclass
|
|
class DelegationTask:
|
|
"""
|
|
A task to be delegated to an expert agent.
|
|
|
|
Represents a unit of work that Tatlock delegates to a specialist.
|
|
Used for tracking and orchestration of multi-expert workflows.
|
|
|
|
Attributes:
|
|
expert_name: Name of the expert agent (e.g., "librarian", "memory")
|
|
task: Clear description of what needs to be done
|
|
context: Additional context from the conversation
|
|
action: Specific action verb (create, search, update, etc.)
|
|
priority: Execution priority (lower = higher priority)
|
|
depends_on: List of task IDs this task depends on
|
|
result: Result from expert after execution
|
|
"""
|
|
|
|
expert_name: str
|
|
task: str
|
|
context: str = ""
|
|
action: str = ""
|
|
priority: int = 0
|
|
depends_on: list[str] = field(default_factory=list)
|
|
result: str | None = None
|
|
task_id: str = ""
|
|
|
|
def __post_init__(self) -> None:
|
|
"""Generate task ID if not provided."""
|
|
if not self.task_id:
|
|
import uuid
|
|
|
|
self.task_id = f"{self.expert_name}_{uuid.uuid4().hex[:8]}"
|
|
|
|
|
|
@dataclass
|
|
class DelegationResult:
|
|
"""
|
|
Result from an expert agent delegation.
|
|
|
|
Attributes:
|
|
expert_name: Which expert handled the task
|
|
task: Original task description
|
|
success: Whether the delegation succeeded
|
|
output: Expert's response/findings. On failure this holds a
|
|
curated, user-safe butler sentence (never exception detail)
|
|
error: Short user-safe error label if failed. Exception detail
|
|
stays in the logs only
|
|
"""
|
|
|
|
expert_name: str
|
|
task: str
|
|
success: bool
|
|
output: str
|
|
error: str | None = None
|
|
|
|
|
|
async def delegate_to_librarian(
|
|
task: str,
|
|
context: str = "",
|
|
) -> DelegationResult:
|
|
"""
|
|
Delegate a research or wiki task to The Librarian.
|
|
|
|
The Librarian handles:
|
|
- Wiki creation (smart_create_wiki_page for topic-based)
|
|
- Wiki updates (update_wiki_page for modifications)
|
|
- Research queries (hybrid_search for comprehensive search)
|
|
- Knowledge graph exploration
|
|
- Document lookups and semantic search
|
|
|
|
This wrapper uses run() not run_stream() to avoid Ollama's
|
|
streaming + tool call bug (PydanticAI issues #1292, #2256).
|
|
|
|
Args:
|
|
task: Clear description of what needs to be done.
|
|
Include the action verb (create, search, update, etc.)
|
|
Example: "Create a wiki page about CI/CD pipelines"
|
|
Example: "Search for information about Docker networking"
|
|
context: Additional context from the user's request or
|
|
conversation history
|
|
|
|
Returns:
|
|
DelegationResult with the Librarian's findings
|
|
|
|
Example:
|
|
>>> result = await delegate_to_librarian(
|
|
... task="Create a wiki page about Kubernetes deployments",
|
|
... context="User is setting up a homelab cluster",
|
|
... )
|
|
>>> if result.success:
|
|
... print(result.output)
|
|
"""
|
|
from src.agents.librarian.agent import run_librarian
|
|
|
|
logger.info(
|
|
"delegation_to_librarian_started",
|
|
task=task[:100],
|
|
has_context=bool(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.
|
|
# One timeout budget for the whole delegation - covers both
|
|
# live paths (steward direct delegation and streaming), which
|
|
# previously had no cap at all (SDK default ~600s per LLM call).
|
|
output = await asyncio.wait_for(
|
|
run_librarian(task=task, context=context),
|
|
timeout=config.LIBRARIAN_TIMEOUT,
|
|
)
|
|
|
|
logger.info(
|
|
"delegation_to_librarian_completed",
|
|
task=task[:50],
|
|
output_length=len(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]
|
|
|
|
return DelegationResult(
|
|
expert_name="librarian",
|
|
task=task,
|
|
success=True,
|
|
output=output,
|
|
)
|
|
|
|
except TimeoutError:
|
|
logger.error(
|
|
"delegation_to_librarian_timeout",
|
|
task=task[:50],
|
|
timeout_seconds=config.LIBRARIAN_TIMEOUT,
|
|
)
|
|
|
|
if span:
|
|
span.metadata["success"] = False
|
|
span.details["error"] = f"timed out after {config.LIBRARIAN_TIMEOUT}s"
|
|
|
|
return DelegationResult(
|
|
expert_name="librarian",
|
|
task=task,
|
|
success=False,
|
|
output=(
|
|
"I'm afraid the research took longer than expected "
|
|
"and had to be abandoned, sir."
|
|
),
|
|
error="The Librarian did not respond within the time budget.",
|
|
)
|
|
|
|
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)
|
|
|
|
# Exception detail stays in the logs; the user-facing output
|
|
# is a curated butler sentence so internals never leak into
|
|
# synthesis.
|
|
return DelegationResult(
|
|
expert_name="librarian",
|
|
task=task,
|
|
success=False,
|
|
output=get_think_message("librarian", task, "error"),
|
|
error="The Librarian was unable to complete the task.",
|
|
)
|
|
|
|
|
|
async def delegate_to_biographer(
|
|
task: str,
|
|
context: str = "",
|
|
) -> DelegationResult:
|
|
"""
|
|
Delegate a memory task to The Biographer.
|
|
|
|
The Biographer handles:
|
|
- Semantic recall ("What car do I drive?", "What's my job?")
|
|
- Recording new facts from conversation
|
|
- Profile updates (name, location, timezone)
|
|
- Preference updates (units, theme)
|
|
- Memory management (forget, list)
|
|
|
|
For direct key-based lookups (get location, get timezone), use
|
|
memory_service directly - it's faster and doesn't require LLM.
|
|
|
|
Args:
|
|
task: Clear description of what needs to be done.
|
|
Include the action verb (recall, remember, forget, etc.)
|
|
Example: "What car do I drive?"
|
|
Example: "Remember that I work at Acme Corp"
|
|
context: Additional context from the user's request or
|
|
conversation history
|
|
|
|
Returns:
|
|
DelegationResult with The Biographer's response
|
|
|
|
Example:
|
|
>>> result = await delegate_to_biographer(
|
|
... task="What do you know about my preferences?",
|
|
... context="User is asking about stored information",
|
|
... )
|
|
>>> if result.success:
|
|
... print(result.output)
|
|
"""
|
|
from src.agents.biographer.agent import run_biographer
|
|
|
|
logger.info(
|
|
"delegation_to_biographer_started",
|
|
task=task[:100],
|
|
has_context=bool(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),
|
|
)
|
|
|
|
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]
|
|
|
|
return DelegationResult(
|
|
expert_name="biographer",
|
|
task=task,
|
|
success=True,
|
|
output=output,
|
|
)
|
|
|
|
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)
|
|
|
|
# Exception detail stays in the logs only.
|
|
return DelegationResult(
|
|
expert_name="biographer",
|
|
task=task,
|
|
success=False,
|
|
output=get_think_message("biographer", task, "error"),
|
|
error="The Biographer was unable to complete the task.",
|
|
)
|
|
|
|
|
|
async def delegate_to_housekeeper(
|
|
task: str,
|
|
context: str = "",
|
|
) -> DelegationResult:
|
|
"""
|
|
Delegate a home automation task to The Housekeeper.
|
|
|
|
The Housekeeper handles:
|
|
- Device control (turn on/off, toggle, brightness, color)
|
|
- Scene activation (movie night, good morning, etc.)
|
|
- Script execution (automation sequences)
|
|
- Automation management (enable/disable rules)
|
|
- Device discovery (list devices by area/type)
|
|
- State queries (get current state, history)
|
|
|
|
Args:
|
|
task: Clear description of what needs to be done.
|
|
Include the action verb (turn on, activate, list, etc.)
|
|
Example: "Turn on the living room lights"
|
|
Example: "Activate the movie night scene"
|
|
Example: "What devices are in the bedroom?"
|
|
context: Additional context from the user's request or
|
|
conversation history
|
|
|
|
Returns:
|
|
DelegationResult with The Housekeeper's response
|
|
|
|
Example:
|
|
>>> result = await delegate_to_housekeeper(
|
|
... task="Turn on the bedroom lights at 50% brightness",
|
|
... context="User is getting ready for bed",
|
|
... )
|
|
>>> if result.success:
|
|
... print(result.output)
|
|
"""
|
|
from src.agents.housekeeper.agent import run_housekeeper
|
|
|
|
logger.info(
|
|
"delegation_to_housekeeper_started",
|
|
task=task[:100],
|
|
has_context=bool(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),
|
|
)
|
|
|
|
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]
|
|
|
|
return DelegationResult(
|
|
expert_name="housekeeper",
|
|
task=task,
|
|
success=True,
|
|
output=output,
|
|
)
|
|
|
|
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)
|
|
|
|
# Exception detail stays in the logs only.
|
|
return DelegationResult(
|
|
expert_name="housekeeper",
|
|
task=task,
|
|
success=False,
|
|
output=get_think_message("housekeeper", task, "error"),
|
|
error="The Housekeeper was unable to complete the task.",
|
|
)
|
|
|
|
|
|
# Future expert delegation wrappers will be added here:
|
|
# - delegate_to_developer(task, context) -> DelegationResult
|
|
# - delegate_to_secretary(task, context) -> DelegationResult
|