- CoordinationEngine for task orchestration between agents - Routing tasks to appropriate expert agents - Sequential and parallel execution support - Result aggregation from multiple agents - Graceful error handling and degradation - Streaming delegation support - Convenience functions: delegate_to_librarian(), delegate_to_librarian_stream() 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
408 lines
12 KiB
Python
408 lines
12 KiB
Python
"""
|
|
Multi-agent coordination engine.
|
|
|
|
Orchestrates delegation from Tatlock to expert agents (Librarian, etc.)
|
|
based on Steward recommendations. Handles:
|
|
- Routing tasks to appropriate agents
|
|
- Parallel and sequential execution
|
|
- Result aggregation
|
|
- Error handling and graceful degradation
|
|
"""
|
|
import asyncio
|
|
import time
|
|
from typing import Any, AsyncGenerator, Optional
|
|
|
|
from src.agents.librarian import run_librarian, run_librarian_stream
|
|
from src.agents.protocol import (
|
|
AgentError,
|
|
AgentRequest,
|
|
AgentResponse,
|
|
AgentTimeoutError,
|
|
AgentUnavailableError,
|
|
CoordinationResult,
|
|
DelegationIntent,
|
|
DelegationReason,
|
|
ToolCallRecord,
|
|
)
|
|
from src.core.household_registry import get_household_registry
|
|
from src.core.logging_config import get_logger
|
|
|
|
logger = get_logger(__name__)
|
|
|
|
|
|
# Agent execution functions registry
|
|
AGENT_EXECUTORS: dict[str, Any] = {
|
|
"librarian": run_librarian,
|
|
}
|
|
|
|
AGENT_STREAM_EXECUTORS: dict[str, Any] = {
|
|
"librarian": run_librarian_stream,
|
|
}
|
|
|
|
|
|
class CoordinationEngine:
|
|
"""
|
|
Coordinates multi-agent task execution.
|
|
|
|
Routes tasks from Tatlock to appropriate expert agents,
|
|
handles execution, and aggregates results.
|
|
"""
|
|
|
|
def __init__(self):
|
|
"""Initialize the coordination engine."""
|
|
self.registry = get_household_registry()
|
|
logger.info("coordination_engine_initialized")
|
|
|
|
def get_available_agents(self) -> list[str]:
|
|
"""
|
|
Get list of available expert agents.
|
|
|
|
Returns:
|
|
List of agent names that can accept delegations
|
|
"""
|
|
available = []
|
|
for name in self.registry.list_members():
|
|
member = self.registry.get_member(name)
|
|
if member and member.agent is not None:
|
|
available.append(name)
|
|
return available
|
|
|
|
def can_delegate_to(self, agent_name: str) -> bool:
|
|
"""
|
|
Check if delegation to an agent is possible.
|
|
|
|
Args:
|
|
agent_name: Name of the target agent
|
|
|
|
Returns:
|
|
True if agent is available and can accept tasks
|
|
"""
|
|
if agent_name not in AGENT_EXECUTORS:
|
|
return False
|
|
|
|
member = self.registry.get_member(agent_name)
|
|
return member is not None and member.agent is not None
|
|
|
|
async def execute_delegation(
|
|
self,
|
|
intent: DelegationIntent,
|
|
context: str = "",
|
|
message_history: Optional[list[Any]] = None,
|
|
) -> AgentResponse:
|
|
"""
|
|
Execute a single delegation to an expert agent.
|
|
|
|
Args:
|
|
intent: The delegation intent with task details
|
|
context: Additional context for the agent
|
|
message_history: Optional conversation history
|
|
|
|
Returns:
|
|
AgentResponse with results
|
|
|
|
Raises:
|
|
AgentUnavailableError: If agent is not available
|
|
AgentTimeoutError: If execution times out
|
|
AgentError: For other execution errors
|
|
"""
|
|
start_time = time.time()
|
|
agent_name = intent.target_agent
|
|
|
|
logger.info(
|
|
"delegation_started",
|
|
agent=agent_name,
|
|
task=intent.task[:100],
|
|
reason=intent.reason.value,
|
|
)
|
|
|
|
# Check if agent is available
|
|
if not self.can_delegate_to(agent_name):
|
|
raise AgentUnavailableError(
|
|
f"Agent '{agent_name}' is not available for delegation",
|
|
agent_name=agent_name,
|
|
)
|
|
|
|
# Get the executor
|
|
executor = AGENT_EXECUTORS.get(agent_name)
|
|
if not executor:
|
|
raise AgentUnavailableError(
|
|
f"No executor found for agent '{agent_name}'",
|
|
agent_name=agent_name,
|
|
)
|
|
|
|
try:
|
|
# Build the request
|
|
request = AgentRequest(
|
|
task=intent.task,
|
|
context=context,
|
|
delegation_reason=intent.reason,
|
|
)
|
|
|
|
# Execute with timeout
|
|
timeout = request.timeout_seconds or 60
|
|
|
|
result = await asyncio.wait_for(
|
|
executor(
|
|
task=request.task,
|
|
context=request.context,
|
|
message_history=message_history,
|
|
),
|
|
timeout=timeout,
|
|
)
|
|
|
|
duration_ms = int((time.time() - start_time) * 1000)
|
|
|
|
logger.info(
|
|
"delegation_completed",
|
|
agent=agent_name,
|
|
duration_ms=duration_ms,
|
|
output_length=len(result),
|
|
)
|
|
|
|
return AgentResponse(
|
|
success=True,
|
|
result=result,
|
|
reasoning=f"Delegated to {agent_name}: {intent.expected_outcome}",
|
|
duration_ms=duration_ms,
|
|
)
|
|
|
|
except asyncio.TimeoutError:
|
|
duration_ms = int((time.time() - start_time) * 1000)
|
|
logger.error(
|
|
"delegation_timeout",
|
|
agent=agent_name,
|
|
duration_ms=duration_ms,
|
|
)
|
|
raise AgentTimeoutError(
|
|
f"Agent '{agent_name}' timed out after {duration_ms}ms",
|
|
agent_name=agent_name,
|
|
)
|
|
|
|
except Exception as e:
|
|
duration_ms = int((time.time() - start_time) * 1000)
|
|
logger.error(
|
|
"delegation_error",
|
|
agent=agent_name,
|
|
error=str(e),
|
|
duration_ms=duration_ms,
|
|
exc_info=True,
|
|
)
|
|
return AgentResponse(
|
|
success=False,
|
|
result="",
|
|
error_message=str(e),
|
|
duration_ms=duration_ms,
|
|
)
|
|
|
|
async def execute_delegation_stream(
|
|
self,
|
|
intent: DelegationIntent,
|
|
context: str = "",
|
|
message_history: Optional[list[Any]] = None,
|
|
) -> AsyncGenerator[str, None]:
|
|
"""
|
|
Execute a delegation with streaming output.
|
|
|
|
Args:
|
|
intent: The delegation intent with task details
|
|
context: Additional context for the agent
|
|
message_history: Optional conversation history
|
|
|
|
Yields:
|
|
Text deltas from the agent
|
|
|
|
Raises:
|
|
AgentUnavailableError: If agent is not available
|
|
"""
|
|
agent_name = intent.target_agent
|
|
|
|
logger.info(
|
|
"delegation_stream_started",
|
|
agent=agent_name,
|
|
task=intent.task[:100],
|
|
)
|
|
|
|
# Check if agent is available
|
|
if agent_name not in AGENT_STREAM_EXECUTORS:
|
|
raise AgentUnavailableError(
|
|
f"Agent '{agent_name}' does not support streaming",
|
|
agent_name=agent_name,
|
|
)
|
|
|
|
executor = AGENT_STREAM_EXECUTORS[agent_name]
|
|
|
|
try:
|
|
async for delta in executor(
|
|
task=intent.task,
|
|
context=context,
|
|
message_history=message_history,
|
|
):
|
|
yield delta
|
|
|
|
logger.info("delegation_stream_completed", agent=agent_name)
|
|
|
|
except Exception as e:
|
|
logger.error(
|
|
"delegation_stream_error",
|
|
agent=agent_name,
|
|
error=str(e),
|
|
exc_info=True,
|
|
)
|
|
yield f"\n\n[Error from {agent_name}: {str(e)}]"
|
|
|
|
async def coordinate(
|
|
self,
|
|
intents: list[DelegationIntent],
|
|
context: str = "",
|
|
message_history: Optional[list[Any]] = None,
|
|
) -> CoordinationResult:
|
|
"""
|
|
Coordinate execution of multiple delegations.
|
|
|
|
Handles parallel execution for independent tasks and
|
|
sequential execution for dependent tasks.
|
|
|
|
Args:
|
|
intents: List of delegation intents to execute
|
|
context: Shared context for all agents
|
|
message_history: Optional conversation history
|
|
|
|
Returns:
|
|
CoordinationResult with aggregated results
|
|
"""
|
|
start_time = time.time()
|
|
agent_responses: dict[str, AgentResponse] = {}
|
|
agents_consulted: list[str] = []
|
|
|
|
logger.info(
|
|
"coordination_started",
|
|
intent_count=len(intents),
|
|
agents=[i.target_agent for i in intents],
|
|
)
|
|
|
|
# Sort by priority
|
|
sorted_intents = sorted(intents, key=lambda x: x.priority)
|
|
|
|
# Group by dependencies (simple version: sequential for now)
|
|
# TODO: Implement parallel execution for independent tasks
|
|
for intent in sorted_intents:
|
|
try:
|
|
response = await self.execute_delegation(
|
|
intent=intent,
|
|
context=context,
|
|
message_history=message_history,
|
|
)
|
|
agent_responses[intent.target_agent] = response
|
|
if response.success:
|
|
agents_consulted.append(intent.target_agent)
|
|
|
|
except AgentError as e:
|
|
agent_responses[intent.target_agent] = AgentResponse(
|
|
success=False,
|
|
result="",
|
|
error_message=str(e),
|
|
)
|
|
|
|
# Aggregate results
|
|
successful_results = [
|
|
r.result for r in agent_responses.values() if r.success and r.result
|
|
]
|
|
|
|
final_response = "\n\n---\n\n".join(successful_results) if successful_results else ""
|
|
|
|
total_duration = int((time.time() - start_time) * 1000)
|
|
|
|
logger.info(
|
|
"coordination_completed",
|
|
total_duration_ms=total_duration,
|
|
agents_consulted=agents_consulted,
|
|
success_count=len(successful_results),
|
|
)
|
|
|
|
return CoordinationResult(
|
|
final_response=final_response,
|
|
agent_responses=agent_responses,
|
|
delegation_intents=intents,
|
|
total_duration_ms=total_duration,
|
|
agents_consulted=agents_consulted,
|
|
)
|
|
|
|
|
|
# Global coordination engine instance
|
|
_coordination_engine: Optional[CoordinationEngine] = None
|
|
|
|
|
|
def get_coordination_engine() -> CoordinationEngine:
|
|
"""Get the global coordination engine instance."""
|
|
global _coordination_engine
|
|
if _coordination_engine is None:
|
|
_coordination_engine = CoordinationEngine()
|
|
return _coordination_engine
|
|
|
|
|
|
async def delegate_to_librarian(
|
|
task: str,
|
|
context: str = "",
|
|
reason: DelegationReason = DelegationReason.DOMAIN_EXPERTISE,
|
|
message_history: Optional[list[Any]] = None,
|
|
) -> AgentResponse:
|
|
"""
|
|
Convenience function to delegate a task to The Librarian.
|
|
|
|
Args:
|
|
task: Research task description
|
|
context: Additional context
|
|
reason: Why delegating to Librarian
|
|
message_history: Optional conversation history
|
|
|
|
Returns:
|
|
AgentResponse with research results
|
|
"""
|
|
engine = get_coordination_engine()
|
|
|
|
intent = DelegationIntent(
|
|
target_agent="librarian",
|
|
task=task,
|
|
reason=reason,
|
|
expected_outcome="Research findings and relevant information",
|
|
)
|
|
|
|
return await engine.execute_delegation(
|
|
intent=intent,
|
|
context=context,
|
|
message_history=message_history,
|
|
)
|
|
|
|
|
|
async def delegate_to_librarian_stream(
|
|
task: str,
|
|
context: str = "",
|
|
message_history: Optional[list[Any]] = None,
|
|
) -> AsyncGenerator[str, None]:
|
|
"""
|
|
Convenience function to delegate to Librarian with streaming.
|
|
|
|
Args:
|
|
task: Research task description
|
|
context: Additional context
|
|
message_history: Optional conversation history
|
|
|
|
Yields:
|
|
Text deltas from The Librarian
|
|
"""
|
|
engine = get_coordination_engine()
|
|
|
|
intent = DelegationIntent(
|
|
target_agent="librarian",
|
|
task=task,
|
|
reason=DelegationReason.DOMAIN_EXPERTISE,
|
|
expected_outcome="Research findings",
|
|
)
|
|
|
|
async for delta in engine.execute_delegation_stream(
|
|
intent=intent,
|
|
context=context,
|
|
message_history=message_history,
|
|
):
|
|
yield delta
|