""" distributed_reasoning_system — Sistema distribuído granulável para LRM. Arquitetura Plan → Evolve → Verify (PEV) com: - Planejamento: decomposição de problema em sub-tarefas (granularidade variável) - Evolução: execução paralela de sub-tarefas via agentes distribuídos - Verificação: checagem de erros, ajustes, e feedback para Planejamento Análise matemática profunda --------------------------- 1. Teoria da decomposição ótima (granularidade): Sea um problema P com complexidade C(P). Decompor em n sub-tarefas {s_1,...,s_n} com complexidades {c_1,...,c_n} onde Σc_i ≥ C(P) (overhead de decomposição). O tempo paralelo é max(c_i) + T_coord. Granularidade ótima: n* = argmin_n [max(c_i(n)) + T_coord(n)] onde c_i(n) ≈ C(P)/n (assumindo divisão equilibrada) e T_coord(n) = α·n + β (custo linear de coordenação). Resolvendo: n* = sqrt(C(P) / (α·(1 + β/α))) 2. Convergência do loop PEV: Seja E_k o erro após k iterações do loop Plan→Evolve→Verify. E_{k+1} = ρ·E_k + ε_verify onde ρ ∈ [0,1) é a taxa de correção e ε_verify é o erro residual do verificador. Convergência garantida quando ρ < 1. Taxa de convergência: E_k ≤ ρ^k · E_0 + ε_verify/(1-ρ) 3. Teoria de jogos (Planejador vs Verificador): Planejador escolhe estratégia σ_P, Verificador escolhe σ_V. Jogo zero-sum: u_P(σ_P, σ_V) = -u_V(σ_P, σ_V). Equilíbrio de Nash: (σ_P*, σ_V*) onde nenhum pode melhorar unilateralmente. Valor do jogo: v = max_σP min_σV u_P(σ_P, σ_V). 4. Álgebra de composição de sub-tarefas: Seja T: Problem → Solution uma transformação. Decompor P em {s_i}: T(P) = compose(T(s_1), ..., T(s_n)) onde compose é a álgebra de combinação (monóide associativo). Propriedade: T(P) = compose(T(s_1), ..., T(s_n)) ⊕ E_decompose onde E_decompose é o erro de decomposição (→ 0 com granularidade fina). 5. Distribuição de probabilidade sobre estratégias: P(estratégia | problema) = softmax(score(estratégia, problema) / τ) onde τ é temperatura (→ 0: greedy, → ∞: uniforme). Atualização Bayesiana após verificação: P(estratégia | problema, resultado) ∝ P(resultado | estratégia) · P(estratégia | problema) """ from __future__ import annotations import math 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 numpy as np # ============================================================================ # Exceções # ============================================================================ class ReasoningError(RuntimeError): """Erro no sistema de raciocínio.""" class PlanError(ReasoningError): """Erro no planejamento.""" class EvolveError(ReasoningError): """Erro na evolução/execução.""" class VerifyError(ReasoningError): """Erro na verificação.""" class ConvergenceFailure(ReasoningError): """Sistema PEV não convergiu após max_iterations.""" # ============================================================================ # Estruturas de dados # ============================================================================ class TaskStatus(Enum): PENDING = "pending" PLANNED = "planned" EXECUTING = "executing" VERIFIED = "verified" FAILED = "failed" NEEDS_REPLAN = "needs_replan" @dataclass class SubTask: """Sub-tarefa decomposta do problema principal.""" id: str description: str complexity: float # custo estimado dependencies: Set[str] = field(default_factory=set) # IDs de sub-tarefas dependentes status: TaskStatus = TaskStatus.PENDING result: Optional[Any] = None error: Optional[str] = None verify_score: Optional[float] = None # score de verificação [0, 1] iterations: int = 0 metadata: Dict[str, Any] = field(default_factory=dict) @dataclass class Problem: """Problema a ser resolvido pelo sistema PEV.""" id: str description: str complexity: float # estimativa de complexidade constraints: Dict[str, Any] = field(default_factory=dict) context: Dict[str, Any] = field(default_factory=dict) max_iterations: int = 5 convergence_threshold: float = 0.95 # score mínimo para aceitar @dataclass class Solution: """Solução produzida pelo sistema PEV.""" problem_id: str result: Any verify_score: float iterations: int subtasks: List[SubTask] total_time_s: float converged: bool error_history: List[Dict[str, Any]] = field(default_factory=list) @dataclass class PEVStats: """Estatísticas do loop Plan→Evolve→Verify.""" n_iterations: int = 0 n_subtasks_total: int = 0 n_subtasks_verified: int = 0 n_subtasks_failed: int = 0 n_replans: int = 0 total_plan_time_s: float = 0.0 total_evolve_time_s: float = 0.0 total_verify_time_s: float = 0.0 error_trajectory: List[float] = field(default_factory=list) # ============================================================================ # 1. Planner — decomposição granular de problemas # ============================================================================ class Planner: """Planejador: decomponhe problema em sub-tarefas granulares. Análise matemática da granularidade ótima: n* = sqrt(C(P) / (α·(1 + β/α))) onde C(P) é complexidade, α é custo por sub-tarefa, β é overhead fixo. Estratégias de decomposição: 1. UNIFORM: divide em n sub-tarefas de complexidade igual 2. ADAPTIVE: divide baseado em estrutura do problema 3. RECURSIVE: sub-tarefas podem ser进一步 decompostas """ def __init__( self, max_subtasks: int = 20, min_subtasks: int = 2, alpha: float = 0.1, # custo por sub-tarefa beta: float = 0.05, # overhead fixo strategy: str = "adaptive", verbose: bool = False, ): self.max_subtasks = max_subtasks self.min_subtasks = min_subtasks self.alpha = alpha self.beta = beta self.strategy = strategy self.verbose = verbose def _compute_optimal_granularity(self, complexity: float) -> int: """Computa granularidade ótima n* = sqrt(C / (α·(1+β/α))).""" denom = self.alpha * (1 + self.beta / self.alpha) if denom < 1e-12: return self.min_subtasks n_star = math.sqrt(complexity / denom) return max(self.min_subtasks, min(self.max_subtasks, int(round(n_star)))) def plan( self, problem: Problem, feedback: Optional[List[Dict[str, Any]]] = None, ) -> List[SubTask]: """Decompõe problema em sub-tarefas. Args: problem: problema a decompor feedback: feedback da iteração anterior (para replanejamento) Retorna: lista de SubTasks """ t0 = time.time() n = self._compute_optimal_granularity(problem.complexity) if self.verbose: print(f" Planner: granularidade ótima n*={n} " f"(C={problem.complexity}, α={self.alpha}, β={self.beta})") # Se há feedback, ajusta n. if feedback: # Aumenta granularidade se houve falhas (mais sub-tarefas = mais isolamento). n_failures = sum(1 for f in feedback if f.get("status") == "failed") n = min(self.max_subtasks, n + n_failures) if self.verbose: print(f" Planner: ajustado para n={n} (feedback: {n_failures} falhas)") subtasks = [] if self.strategy == "uniform": subtasks = self._decompose_uniform(problem, n) elif self.strategy == "adaptive": subtasks = self._decompose_adaptive(problem, n) else: subtasks = self._decompose_uniform(problem, n) # Adiciona dependências (grafo DAG). for i, st in enumerate(subtasks): if i > 0: st.dependencies.add(subtasks[i-1].id) if self.verbose: print(f" Planner: {len(subtasks)} sub-tarefas criadas em " f"{time.time()-t0:.3f}s") return subtasks def _decompose_uniform(self, problem: Problem, n: int) -> List[SubTask]: """Decomposição uniforme: n sub-tarefas de complexidade igual.""" c_per = problem.complexity / n return [ SubTask( id=f"{problem.id}_st_{i}", description=f"Sub-tarefa {i}/{n}: parte {i+1} de '{problem.description}'", complexity=c_per, ) for i in range(n) ] def _decompose_adaptive(self, problem: Problem, n: int) -> List[SubTask]: """Decomposição adaptativa: complexidade varia baseado em estrutura. Heurística: primeiro 30% das sub-tarefas têm complexidade maior (mais difíceis), resto é mais simples. """ subtasks = [] # Distribui complexidade: primeiras sub-tarefas mais pesadas. weights = np.array([1.5 if i < n * 0.3 else 0.8 for i in range(n)]) weights = weights / weights.sum() complexities = problem.complexity * weights for i in range(n): subtasks.append(SubTask( id=f"{problem.id}_st_{i}", description=f"Sub-tarefa adaptativa {i}/{n}", complexity=float(complexities[i]), metadata={"weight": float(weights[i])}, )) return subtasks # ============================================================================ # 2. Evolver — execução paralela de sub-tarefas # ============================================================================ class Evolver: """Evoluidor: executa sub-tarefas em paralelo via agentes distribuídos. Usa ThreadPoolExecutor para paralelismo. Respeita dependências (DAG). """ def __init__( self, n_workers: int = 4, timeout_s: float = 30.0, verbose: bool = False, ): self.n_workers = n_workers self.timeout_s = timeout_s self.verbose = verbose def evolve( self, subtasks: List[SubTask], executor_fn: Callable[[SubTask], Any], ) -> List[SubTask]: """Executa sub-tarefas em paralelo, respeitando dependências. Args: subtasks: lista de sub-tarefas executor_fn: função que executa uma sub-tarefa Retorna: sub-tarefas atualizadas com resultados """ t0 = time.time() task_map = {st.id: st for st in subtasks} completed: Set[str] = set() remaining = set(st.id for st in subtasks) # Executa em rounds (respeita DAG). with ThreadPoolExecutor(max_workers=self.n_workers) as pool: while remaining: # Encontra sub-tarefas prontas (dependências satisfeitas). ready = [ sid for sid in remaining if task_map[sid].dependencies.issubset(completed) ] if not ready: # Deadlock: dependências circulares. raise EvolveError( f"Deadlock: sub-tarefas {remaining} têm dependências circulares" ) # Submete prontas. futures = { pool.submit(self._execute_task, task_map[sid], executor_fn): sid for sid in ready } # Coleta. for fut in as_completed(futures, timeout=self.timeout_s): sid = futures[fut] try: result = fut.result() task_map[sid].result = result task_map[sid].status = TaskStatus.EXECUTING completed.add(sid) remaining.discard(sid) if self.verbose: print(f" Evolver: {sid} OK") except Exception as e: task_map[sid].error = str(e) task_map[sid].status = TaskStatus.FAILED completed.add(sid) # marca como processada remaining.discard(sid) if self.verbose: print(f" Evolver: {sid} FALHOU: {e}") if self.verbose: print(f" Evolver: {len(completed)} sub-tarefas em {time.time()-t0:.3f}s") return list(task_map.values()) def _execute_task( self, task: SubTask, executor_fn: Callable, ) -> Any: """Executa uma sub-tarefa com timeout.""" task.status = TaskStatus.EXECUTING task.iterations += 1 return executor_fn(task) # ============================================================================ # 3. Verifier — verificação de resultados # ============================================================================ class Verifier: """Verificador: checa resultados de sub-tarefas. Análise matemática da convergência: E_{k+1} = ρ·E_k + ε_verify E_k ≤ ρ^k · E_0 + ε_verify/(1-ρ) Tipos de verificação: 1. CONSISTENCY: checa consistência entre sub-tarefas 2. BOUNDS: checa se resultado está dentro de limites esperados 3. CROSS_VALIDATION: valida cruzando resultados """ def __init__( self, verify_fn: Optional[Callable[[SubTask, Any], float]] = None, consistency_threshold: float = 0.8, verbose: bool = False, ): self.verify_fn = verify_fn or self._default_verify self.consistency_threshold = consistency_threshold self.verbose = verbose def verify( self, subtasks: List[SubTask], problem: Problem, ) -> Tuple[List[SubTask], float, List[Dict[str, Any]]]: """Verifica resultados das sub-tarefas. Retorna: (subtasks_atualizadas, score_global, feedback) """ t0 = time.time() feedback = [] scores = [] for st in subtasks: if st.status == TaskStatus.FAILED: st.verify_score = 0.0 feedback.append({ "id": st.id, "status": "failed", "error": st.error, "action": "replan", }) scores.append(0.0) continue if st.result is None: st.verify_score = 0.0 feedback.append({ "id": st.id, "status": "no_result", "action": "reexecute", }) scores.append(0.0) continue # Verifica. score = self.verify_fn(st, st.result) st.verify_score = score scores.append(score) if score >= self.consistency_threshold: st.status = TaskStatus.VERIFIED feedback.append({ "id": st.id, "status": "verified", "score": score, "action": "accept", }) else: st.status = TaskStatus.NEEDS_REPLAN feedback.append({ "id": st.id, "status": "needs_replan", "score": score, "action": "replan", }) # Score global: média harmônica (penaliza sub-tarefas com score baixo). # Sub-tarefas falhadas (score=0) contam como epsilon para evitar divisão por zero. eps_score = 0.01 all_scores = [max(s, eps_score) for s in scores] global_score = len(all_scores) / sum(1/s for s in all_scores) if self.verbose: print(f" Verifier: score global={global_score:.4f}, " f"feedback={len(feedback)} itens em {time.time()-t0:.3f}s") return subtasks, global_score, feedback def _default_verify(self, task: SubTask, result: Any) -> float: """Verificação default: checa se resultado não é None e é razoável.""" if result is None: return 0.0 if isinstance(result, (int, float)): return 1.0 if not math.isnan(result) else 0.0 if isinstance(result, str): return 1.0 if len(result) > 0 else 0.0 if isinstance(result, (list, dict)): return 1.0 if len(result) > 0 else 0.5 return 0.8 # default para tipos desconhecidos # ============================================================================ # 4. PEVSystem — sistema distribuído Plan→Evolve→Verify # ============================================================================ class PEVSystem: """Sistema distribuído granulável Plan→Evolve→Verify. Loop principal: 1. PLAN: Planner decompõe problema em sub-tarefas 2. EVOLVE: Evolver executa sub-tarefas em paralelo 3. VERIFY: Verifier checa resultados 4. Se score < threshold: feedback → PLAN (replanejamento) 5. Repete até convergência ou max_iterations Análise de convergência: E_{k+1} = ρ·E_k + ε_verify Convergência garantida quando ρ < 1 (verificador corrige > 0% do erro). k_max = ceil(log((E_0 - ε/(1-ρ)) / (threshold - ε/(1-ρ))) / log(ρ)) Uso: >>> system = PEVSystem() >>> solution = system.solve(problem, executor_fn=my_executor) """ def __init__( self, n_workers: int = 4, max_subtasks: int = 20, planner_strategy: str = "adaptive", alpha: float = 0.1, beta: float = 0.05, consistency_threshold: float = 0.8, verify_fn: Optional[Callable] = None, verbose: bool = False, ): self.planner = Planner( max_subtasks=max_subtasks, alpha=alpha, beta=beta, strategy=planner_strategy, verbose=verbose, ) self.evolver = Evolver( n_workers=n_workers, verbose=verbose, ) self.verifier = Verifier( verify_fn=verify_fn, consistency_threshold=consistency_threshold, verbose=verbose, ) self.verbose = verbose self.stats = PEVStats() def solve( self, problem: Problem, executor_fn: Callable[[SubTask], Any], ) -> Solution: """Resolve problema via loop PEV. Args: problem: problema a resolver executor_fn: função que executa uma sub-tarefa Retorna: Solution """ t0 = time.time() feedback = None all_subtasks: List[SubTask] = [] error_history = [] converged = False for iteration in range(1, problem.max_iterations + 1): self.stats.n_iterations = iteration if self.verbose: print(f"\n--- PEV Iteração {iteration}/{problem.max_iterations} ---") # 1. PLAN. t_plan = time.time() subtasks = self.planner.plan(problem, feedback) self.stats.total_plan_time_s += time.time() - t_plan self.stats.n_subtasks_total += len(subtasks) # 2. EVOLVE. t_evolve = time.time() subtasks = self.evolver.evolve(subtasks, executor_fn) self.stats.total_evolve_time_s += time.time() - t_evolve # 3. VERIFY. t_verify = time.time() subtasks, global_score, feedback = self.verifier.verify(subtasks, problem) self.stats.total_verify_time_s += time.time() - t_verify self.stats.error_trajectory.append(1.0 - global_score) # Conta status. n_verified = sum(1 for st in subtasks if st.status == TaskStatus.VERIFIED) n_failed = sum(1 for st in subtasks if st.status == TaskStatus.FAILED) n_replan = sum(1 for st in subtasks if st.status == TaskStatus.NEEDS_REPLAN) self.stats.n_subtasks_verified += n_verified self.stats.n_subtasks_failed += n_failed self.stats.n_replans += n_replan all_subtasks = subtasks if self.verbose: print(f" Score: {global_score:.4f} " f"(verified={n_verified}, failed={n_failed}, replan={n_replan})") # 4. Convergência? if global_score >= problem.convergence_threshold: converged = True if self.verbose: print(f" ✓ Convergiu em {iteration} iterações") break # 5. Prepara feedback para próxima iteração. error_history.append({ "iteration": iteration, "score": global_score, "n_verified": n_verified, "n_failed": n_failed, "n_replan": n_replan, }) if n_replan > 0: feedback = [f for f in feedback if f.get("action") == "replan"] if not converged: if self.verbose: print(f" ✗ Não convergiu após {problem.max_iterations} iterações") # Compõe solução final. result = self._compose_solution(all_subtasks) total_time = time.time() - t0 return Solution( problem_id=problem.id, result=result, verify_score=self.stats.error_trajectory[-1] if self.stats.error_trajectory else 0.0, iterations=self.stats.n_iterations, subtasks=all_subtasks, total_time_s=total_time, converged=converged, error_history=error_history, ) def _compose_solution(self, subtasks: List[SubTask]) -> Any: """Compõe solução final a partir de sub-tarefas. Álgebra de composição: T(P) = compose(T(s_1), ..., T(s_n)) """ results = [] for st in subtasks: if st.status == TaskStatus.VERIFIED and st.result is not None: results.append(st.result) if not results: return None if len(results) == 1: return results[0] # Tenta concatenar se possível. if all(isinstance(r, (list,)) for r in results): return [item for r in results for item in r] if all(isinstance(r, dict) for r in results): merged = {} for r in results: merged.update(r) return merged # Retorna lista de resultados. return results def get_stats(self) -> PEVStats: return self.stats def summary(self) -> str: s = self.stats return ( f"PEVSystem Stats:\n" f" Iterações: {s.n_iterations}\n" f" Sub-tarefas: {s.n_subtasks_total} total, " f"{s.n_subtasks_verified} verificadas, {s.n_subtasks_failed} falhas\n" f" Replanejamentos: {s.n_replans}\n" f" Tempo: plan={s.total_plan_time_s:.3f}s, " f"evolve={s.total_evolve_time_s:.3f}s, " f"verify={s.total_verify_time_s:.3f}s\n" f" Trajetória de erro: {[f'{e:.4f}' for e in s.error_trajectory]}" ) # ============================================================================ # 5. BayesianStrategySelector — seleção Bayesiana de estratégias # ============================================================================ class BayesianStrategySelector: """Seleção Bayesiana de estratégias de decomposição. Mantém posterior sobre estratégias dado histórico de resultados. P(estratégia | problema, resultado) ∝ P(resultado | estratégia) · P(estratégia | problema) Estratégias: ['uniform', 'adaptive', 'recursive'] Para cada estratégia, mantém Beta(α, β) sobre probabilidade de sucesso. """ def __init__( self, strategies: List[str] = None, prior_alpha: float = 1.0, prior_beta: float = 1.0, temperature: float = 1.0, ): self.strategies = strategies or ["uniform", "adaptive", "recursive"] self.temperature = temperature # Posterior Beta(α, β) para cada estratégia. self.posteriors: Dict[str, Tuple[float, float]] = { s: (prior_alpha, prior_beta) for s in self.strategies } # Histórico de resultados. self.history: List[Dict[str, Any]] = [] def select(self) -> str: """Seleciona estratégia via Thompson sampling.""" samples = {} for s in self.strategies: a, b = self.posteriors[s] samples[s] = np.random.beta(a, b) return max(samples, key=samples.get) def update(self, strategy: str, success: bool, score: float = 0.0) -> None: """Atualiza posterior após observação.""" a, b = self.posteriors[strategy] if success: a += 1 + score # score pondera o sucesso else: b += 1 self.posteriors[strategy] = (a, b) self.history.append({ "strategy": strategy, "success": success, "score": score, "posterior": (a, b), }) def get_probabilities(self) -> Dict[str, float]: """Retorna probabilidade esperada de sucesso por estratégia.""" return { s: a / (a + b) for s, (a, b) in self.posteriors.items() } def summary(self) -> str: lines = ["BayesianStrategySelector:"] for s in self.strategies: a, b = self.posteriors[s] p = a / (a + b) lines.append(f" {s}: Beta({a:.1f}, {b:.1f}), P(sucesso)={p:.3f}") return "\n".join(lines) # ============================================================================ # 6. Análise matemática — funções de análise # ============================================================================ def analyze_convergence( error_trajectory: List[float], rho: float = 0.5, epsilon: float = 0.01, ) -> Dict[str, Any]: """Analisa convergência do loop PEV. E_{k+1} = ρ·E_k + ε E_k ≤ ρ^k · E_0 + ε/(1-ρ) Args: error_trajectory: lista de erros por iteração rho: taxa de correção estimada epsilon: erro residual do verificador Retorna: análise de convergência """ if not error_trajectory: return {"converged": False, "reason": "sem dados"} n = len(error_trajectory) final_error = error_trajectory[-1] initial_error = error_trajectory[0] # Estima ρ empiricamente. if n > 1: rhos = [] for i in range(n - 1): if error_trajectory[i] > epsilon: rho_est = (error_trajectory[i+1] - epsilon) / max(error_trajectory[i], 1e-8) rhos.append(max(0, min(1, rho_est))) rho_empirical = np.mean(rhos) if rhos else rho else: rho_empirical = rho # Prediz k_max para threshold. threshold = 0.05 # 5% de erro if rho_empirical < 1 and (initial_error - epsilon/(1-rho_empirical)) > 0: k_max_pred = math.ceil( math.log((threshold - epsilon/(1-rho_empirical)) / max(initial_error - epsilon/(1-rho_empirical), 1e-8)) / math.log(max(rho_empirical, 1e-8)) ) else: k_max_pred = float('inf') return { "n_iterations": n, "initial_error": initial_error, "final_error": final_error, "rho_empirical": float(rho_empirical), "epsilon": epsilon, "converged": final_error < threshold, "k_max_predicted": k_max_pred if k_max_pred != float('inf') else -1, "error_ratio": final_error / max(initial_error, 1e-8), "trajectory": error_trajectory, } def analyze_granularity( complexity: float, alpha: float = 0.1, beta: float = 0.05, n_range: List[int] = None, ) -> Dict[str, Any]: """Analisa granularidade ótima teórica. n* = sqrt(C / (α·(1+β/α))) Para cada n, computa: - T_parallel(n) = C/n + α·n + β - S(n) = C / T_parallel(n) (speedup vs sequencial) """ if n_range is None: n_range = list(range(1, 21)) results = [] for n in n_range: t_seq = complexity t_par = complexity / n + alpha * n + beta speedup = t_seq / t_par if t_par > 0 else 0 results.append({ "n": n, "t_sequential": t_seq, "t_parallel": t_par, "speedup": speedup, }) # n ótimo teórico. denom = alpha * (1 + beta / alpha) n_star = math.sqrt(complexity / denom) if denom > 0 else 1 # n ótimo empírico. best_n = max(results, key=lambda x: x["speedup"]) return { "complexity": complexity, "alpha": alpha, "beta": beta, "n_star_theoretical": n_star, "n_star_empirical": best_n["n"], "max_speedup": best_n["speedup"], "results": results, }