import asyncio import time from typing import Callable, Optional from loguru import logger from pipecat.frames.frames import ( CancelFrame, EndFrame, Frame, InterruptionFrame, LLMContextFrame, LLMFullResponseEndFrame, LLMFullResponseStartFrame, LLMTextFrame, StartFrame, ) from pipecat.processors.frame_processor import FrameDirection, FrameProcessor class DualEngineProcessor(FrameProcessor): """Dual-Engine Orchestrator. Combines a fast local engine (macOS Foundation Model / Apple MLX) for instant sub-400ms voice feedback with a deep engine (Hermes Agent / Luna / Gemma) for deep reasoning, tool execution, and workspace memory. """ def __init__( self, *, fast_llm: Optional[FrameProcessor] = None, deep_llm: FrameProcessor, observer: Optional[Callable[[str], None]] = None, **kwargs, ): super().__init__(**kwargs) self._fast_llm = fast_llm self._deep_llm = deep_llm self._observer = observer self._current_user_text: str = "" self._fast_task: Optional[asyncio.Task] = None self._deep_task: Optional[asyncio.Task] = None self._fast_spoken: bool = False self._deep_spoken: bool = False self._last_tool_phrase: str = "" if hasattr(self._deep_llm, "_on_tool_event"): self._deep_llm._on_tool_event = self.handle_tool_signal def handle_tool_signal(self, detail: str): if not detail or self._deep_spoken: return detail_lower = detail.lower() if "read" in detail_lower or "view" in detail_lower or "cat" in detail_lower: phrase = "Inspecting project files." elif "search" in detail_lower or "grep" in detail_lower or "find" in detail_lower: phrase = "Searching the codebase." elif "exec" in detail_lower or "run" in detail_lower or "command" in detail_lower: phrase = "Running command." else: phrase = "Working on that." if phrase == self._last_tool_phrase: return self._last_tool_phrase = phrase logger.info(f"🗣 [DualEngine Voice Signal]: {phrase!r} (from tool event: {detail[:60]!r})") asyncio.create_task(self._speak_tool_update(phrase)) async def _speak_tool_update(self, phrase: str): try: await self.push_frame(LLMFullResponseStartFrame()) await self.push_frame(LLMTextFrame(phrase)) await self.push_frame(LLMFullResponseEndFrame()) except Exception as e: logger.debug(f"Tool voice update error: {e}") async def setup(self, task_manager): await super().setup(task_manager) if self._fast_llm and hasattr(self._fast_llm, "setup"): await self._fast_llm.setup(task_manager) if self._deep_llm and hasattr(self._deep_llm, "setup"): await self._deep_llm.setup(task_manager) def set_task_manager(self, task_manager): super().set_task_manager(task_manager) if self._fast_llm and hasattr(self._fast_llm, "set_task_manager"): self._fast_llm.set_task_manager(task_manager) if self._deep_llm and hasattr(self._deep_llm, "set_task_manager"): self._deep_llm.set_task_manager(task_manager) def link(self, processor: "FrameProcessor"): super().link(processor) if self._fast_llm and hasattr(self._fast_llm, "link"): self._fast_llm.link(processor) if self._deep_llm and hasattr(self._deep_llm, "link"): self._deep_llm.link(processor) async def process_frame(self, frame: Frame, direction: FrameDirection): await super().process_frame(frame, direction) if isinstance(frame, StartFrame): if self._fast_llm: await self._fast_llm.process_frame(frame, direction) await self._deep_llm.process_frame(frame, direction) await self.push_frame(frame, direction) elif isinstance(frame, (EndFrame, CancelFrame)): await self._cancel_active_tasks() if self._fast_llm: await self._fast_llm.process_frame(frame, direction) await self._deep_llm.process_frame(frame, direction) await self.push_frame(frame, direction) elif isinstance(frame, InterruptionFrame): await self._cancel_active_tasks() if self._fast_llm: await self._fast_llm.process_frame(frame, direction) await self._deep_llm.process_frame(frame, direction) await self.push_frame(frame, direction) elif isinstance(frame, LLMContextFrame): text = self._extract_user_text(frame.context) if text: await self.start_dual_turn(text) else: await self.push_frame(frame, direction) else: await self.push_frame(frame, direction) def start_turn_direct(self, text: str): utterance = text.strip() if not utterance: return asyncio.create_task(self.start_dual_turn(utterance)) async def start_dual_turn(self, text: str): utterance = text.strip() if not utterance: return await self._cancel_active_tasks() self._current_user_text = utterance self._fast_spoken = False self._deep_spoken = False self._suppress_deep = False logger.info(f"⚡ [DualEngine] Starting turn for prompt: {utterance!r}") t0 = time.perf_counter() # Start deep Hermes processing in background self._deep_task = asyncio.create_task(self._run_deep_path(utterance, t0)) # Dispatch fast-path acknowledgment concurrently if self._fast_llm and hasattr(self._fast_llm, "_run_turn"): self._fast_task = asyncio.create_task(self._run_fast_path(utterance, t0)) async def _run_fast_path(self, utterance: str, t0: float): try: fast_prompt = ( "You are a fast voice assistant.\n" "Rules:\n" "1. If the prompt is a simple greeting or fully answered by a short sentence, " "end your answer with [COMPLETE].\n" "2. If it requires deep search/code/tools, use a soft natural human filler " '(e.g., "Ah, let me check that...", "Hmm, let me look into that.") and end with [NEEDS_DEEP].\n' "3. Keep output under 15 words.\n\n" f"User prompt: {utterance!r}" ) chunks: list[str] = [] if hasattr(self._fast_llm, "_run_turn_cli"): await self._fast_llm._run_turn_cli(fast_prompt, chunks) elif hasattr(self._fast_llm, "_run_turn"): await self._fast_llm._run_turn(fast_prompt) t1 = time.perf_counter() raw_text = " ".join(chunks).strip() is_complete = "[COMPLETE]" in raw_text cleaned_text = raw_text.replace("[COMPLETE]", "").replace("[NEEDS_DEEP]", "").strip() if cleaned_text and not self._deep_spoken: self._fast_spoken = True if is_complete: self._suppress_deep = True logger.info(f"⚡ [DualEngine Speculative Routing]: Query marked COMPLETE by fast model. Suppressing redundant deep response.") logger.info(f"⏱ [DualEngine Fast-Path ({int((t1-t0)*1000)}ms)]: {cleaned_text!r} (Complete: {is_complete})") try: import web_server web_server.broadcast_event("fast_reply", { "text": cleaned_text, "is_complete": is_complete, }) except Exception: pass await self.push_frame(LLMFullResponseStartFrame()) await self.push_frame(LLMTextFrame(cleaned_text)) await self.push_frame(LLMFullResponseEndFrame()) except asyncio.CancelledError: pass except Exception as e: logger.debug(f"DualEngine fast-path error: {e}") async def _run_deep_path(self, utterance: str, t0: float): try: # If fast model marked turn COMPLETE, run deep Hermes in background history mode if self._suppress_deep: logger.info("Hermes deep path running silently in background history sync mode...") if hasattr(self._deep_llm, "_run_turn"): try: await self._deep_llm._run_turn(utterance, suppress_output=self._suppress_deep) except TypeError: await self._deep_llm._run_turn(utterance) t1 = time.perf_counter() self._deep_spoken = True logger.info(f"⏱ [DualEngine Deep-Path ({int((t1-t0)*1000)}ms)] turn complete.") try: import web_server web_server.broadcast_event("profiling", { "mode": "Dual-Engine (Fast + Deep)", "total_ms": int((t1 - t0) * 1000), }) except Exception: pass except asyncio.CancelledError: pass except Exception as e: logger.error(f"DualEngine deep-path error: {e}") async def _cancel_active_tasks(self): for task in (self._fast_task, self._deep_task): if task and not task.done(): task.cancel() try: await task except asyncio.CancelledError: pass self._fast_task = None self._deep_task = None def _extract_user_text(self, context) -> str: if not context or not hasattr(context, "messages"): return "" for msg in reversed(context.messages): if isinstance(msg, dict) and msg.get("role") == "user": content = msg.get("content", "") if isinstance(content, str): return content elif isinstance(content, list): return " ".join([c.get("text", "") for c in content if isinstance(c, dict)]) return ""