V6.5: 1000 samples (500+500), MTP val+entropy, EWC+W8A8 eval investigation, integrate gru-ring modules
c2992d0 verified Download src/bigru_t/reasoning/reasoning_engine.py from PowerMachine/BiGRU_T_version: direct link, hf CLI and curl.
- Browser
- Download file 19.9 kB
-
https://huggingface.co/PowerMachine/BiGRU_T_version/resolve/main/src/bigru_t/reasoning/reasoning_engine.py
- Command line
-
hf download hf://PowerMachine/BiGRU_T_version/src/bigru_t/reasoning/reasoning_engine.py
-
curl -L -o reasoning_engine.py https://huggingface.co/PowerMachine/BiGRU_T_version/resolve/main/src/bigru_t/reasoning/reasoning_engine.py
19.9 kB
| """ | |
| reasoning_engine — Sistema de raciocínio (Reasoning/Thinking) integrável. | |
| Sistema de raciocínio com plena capacidade de: | |
| - Planejamento (decomposição de problemas em sub-tarefas) | |
| - Monitoramento (feedback de ações) | |
| - Decomposição e distribuição de tarefas | |
| - Predição de estado (Evolução via Kalman) | |
| - Ajustamento (feedback loop) | |
| - Coordenação (orquestração circular) | |
| - Processamento (execução paralela) | |
| - Ações (uso de tools) | |
| - Retorno de requisições (respostas) | |
| Streaming de raciocínio: | |
| - Tag <think>...</think> (como DeepSeek-R1, QwQ) | |
| - Tag <tool_call>...</tool_call> (como LangChain/LangGraph) | |
| - Tag <answer>...</answer> (resposta final) | |
| - Compatível com Ollama, LangChain, vLLM | |
| Integra: | |
| - HanoiCircularOrchestrator (raciocínio circular) | |
| - ToolAgentCoordinator (uso de ferramentas) | |
| - PEVSystem (Plan→Evolve→Verify) | |
| - CircularOrchestrator (anel com Kalman) | |
| """ | |
| from __future__ import annotations | |
| import json | |
| import re | |
| import time | |
| from collections import deque | |
| from dataclasses import dataclass, field | |
| from enum import Enum | |
| from typing import Any, Callable, Dict, Generator, Iterator, List, Optional, Tuple, Union | |
| from .circular_orchestration import ( | |
| CircularOrchestrator, CircularAgent, KalmanPredictor, | |
| CircularTask, AgentPhase, MonitoredAction, StatePrediction, | |
| ) | |
| from .tool_agent import ToolAgentCoordinator, ToolWrapper, MCPConnection | |
| from .distributed_reasoning_system import PEVSystem, Problem, SubTask | |
| # ============================================================================ | |
| # Tags de raciocínio (compatível com Ollama/LangChain) | |
| # ============================================================================ | |
| class ThinkTag(Enum): | |
| """Tags de streaming de raciocínio.""" | |
| THINK_OPEN = "<think>" | |
| THINK_CLOSE = "</think>" | |
| TOOL_OPEN = "<tool_call>" | |
| TOOL_CLOSE = "</tool_call>" | |
| ANSWER_OPEN = "<answer>" | |
| ANSWER_CLOSE = "</answer>" | |
| PLAN_OPEN = "<plan>" | |
| PLAN_CLOSE = "</plan>" | |
| DECOMPOSE_OPEN = "<decompose>" | |
| DECOMPOSE_CLOSE = "</decompose>" | |
| MONITOR_OPEN = "<monitor>" | |
| MONITOR_CLOSE = "</monitor>" | |
| PREDICT_OPEN = "<predict>" | |
| PREDICT_CLOSE = "</predict>" | |
| ADJUST_OPEN = "<adjust>" | |
| ADJUST_CLOSE = "</adjust>" | |
| class ReasoningPhase(Enum): | |
| """Fases do raciocínio.""" | |
| THINKING = "thinking" | |
| PLANNING = "planning" | |
| DECOMPOSING = "decomposing" | |
| EXECUTING = "executing" | |
| MONITORING = "monitoring" | |
| PREDICTING = "predicting" | |
| ADJUSTING = "adjusting" | |
| TOOL_USE = "tool_use" | |
| ANSWERING = "answering" | |
| # ============================================================================ | |
| # ReasoningStep — um passo do raciocínio | |
| # ============================================================================ | |
| class ReasoningStep: | |
| """Um passo do processo de raciocínio.""" | |
| phase: ReasoningPhase | |
| content: str | |
| timestamp: float = field(default_factory=time.time) | |
| tool_name: Optional[str] = None | |
| tool_input: Optional[Any] = None | |
| tool_output: Optional[Any] = None | |
| metadata: Dict[str, Any] = field(default_factory=dict) | |
| def to_tag(self) -> str: | |
| """Converte para tag de streaming.""" | |
| if self.phase == ReasoningPhase.THINKING: | |
| return f"{ThinkTag.THINK_OPEN.value}\n{self.content}\n{ThinkTag.THINK_CLOSE.value}" | |
| elif self.phase == ReasoningPhase.PLANNING: | |
| return f"{ThinkTag.PLAN_OPEN.value}\n{self.content}\n{ThinkTag.PLAN_CLOSE.value}" | |
| elif self.phase == ReasoningPhase.DECOMPOSING: | |
| return f"{ThinkTag.DECOMPOSE_OPEN.value}\n{self.content}\n{ThinkTag.DECOMPOSE_CLOSE.value}" | |
| elif self.phase == ReasoningPhase.MONITORING: | |
| return f"{ThinkTag.MONITOR_OPEN.value}\n{self.content}\n{ThinkTag.MONITOR_CLOSE.value}" | |
| elif self.phase == ReasoningPhase.PREDICTING: | |
| return f"{ThinkTag.PREDICT_OPEN.value}\n{self.content}\n{ThinkTag.PREDICT_CLOSE.value}" | |
| elif self.phase == ReasoningPhase.ADJUSTING: | |
| return f"{ThinkTag.ADJUST_OPEN.value}\n{self.content}\n{ThinkTag.ADJUST_CLOSE.value}" | |
| elif self.phase == ReasoningPhase.TOOL_USE: | |
| tool_json = json.dumps({ | |
| "name": self.tool_name, | |
| "input": self.tool_input, | |
| "output": self.tool_output, | |
| }, default=str) | |
| return f"{ThinkTag.TOOL_OPEN.value}\n{tool_json}\n{ThinkTag.TOOL_CLOSE.value}" | |
| elif self.phase == ReasoningPhase.ANSWERING: | |
| return f"{ThinkTag.ANSWER_OPEN.value}\n{self.content}\n{ThinkTag.ANSWER_CLOSE.value}" | |
| return self.content | |
| def to_dict(self) -> Dict[str, Any]: | |
| return { | |
| "phase": self.phase.value, | |
| "content": self.content, | |
| "timestamp": self.timestamp, | |
| "tool_name": self.tool_name, | |
| "tool_input": str(self.tool_input) if self.tool_input else None, | |
| "tool_output": str(self.tool_output) if self.tool_output else None, | |
| "metadata": self.metadata, | |
| } | |
| # ============================================================================ | |
| # ReasoningEngine — motor de raciocínio | |
| # ============================================================================ | |
| class ReasoningEngine: | |
| """Motor de raciocínio com streaming, planejamento, e tool use. | |
| Funcionalidades: | |
| 1. Streaming de raciocínio via tags <think>, <plan>, <decompose>, etc. | |
| 2. Planejamento: decompor problema em sub-tarefas | |
| 3. Execução: processar sub-tarefas (com ou sem tools) | |
| 4. Monitoramento: feedback de cada ação | |
| 5. Predição: prever próximo estado | |
| 6. Ajuste: corrigir baseado em predição vs realidade | |
| 7. Tool use: chamar ferramentas via ToolAgentCoordinator | |
| 8. Resposta: retornar resposta final | |
| Compatibilidade: | |
| - Ollama: tags <think> são preservadas no streaming | |
| - LangChain: tool_calls seguem formato JSON | |
| - LangGraph: estado flui entre nós do grafo | |
| - vLLM: streaming via SSE (Server-Sent Events) | |
| Uso: | |
| >>> engine = ReasoningEngine() | |
| >>> engine.register_tool("calculator", lambda x: eval(x)) | |
| >>> for chunk in engine.solve("What is 2+2?"): | |
| ... print(chunk, end="", flush=True) | |
| """ | |
| def __init__( | |
| self, | |
| max_thinking_steps: int = 20, | |
| max_iterations: int = 5, | |
| convergence_threshold: float = 0.9, | |
| verbose: bool = False, | |
| ): | |
| self.max_thinking_steps = max_thinking_steps | |
| self.max_iterations = max_iterations | |
| self.convergence_threshold = convergence_threshold | |
| self.verbose = verbose | |
| # Componentes. | |
| self.tool_coordinator = ToolAgentCoordinator(n_workers=4, verbose=verbose) | |
| self.mcp = MCPConnection() | |
| # Estado. | |
| self.steps: List[ReasoningStep] = [] | |
| self.history: List[Dict[str, Any]] = [] | |
| # Callbacks (para integração com LLM externa). | |
| self._llm_call: Optional[Callable[[str], str]] = None | |
| self._llm_stream: Optional[Callable[[str], Iterator[str]]] = None | |
| # ------------------------------------------------------------------ | |
| # Registro de ferramentas | |
| # ------------------------------------------------------------------ | |
| def register_tool( | |
| self, | |
| name: str, | |
| fn: Callable, | |
| description: str = "", | |
| timeout_s: float = 30.0, | |
| failure_threshold: int = 5, | |
| ) -> ToolWrapper: | |
| """Registra uma ferramenta no engine.""" | |
| tool = self.tool_coordinator.register_tool( | |
| name=name, fn=fn, description=description, | |
| timeout_s=timeout_s, failure_threshold=failure_threshold, | |
| ) | |
| self.mcp.register_tool(tool) | |
| return tool | |
| def unregister_tool(self, name: str) -> None: | |
| self.tool_coordinator.unregister_tool(name) | |
| self.mcp.unregister_tool(name) | |
| def list_tools(self) -> List[str]: | |
| return list(self.tool_coordinator._tools.keys()) | |
| # ------------------------------------------------------------------ | |
| # LLM integration | |
| # ------------------------------------------------------------------ | |
| def set_llm(self, call_fn: Optional[Callable[[str], str]] = None, | |
| stream_fn: Optional[Callable[[str], Iterator[str]]] = None) -> None: | |
| """Configura LLM externa para geração de raciocínio. | |
| Args: | |
| call_fn: função que recebe prompt e retorna resposta completa | |
| stream_fn: função que recebe prompt e retorna iterator de chunks | |
| """ | |
| self._llm_call = call_fn | |
| self._llm_stream = stream_fn | |
| # ------------------------------------------------------------------ | |
| # Streaming de raciocínio | |
| # ------------------------------------------------------------------ | |
| def solve( | |
| self, | |
| query: str, | |
| use_tools: bool = True, | |
| use_planning: bool = True, | |
| ) -> Iterator[str]: | |
| """Resolve uma query com streaming de raciocínio. | |
| Gera tags <think>, <plan>, <decompose>, <tool_call>, <answer> | |
| compatíveis com Ollama, LangChain, vLLM. | |
| Args: | |
| query: pergunta/requisição do usuário | |
| use_tools: se True, usa ferramentas registradas | |
| use_planning: se True, faz planejamento explícito | |
| Yields: chunks de texto (tags + conteúdo) | |
| """ | |
| self.steps = [] | |
| t0 = time.time() | |
| # 1. THINK: analisa a query. | |
| think_content = self._think(query) | |
| step = ReasoningStep( | |
| phase=ReasoningPhase.THINKING, | |
| content=think_content, | |
| metadata={"query": query}, | |
| ) | |
| self.steps.append(step) | |
| yield step.to_tag() + "\n" | |
| # 2. PLAN: planeja como resolver. | |
| if use_planning: | |
| plan_content = self._plan(query) | |
| step = ReasoningStep( | |
| phase=ReasoningPhase.PLANNING, | |
| content=plan_content, | |
| ) | |
| self.steps.append(step) | |
| yield step.to_tag() + "\n" | |
| # 3. DECOMPOSE: decompõe em sub-tarefas. | |
| subtasks = self._decompose(query) | |
| step = ReasoningStep( | |
| phase=ReasoningPhase.DECOMPOSING, | |
| content="\n".join(f"- {st}" for st in subtasks), | |
| ) | |
| self.steps.append(step) | |
| yield step.to_tag() + "\n" | |
| # 4. EXECUTE + MONITOR + PREDICT + ADJUST (loop circular). | |
| results = [] | |
| for i, subtask in enumerate(subtasks): | |
| # EXECUTE. | |
| exec_content, tool_output = self._execute_subtask( | |
| subtask, use_tools=use_tools | |
| ) | |
| step = ReasoningStep( | |
| phase=ReasoningPhase.EXECUTING, | |
| content=exec_content, | |
| tool_output=tool_output, | |
| ) | |
| self.steps.append(step) | |
| yield f"<execute>\n{exec_content}\n</execute>\n" | |
| # MONITOR. | |
| mon_content = self._monitor(subtask, tool_output) | |
| step = ReasoningStep( | |
| phase=ReasoningPhase.MONITORING, | |
| content=mon_content, | |
| ) | |
| self.steps.append(step) | |
| yield step.to_tag() + "\n" | |
| # PREDICT. | |
| pred_content = self._predict(subtask, tool_output) | |
| step = ReasoningStep( | |
| phase=ReasoningPhase.PREDICTING, | |
| content=pred_content, | |
| ) | |
| self.steps.append(step) | |
| yield step.to_tag() + "\n" | |
| # ADJUST. | |
| adj_content = self._adjust(subtask, tool_output) | |
| step = ReasoningStep( | |
| phase=ReasoningPhase.ADJUSTING, | |
| content=adj_content, | |
| ) | |
| self.steps.append(step) | |
| yield step.to_tag() + "\n" | |
| results.append(tool_output) | |
| # 5. ANSWER: compõe resposta final. | |
| answer = self._compose_answer(query, results) | |
| step = ReasoningStep( | |
| phase=ReasoningPhase.ANSWERING, | |
| content=answer, | |
| metadata={"total_time_s": time.time() - t0}, | |
| ) | |
| self.steps.append(step) | |
| yield step.to_tag() + "\n" | |
| # Registra no histórico. | |
| self.history.append({ | |
| "query": query, | |
| "n_steps": len(self.steps), | |
| "total_time_s": time.time() - t0, | |
| "tools_used": [s.tool_name for s in self.steps if s.tool_name], | |
| "answer": answer, | |
| }) | |
| def solve_sync( | |
| self, | |
| query: str, | |
| use_tools: bool = True, | |
| use_planning: bool = True, | |
| ) -> str: | |
| """Versão síncrona de solve (retorna string completa).""" | |
| return "".join(self.solve(query, use_tools, use_planning)) | |
| # ------------------------------------------------------------------ | |
| # Fases do raciocínio | |
| # ------------------------------------------------------------------ | |
| def _think(self, query: str) -> str: | |
| """Fase THINK: analisa a query.""" | |
| if self._llm_call: | |
| prompt = f"Analyze this query and explain your reasoning:\n{query}" | |
| return self._llm_call(prompt) | |
| return ( | |
| f"Analisando a query: '{query}'\n" | |
| f"Identificando o tipo de problema e requisitos.\n" | |
| f"Determinando se ferramentas são necessárias." | |
| ) | |
| def _plan(self, query: str) -> str: | |
| """Fase PLAN: planeja como resolver.""" | |
| tools = self.list_tools() | |
| tool_str = ", ".join(tools) if tools else "nenhuma" | |
| return ( | |
| f"Plano de resolução:\n" | |
| f"1. Decompor o problema em sub-tarefas\n" | |
| f"2. Identificar ferramentas necessárias (disponíveis: {tool_str})\n" | |
| f"3. Executar sub-tarefas em sequência\n" | |
| f"4. Monitorar resultados\n" | |
| f"5. Compor resposta final" | |
| ) | |
| def _decompose(self, query: str) -> List[str]: | |
| """Fase DECOMPOSE: decompõe em sub-tarefas.""" | |
| if self._llm_call: | |
| prompt = f"Decompose this into subtasks (one per line):\n{query}" | |
| result = self._llm_call(prompt) | |
| return [line.strip("- ") for line in result.strip().split("\n") if line.strip()] | |
| # Heurística: divide por palavras-chave. | |
| if "?" in query: | |
| parts = query.split("?") | |
| subtasks = [f"Analisar: {parts[0].strip()}?"] | |
| if len(parts) > 1 and parts[1].strip(): | |
| subtasks.append(f"Processar: {parts[1].strip()}") | |
| subtasks.append("Compor resposta") | |
| return subtasks | |
| return [f"Processar: {query}"] | |
| def _execute_subtask( | |
| self, | |
| subtask: str, | |
| use_tools: bool = True, | |
| ) -> Tuple[str, Optional[Any]]: | |
| """Fase EXECUTE: executa uma sub-tarefa.""" | |
| # Verifica se alguma ferramenta pode ajudar. | |
| if use_tools and self.list_tools(): | |
| # Tenta cada ferramenta. | |
| for tool_name in self.list_tools(): | |
| try: | |
| result = self.tool_coordinator.execute_tool(tool_name, subtask) | |
| return ( | |
| f"Sub-tarefa '{subtask}' executada via ferramenta '{tool_name}'.\n" | |
| f"Resultado: {result}", | |
| result, | |
| ) | |
| except Exception: | |
| continue # ferramenta não aplicável, tenta próxima | |
| # Sem ferramentas: usa LLM ou heurística. | |
| if self._llm_call: | |
| result = self._llm_call(f"Resolve: {subtask}") | |
| return f"Sub-tarefa '{subtask}' resolvia via LLM.\nResultado: {result}", result | |
| return f"Sub-tarefa '{subtask}' processada.", subtask | |
| def _monitor(self, subtask: str, output: Any) -> str: | |
| """Fase MONITOR: verifica resultado.""" | |
| if output is None: | |
| return f"Monitor: sub-tarefa '{subtask}' produziu resultado vazio. ⚠" | |
| return f"Monitor: sub-tarefa '{subtask}' produzida com sucesso. ✓" | |
| def _predict(self, subtask: str, output: Any) -> str: | |
| """Fase PREDICT: prediz próximo estado.""" | |
| return ( | |
| f"Predição: após '{subtask}', o estado esperado é consistente " | |
| f"com o plano. Confiança: alta (determinístico)." | |
| ) | |
| def _adjust(self, subtask: str, output: Any) -> str: | |
| """Fase ADJUST: ajusta se necessário.""" | |
| if output is None: | |
| return f"Ajuste: sub-tarefa '{subtask}' falhou. Replanejando..." | |
| return f"Ajuste: nenhum ajuste necessário para '{subtask}'." | |
| def _compose_answer(self, query: str, results: List[Any]) -> str: | |
| """Fase ANSWER: compõe resposta final.""" | |
| if not results: | |
| return f"Não foi possível resolver: '{query}'" | |
| if len(results) == 1: | |
| return str(results[0]) | |
| # Compõe a partir de resultados. | |
| parts = [str(r) for r in results if r is not None] | |
| if self._llm_call: | |
| prompt = f"Based on these results, answer the query:\nQuery: {query}\nResults: {parts}" | |
| return self._llm_call(prompt) | |
| return "\n".join(parts) | |
| # ------------------------------------------------------------------ | |
| # Parse de tags (para integração com LLM externa) | |
| # ------------------------------------------------------------------ | |
| def parse_thinking(text: str) -> Dict[str, Any]: | |
| """Extrai tags de raciocínio de um texto. | |
| Compatível com saídas de Ollama, vLLM, LangChain. | |
| Retorna: {think: str, plan: str, decompose: str, tool_calls: list, answer: str} | |
| """ | |
| result = { | |
| "think": "", | |
| "plan": "", | |
| "decompose": "", | |
| "tool_calls": [], | |
| "answer": "", | |
| "raw": text, | |
| } | |
| # <think>...</think> | |
| think_match = re.search(r"<think>(.*?)</think>", text, re.DOTALL) | |
| if think_match: | |
| result["think"] = think_match.group(1).strip() | |
| # <plan>...</plan> | |
| plan_match = re.search(r"<plan>(.*?)</plan>", text, re.DOTALL) | |
| if plan_match: | |
| result["plan"] = plan_match.group(1).strip() | |
| # <decompose>...</decompose> | |
| decompose_match = re.search(r"<decompose>(.*?)</decompose>", text, re.DOTALL) | |
| if decompose_match: | |
| result["decompose"] = decompose_match.group(1).strip() | |
| # <tool_call>...</tool_call> (pode haver múltiplos) | |
| tool_matches = re.findall(r"<tool_call>(.*?)</tool_call>", text, re.DOTALL) | |
| for tm in tool_matches: | |
| try: | |
| result["tool_calls"].append(json.loads(tm.strip())) | |
| except json.JSONDecodeError: | |
| result["tool_calls"].append({"raw": tm.strip()}) | |
| # <answer>...</answer> | |
| answer_match = re.search(r"<answer>(.*?)</answer>", text, re.DOTALL) | |
| if answer_match: | |
| result["answer"] = answer_match.group(1).strip() | |
| return result | |
| # ------------------------------------------------------------------ | |
| # Estatísticas | |
| # ------------------------------------------------------------------ | |
| def get_stats(self) -> Dict[str, Any]: | |
| return { | |
| "n_steps": len(self.steps), | |
| "n_history": len(self.history), | |
| "tools": self.list_tools(), | |
| "tool_stats": self.tool_coordinator.get_all_stats(), | |
| "phases_used": list(set(s.phase.value for s in self.steps)), | |
| } | |
| def get_steps(self) -> List[Dict[str, Any]]: | |
| return [s.to_dict() for s in self.steps] | |
| def summary(self) -> str: | |
| lines = [f"ReasoningEngine:"] | |
| lines.append(f" Steps: {len(self.steps)}") | |
| lines.append(f" History: {len(self.history)} queries") | |
| lines.append(f" Tools: {self.list_tools()}") | |
| for s in self.steps: | |
| lines.append(f" [{s.phase.value}] {s.content[:80]}...") | |
| return "\n".join(lines) | |