""" tool_agent — Agente Especializado para Ferramentas (Tools). Implementa: - 2.1. Memória individualizada de instruções e dados por ferramenta - 2.2. Suporte Model Context Protocol (MCP) - 2.3. Resiliência: falha isolada, restart, circuit breaker - 2.3.1. Arquitetura distribuída otimizada (análise matemática) Arquitetura recomendada (análise): - Hybrid: 1 Agente Coordenador + N Agentes Especialistas (1 por ferramenta) - Vantagens sobre Multi-Agent (1 agente p/ tudo) e Single-Agent (1 p/ todas): * Isolamento de falhas (1 ferramenta falha → só aquele agente reinicia) * Especialização (cada agente tem memória otimizada para sua ferramenta) * Paralelismo (agentes independentes rodam em paralelo) * Escalabilidade (adicionar ferramenta = adicionar agente) Matemática da arquitetura: - Latência esperada: T = max_i(T_i) + T_coord (paralelo + coordenação) - Vs Single-Agent: T_single = Σ T_i (sequencial) - Speedup: S = Σ T_i / (max T_i + T_coord) ≈ N / (1 + T_coord/max T_i) - Quando T_coord << max T_i: S → N (speedup linear) """ from __future__ import annotations import asyncio import hashlib import json import os import time import threading from collections import defaultdict, deque from concurrent.futures import Future, ThreadPoolExecutor, as_completed from dataclasses import dataclass, field from enum import Enum from typing import Any, Callable, Dict, List, Optional, Set, Tuple, Union import torch import torch.nn as nn # ============================================================================ # Exceções # ============================================================================ class ToolError(RuntimeError): """Erro genérico de ferramenta.""" class ToolTimeoutError(ToolError): """Timeout ao executar ferramenta.""" class ToolCircuitBreakerOpen(ToolError): """Circuit breaker aberto — ferramenta em estado de falha.""" class MCPConnectionError(ToolError): """Erro de conexão MCP.""" class ToolMemoryError(ToolError): """Erro de memória da ferramenta.""" # ============================================================================ # 2.1. ToolMemory — memória individualizada por ferramenta # ============================================================================ class ToolMemory: """Memória individualizada de instruções e dados por ferramenta. Cada ferramenta tem sua própria ToolMemory, contendo: - instructions: instruções de uso (templates, exemplos) - data_cache: cache de resultados anteriores (key → value) - error_history: histórico de erros (para diagnóstico) - usage_stats: estatísticas de uso (n_calls, n_success, n_failures) Persistência: salva/carrega de disco (JSON). """ def __init__( self, tool_name: str, max_cache_size: int = 1000, max_instructions: int = 100, max_error_history: int = 50, persist_path: Optional[str] = None, ): self.tool_name = tool_name self.max_cache_size = max_cache_size self.max_instructions = max_instructions self.max_error_history = max_error_history self.persist_path = persist_path self._lock = threading.Lock() # Estado. self.instructions: List[Dict[str, Any]] = [] self.data_cache: Dict[str, Any] = {} # key → cached result self.error_history: deque = deque(maxlen=max_error_history) self.usage_stats = { "n_calls": 0, "n_success": 0, "n_failures": 0, "total_time_s": 0.0, "avg_time_s": 0.0, } # Carrega do disco se disponível. if persist_path and os.path.exists(persist_path): self.load() def add_instruction( self, instruction: str, context: Optional[Dict[str, Any]] = None, priority: int = 0, ) -> None: """Adiciona instrução de uso.""" with self._lock: entry = { "instruction": instruction, "context": context or {}, "priority": priority, "timestamp": time.time(), } self.instructions.append(entry) # Ordena por prioridade (maior = mais importante). self.instructions.sort(key=lambda x: -x["priority"]) # Limita tamanho. if len(self.instructions) > self.max_instructions: self.instructions = self.instructions[:self.max_instructions] self._maybe_persist() def get_instructions( self, context_filter: Optional[Dict[str, Any]] = None, top_k: int = 10, ) -> List[Dict[str, Any]]: """Recupera instruções relevantes.""" with self._lock: if context_filter is None: return self.instructions[:top_k] # Filtra por contexto. filtered = [ inst for inst in self.instructions if all(inst["context"].get(k) == v for k, v in context_filter.items()) ] return filtered[:top_k] def cache_result(self, key: str, value: Any, ttl: float = 3600) -> None: """Cacheia resultado com TTL (time-to-live).""" with self._lock: if len(self.data_cache) >= self.max_cache_size: # Remove entrada mais antiga (FIFO). oldest_key = next(iter(self.data_cache)) del self.data_cache[oldest_key] self.data_cache[key] = { "value": value, "timestamp": time.time(), "ttl": ttl, } self._maybe_persist() def get_cached(self, key: str) -> Optional[Any]: """Recupera resultado do cache se não expirou.""" with self._lock: if key not in self.data_cache: return None entry = self.data_cache[key] if time.time() - entry["timestamp"] > entry["ttl"]: del self.data_cache[key] return None return entry["value"] def record_error(self, error: str, context: Optional[Dict] = None) -> None: """Registra erro no histórico.""" with self._lock: self.error_history.append({ "error": error, "context": context or {}, "timestamp": time.time(), }) def record_call(self, success: bool, duration_s: float) -> None: """Registra estatísticas de uso.""" with self._lock: self.usage_stats["n_calls"] += 1 if success: self.usage_stats["n_success"] += 1 else: self.usage_stats["n_failures"] += 1 self.usage_stats["total_time_s"] += duration_s n = self.usage_stats["n_calls"] self.usage_stats["avg_time_s"] = self.usage_stats["total_time_s"] / n def get_stats(self) -> Dict[str, Any]: """Retorna estatísticas de uso.""" with self._lock: return { **self.usage_stats, "cache_size": len(self.data_cache), "instructions_count": len(self.instructions), "error_count": len(self.error_history), "success_rate": ( self.usage_stats["n_success"] / max(self.usage_stats["n_calls"], 1) ), } def clear_cache(self) -> None: """Limpa cache.""" with self._lock: self.data_cache.clear() def save(self) -> None: """Salva estado em disco.""" if not self.persist_path: return os.makedirs(os.path.dirname(self.persist_path), exist_ok=True) data = { "tool_name": self.tool_name, "instructions": self.instructions, "usage_stats": self.usage_stats, "error_history": list(self.error_history), } with open(self.persist_path, "w") as f: json.dump(data, f, indent=2, default=str) def load(self) -> None: """Carrega estado do disco.""" if not self.persist_path or not os.path.exists(self.persist_path): return with open(self.persist_path) as f: data = json.load(f) self.instructions = data.get("instructions", []) self.usage_stats = data.get("usage_stats", self.usage_stats) self.error_history = deque( data.get("error_history", []), maxlen=self.max_error_history, ) def _maybe_persist(self) -> None: if self.persist_path: self.save() # ============================================================================ # 2.2. MCP (Model Context Protocol) — conexão e contexto # ============================================================================ @dataclass class MCPMessage: """Mensagem no formato Model Context Protocol.""" role: str # 'user', 'assistant', 'tool', 'system' content: str tool_name: Optional[str] = None tool_call_id: Optional[str] = None metadata: Dict[str, Any] = field(default_factory=dict) @dataclass class MCPContext: """Contexto MCP para uma sessão.""" session_id: str messages: List[MCPMessage] = field(default_factory=list) active_tools: Set[str] = field(default_factory=set) shared_state: Dict[str, Any] = field(default_factory=dict) def add_message(self, msg: MCPMessage) -> None: self.messages.append(msg) def get_context_window(self, max_messages: int = 20) -> List[MCPMessage]: """Retorna as últimas max_messages mensagens.""" return self.messages[-max_messages:] def set_shared(self, key: str, value: Any) -> None: self.shared_state[key] = value def get_shared(self, key: str, default: Any = None) -> Any: return self.shared_state.get(key, default) class MCPConnection: """Conexão MCP entre agente e ferramentas. Permite: - Registrar ferramentas (tools) - Enviar/receber mensagens - Compartilhar estado entre ferramentas - Gerenciar contexto da sessão """ def __init__(self, session_id: Optional[str] = None): self.session_id = session_id or hashlib.md5( str(time.time()).encode() ).hexdigest()[:12] self.context = MCPContext(session_id=self.session_id) self._tools: Dict[str, "ToolWrapper"] = {} self._lock = threading.Lock() def register_tool(self, tool: "ToolWrapper") -> None: with self._lock: self._tools[tool.name] = tool self.context.active_tools.add(tool.name) def unregister_tool(self, name: str) -> None: with self._lock: self._tools.pop(name, None) self.context.active_tools.discard(name) def get_tool(self, name: str) -> Optional["ToolWrapper"]: return self._tools.get(name) def list_tools(self) -> List[str]: return list(self._tools.keys()) def send_message( self, role: str, content: str, tool_name: Optional[str] = None, metadata: Optional[Dict] = None, ) -> MCPMessage: """Envia mensagem no contexto MCP.""" msg = MCPMessage( role=role, content=content, tool_name=tool_name, metadata=metadata or {}, ) self.context.add_message(msg) return msg def get_context_window(self, max_messages: int = 20) -> List[MCPMessage]: return self.context.get_context_window(max_messages) def set_shared_state(self, key: str, value: Any) -> None: self.context.set_shared(key, value) def get_shared_state(self, key: str, default: Any = None) -> Any: return self.context.get_shared(key, default) # ============================================================================ # 2.3. CircuitBreaker — resiliência à falhas # ============================================================================ class CircuitState(Enum): CLOSED = "closed" # Funcionando normalmente OPEN = "open" # Falhando, rejeita chamadas HALF_OPEN = "half_open" # Testando se recuperou class CircuitBreaker: """Circuit breaker para uma ferramenta. Estados: - CLOSED: chamadas passam normalmente. Conta falhas. - OPEN: após N falhas consecutivas, rejeita todas as chamadas. - HALF_OPEN: após timeout, permite 1 chamada de teste. Se sucesso → CLOSED. Se falha → OPEN. Matemática: - P(falha) = n_failures / n_calls - Threshold: se P(falha) > θ → OPEN - Recovery: após T_recovery segundos → HALF_OPEN """ def __init__( self, failure_threshold: int = 5, recovery_timeout_s: float = 30.0, half_open_max_calls: int = 1, verbose: bool = False, ): self.failure_threshold = failure_threshold self.recovery_timeout_s = recovery_timeout_s self.half_open_max_calls = half_open_max_calls self.verbose = verbose self._state = CircuitState.CLOSED self._failure_count = 0 self._success_count = 0 self._last_failure_time = 0.0 self._half_open_calls = 0 self._lock = threading.Lock() @property def state(self) -> CircuitState: with self._lock: if self._state == CircuitState.OPEN: # Verifica se deve transitar para HALF_OPEN. if time.time() - self._last_failure_time > self.recovery_timeout_s: self._state = CircuitState.HALF_OPEN self._half_open_calls = 0 if self.verbose: print(f" CircuitBreaker: OPEN → HALF_OPEN") return self._state def can_execute(self) -> bool: """Verifica se a chamada pode prosseguir.""" state = self.state if state == CircuitState.CLOSED: return True if state == CircuitState.HALF_OPEN: with self._lock: if self._half_open_calls < self.half_open_max_calls: self._half_open_calls += 1 return True return False return False # OPEN def record_success(self) -> None: """Registra sucesso.""" with self._lock: if self._state == CircuitState.HALF_OPEN: self._success_count += 1 self._state = CircuitState.CLOSED self._failure_count = 0 if self.verbose: print(f" CircuitBreaker: HALF_OPEN → CLOSED (recuperou)") elif self._state == CircuitState.CLOSED: self._failure_count = 0 # reset consecutive failures def record_failure(self) -> None: """Registra falha.""" with self._lock: self._failure_count += 1 self._last_failure_time = time.time() if self._state == CircuitState.HALF_OPEN: self._state = CircuitState.OPEN if self.verbose: print(f" CircuitBreaker: HALF_OPEN → OPEN (falhou no teste)") elif self._state == CircuitState.CLOSED: if self._failure_count >= self.failure_threshold: self._state = CircuitState.OPEN if self.verbose: print(f" CircuitBreaker: CLOSED → OPEN " f"({self._failure_count} falhas consecutivas)") def reset(self) -> None: """Reseta para CLOSED.""" with self._lock: self._state = CircuitState.CLOSED self._failure_count = 0 self._success_count = 0 self._half_open_calls = 0 def get_state(self) -> Dict[str, Any]: return { "state": self.state.value, "failure_count": self._failure_count, "success_count": self._success_count, "failure_threshold": self.failure_threshold, "recovery_timeout_s": self.recovery_timeout_s, } # ============================================================================ # 2.3. ToolWrapper — envolve ferramenta com memória + circuit breaker # ============================================================================ class ToolWrapper: """Envolve uma ferramenta (callable) com: - ToolMemory (memória individualizada) - CircuitBreaker (resiliência) - Timeout (evita hang) - Cache (evita recomputação) Uma ferramenta é qualquer callable: fn(input) -> output. """ def __init__( self, name: str, fn: Callable, description: str = "", timeout_s: float = 30.0, enable_cache: bool = True, enable_circuit_breaker: bool = True, failure_threshold: int = 5, recovery_timeout_s: float = 30.0, persist_path: Optional[str] = None, verbose: bool = False, ): self.name = name self.fn = fn self.description = description self.timeout_s = timeout_s self.verbose = verbose # Componentes. self.memory = ToolMemory( tool_name=name, persist_path=persist_path, ) self.circuit_breaker = CircuitBreaker( failure_threshold=failure_threshold, recovery_timeout_s=recovery_timeout_s, verbose=verbose, ) if enable_circuit_breaker else None self.enable_cache = enable_cache def execute( self, input_data: Any, cache_key: Optional[str] = None, use_cache: bool = True, ) -> Any: """Executa a ferramenta com resiliência. Args: input_data: input para a ferramenta cache_key: chave para cache (None = sem cache) use_cache: se True, tenta cache antes de executar Returns: resultado da ferramenta Raises: ToolCircuitBreakerOpen: se circuit breaker está OPEN ToolTimeoutError: se timeout ToolError: se a ferramenta falha """ # 1. Verifica circuit breaker. if self.circuit_breaker and not self.circuit_breaker.can_execute(): raise ToolCircuitBreakerOpen( f"Ferramenta '{self.name}' em circuit breaker OPEN " f"(estado: {self.circuit_breaker.state.value})" ) # 2. Verifica cache. if self.enable_cache and use_cache and cache_key: cached = self.memory.get_cached(cache_key) if cached is not None: if self.verbose: print(f" Tool '{self.name}': cache hit ({cache_key})") return cached # 3. Executa com timeout. t0 = time.time() success = False try: result = self._execute_with_timeout(input_data) success = True duration = time.time() - t0 # 4. Cache resultado. if self.enable_cache and cache_key: self.memory.cache_result(cache_key, result) # 5. Registra sucesso. if self.circuit_breaker: self.circuit_breaker.record_success() self.memory.record_call(success=True, duration_s=duration) return result except ToolTimeoutError: duration = time.time() - t0 self.memory.record_error("timeout", {"input": str(input_data)[:100]}) self.memory.record_call(success=False, duration_s=duration) if self.circuit_breaker: self.circuit_breaker.record_failure() raise except Exception as e: duration = time.time() - t0 self.memory.record_error(str(e), {"input": str(input_data)[:100]}) self.memory.record_call(success=False, duration_s=duration) if self.circuit_breaker: self.circuit_breaker.record_failure() raise ToolError(f"Ferramenta '{self.name}' falhou: {e}") from e def _execute_with_timeout(self, input_data: Any) -> Any: """Executa fn com timeout via ThreadPoolExecutor.""" future = ThreadPoolExecutor(max_workers=1).submit(self.fn, input_data) try: return future.result(timeout=self.timeout_s) except TimeoutError: raise ToolTimeoutError( f"Ferramenta '{self.name}' timeout após {self.timeout_s}s" ) def get_stats(self) -> Dict[str, Any]: stats = self.memory.get_stats() if self.circuit_breaker: stats["circuit_breaker"] = self.circuit_breaker.get_state() return stats def restart(self) -> None: """Reinicia a ferramenta (reset circuit breaker + clear cache).""" if self.circuit_breaker: self.circuit_breaker.reset() self.memory.clear_cache() if self.verbose: print(f" Tool '{self.name}': reiniciada") # ============================================================================ # 2.3.1. ToolAgentCoordinator — arquitetura distribuída otimizada # ============================================================================ class ToolAgentCoordinator: """Coordenador de Agentes de Ferramentas (arquitetura híbrida). Arquitetura: 1 Coordenador + N Agentes Especialistas (1 por ferramenta). Análise matemática da arquitetura: - Single-Agent (1 agente, N ferramentas): T = Σ T_i (sequencial) - Multi-Agent (N agentes, 1 ferramenta cada): T = max(T_i) + T_coord - Hybrid (este): T = max(T_i) + T_coord, com isolamento de falhas Quando T_coord << max(T_i): speedup ≈ N. Isolamento: se ferramenta i falha, só agente i reinicia (não afeta outros). O Coordenador: 1. Roteia requisições para o agente especialista apropriado 2. Agrega resultados 3. Gerencia resiliência (restart de agentes falhos) 4. Mantém contexto MCP compartilhado """ def __init__( self, n_workers: Optional[int] = None, verbose: bool = False, ): self.n_workers = n_workers or min(4, os.cpu_count() or 4) self.verbose = verbose self._tools: Dict[str, ToolWrapper] = {} self._mcp = MCPConnection() self._executor = ThreadPoolExecutor(max_workers=self.n_workers) self._lock = threading.Lock() def register_tool( self, name: str, fn: Callable, description: str = "", timeout_s: float = 30.0, failure_threshold: int = 5, recovery_timeout_s: float = 30.0, ) -> ToolWrapper: """Registra uma ferramenta.""" tool = ToolWrapper( name=name, fn=fn, description=description, timeout_s=timeout_s, failure_threshold=failure_threshold, recovery_timeout_s=recovery_timeout_s, verbose=self.verbose, ) with self._lock: self._tools[name] = tool self._mcp.register_tool(tool) if self.verbose: print(f" Coordenador: ferramenta '{name}' registrada") return tool def unregister_tool(self, name: str) -> None: with self._lock: self._tools.pop(name, None) self._mcp.unregister_tool(name) def execute_tool( self, name: str, input_data: Any, cache_key: Optional[str] = None, ) -> Any: """Executa uma ferramenta específica.""" tool = self._tools.get(name) if tool is None: raise ToolError(f"Ferramenta '{name}' não registrada") return tool.execute(input_data, cache_key=cache_key) def execute_parallel( self, requests: List[Tuple[str, Any]], fail_fast: bool = False, ) -> Dict[str, Any]: """Executa múltiplas ferramentas em paralelo. Args: requests: lista de (tool_name, input_data) fail_fast: se True, primeira falha cancela todas Retorna: {tool_name: result} (ou {tool_name: error}) """ results: Dict[str, Any] = {} futures: Dict[Future, str] = {} for tool_name, input_data in requests: tool = self._tools.get(tool_name) if tool is None: results[tool_name] = ToolError(f"Ferramenta '{tool_name}' não registrada") continue # Submete. fut = self._executor.submit(self._safe_execute, tool_name, input_data) futures[fut] = tool_name # Coleta. for fut in as_completed(futures): tool_name = futures[fut] try: result = fut.result() results[tool_name] = result except Exception as e: results[tool_name] = e if fail_fast: # Cancela restantes. for f in futures: f.cancel() break return results def _safe_execute(self, tool_name: str, input_data: Any) -> Any: """Executa ferramenta com tratamento de erro.""" try: return self.execute_tool(tool_name, input_data) except ToolCircuitBreakerOpen: # Tenta restart. if self.verbose: print(f" Coordenador: restart da ferramenta '{tool_name}'") tool = self._tools[tool_name] tool.restart() # Tenta novamente após restart. return self.execute_tool(tool_name, input_data) except Exception as e: if self.verbose: print(f" Coordenador: ferramenta '{tool_name}' falhou: {e}") raise def execute_with_fallback( self, primary_tool: str, fallback_tool: str, input_data: Any, ) -> Any: """Executa ferramenta primária; se falhar, usa fallback.""" try: return self.execute_tool(primary_tool, input_data) except (ToolError, ToolCircuitBreakerOpen, ToolTimeoutError) as e: if self.verbose: print(f" Coordenador: '{primary_tool}' falhou ({e}), " f"usando fallback '{fallback_tool}'") return self.execute_tool(fallback_tool, input_data) def get_all_stats(self) -> Dict[str, Dict[str, Any]]: """Estatísticas de todas as ferramentas.""" return {name: tool.get_stats() for name, tool in self._tools.items()} def get_mcp_context(self, max_messages: int = 20) -> List[MCPMessage]: """Retorna contexto MCP.""" return self._mcp.get_context_window(max_messages) def restart_all(self) -> None: """Reinicia todas as ferramentas.""" for tool in self._tools.values(): tool.restart() def summary(self) -> str: lines = [f"ToolAgentCoordinator ({len(self._tools)} ferramentas):"] for name, tool in self._tools.items(): stats = tool.get_stats() cb_state = stats.get("circuit_breaker", {}).get("state", "N/A") lines.append( f" {name}: calls={stats['n_calls']}, " f"success_rate={stats['success_rate']:.1%}, " f"avg_time={stats['avg_time_time_s' if 'avg_time_time_s' in stats else 'avg_time_s']:.3f}s, " f"circuit={cb_state}" ) return "\n".join(lines) def close(self) -> None: """Fecha executor.""" self._executor.shutdown(wait=True) # ============================================================================ # 2.3.1. Arquitetura — análise matemática # ============================================================================ def analyze_architecture( n_tools: int, tool_latencies: List[float], coord_overhead_s: float = 0.01, ) -> Dict[str, Any]: """Analisa 3 arquiteturas e recomenda a melhor. Arquiteturas: 1. Single-Agent: 1 agente executa todas N ferramentas sequencialmente T = Σ T_i 2. Multi-Agent: N agentes, 1 ferramenta cada, em paralelo T = max(T_i) + T_coord 3. Hybrid (Coordinator + N Specialists): mesmo que Multi-Agent mas com isolamento de falhas e memória individualizada T = max(T_i) + T_coord (mesmo tempo, mais resiliência) Args: n_tools: número de ferramentas tool_latencies: latência esperada de cada ferramenta (s) coord_overhead_s: overhead de coordenação (s) Retorna: análise com recomendação """ assert len(tool_latencies) == n_tools total_seq = sum(tool_latencies) max_lat = max(tool_latencies) # Single-Agent. t_single = total_seq # Multi-Agent / Hybrid. t_parallel = max_lat + coord_overhead_s speedup = t_single / t_parallel if t_parallel > 0 else 0 # Recomendação. if speedup > 1.5: recommendation = "Hybrid (Coordinator + N Specialists)" reason = f"Speedup {speedup:.1f}× justifica overhead de coordenação" elif n_tools <= 2: recommendation = "Single-Agent" reason = f"Poucas ferramentas ({n_tools}), overhead não compensa" else: recommendation = "Hybrid" reason = f"Isolamento de falhas importante para {n_tools} ferramentas" # Análise de resiliência. # Em Single-Agent: 1 falha derruba tudo (P(falha total) = 1 - ∏(1-p_i)) # Em Hybrid: 1 falha afeta só 1 agente (P(falha total) = ∏ p_i, muito menor) return { "n_tools": n_tools, "tool_latencies": tool_latencies, "single_agent_time_s": t_single, "parallel_time_s": t_parallel, "speedup": speedup, "coord_overhead_s": coord_overhead_s, "recommendation": recommendation, "reason": reason, "resilience_advantage": ( "Hybrid: 1 falha afeta 1 agente (isolamento). " "Single-Agent: 1 falha derruba todas." ), }