Tribh/devops-copilot
0
1"""2engine.py — WorkflowEngine with async persistence, approval gate, and telemetry.3"""4from typing import Optional5import json6import uuid7 8from devops_copilot.agents.workflow_agents import PlannerAgent, ExecutorAgent9from devops_copilot.agents.base import AgentState10from devops_copilot.core.memory import MemorySystem11from devops_copilot.core.persistence import PersistenceLayer12from devops_copilot.core.observability import start_metrics_server13from devops_copilot.core.telemetry import tracer14from devops_copilot.utils.logger import logger15 16 17class WorkflowEngine:18 """Orchestrates the Multi-Agent ReAct loop with async persistence."""19 20 def __init__(self, db_path: str = "agentnexus_state.db",21 memory_dir: str = "./chroma_db", run_metrics: bool = True):22 self.planner = PlannerAgent()23 self.executor = ExecutorAgent()24 self.memory = MemorySystem(persist_directory=memory_dir)25 self.persistence = PersistenceLayer(db_path=db_path)26 27 if run_metrics:28 try:29 start_metrics_server()30 except Exception as e:31 logger.warning(f"Could not start metrics server: {e}")32 33 async def run(self, user_request: str, session_id: Optional[str] = None,34 max_steps: int = 5) -> str:35 """Run the incremental Plan → Execute → Reflect loop.36 37 Automatically awaits setup() on first use.38 Blocks on PENDING_APPROVAL and resumes when the API sets human_approved=True.39 """40 if not session_id:41 session_id = str(uuid.uuid4())42 43 # Ensure DB tables exist (idempotent)44 await self.persistence.setup()45 46 # Load or create session state47 stored_state = await self.persistence.load_session(session_id)48 if stored_state:49 state = AgentState(**stored_state)50 logger.info(f"Resumed session {session_id}")51 else:52 state = AgentState(session_id=session_id)53 logger.info(f"Started new session {session_id}")54 55 results = []56 for i in range(max_steps):57 # 0. Context retrieval58 memories = self.memory.search_memories(user_request)59 context_str = str(memories.get("documents", []))60 61 # Start OTel-style trace for this turn62 turn_trace = tracer.start_trace(f"Turn {i+1}", parent_id=session_id)63 64 # 1. Incremental Planning — one step at a time65 logger.info(f"--- Step {i+1} Planning ---")66 plan = await self.planner.chat(67 f"Context: {context_str}\nRequest: {user_request}\nPrevious results: {results}",68 state69 )70 71 if not plan.steps:72 logger.info("Planner signals completion (empty steps).")73 turn_trace.finish(status="completed", metadata={"info": "no more steps"})74 break75 76 step = plan.steps[0]77 78 # 2. Human-in-the-loop gate79 logger.info(f"Executing step: {step.tool_name}")80 needs_approval = "REQUIRES_APPROVAL" in step.thought.upper()81 82 if needs_approval and not state.metadata.get("human_approved"):83 logger.warning(f"⏸ ACTION BLOCKED pending approval: {step.tool_name}")84 results.append({85 "step": i + 1,86 "status": "PENDING_APPROVAL",87 "tool": step.tool_name,88 "thought": step.thought,89 "approve_endpoint": f"POST /sessions/{session_id}/approve"90 })91 await self.persistence.save_session(session_id, state.model_dump())92 return json.dumps(results, indent=2)93 94 result = await self.executor.chat(step, state)95 results.append({"step": i + 1, "status": "executed",96 "tool": step.tool_name, "result": result})97 98 # Reset approval flag after single use99 state.metadata["human_approved"] = False100 101 # 3. Persist state102 await self.persistence.save_session(session_id, state.model_dump())103 104 # Finish trace105 turn_trace.finish(metadata={106 "tool": step.tool_name,107 "success": "Error" not in str(result)108 })109 110 if "FINISH" in step.thought.upper():111 break112 113 # 4. Long-term memory114 self.memory.add_memories(115 documents=[f"Request: {user_request}\nLog: {results}"],116 metadatas=[{"session_id": session_id}],117 ids=[str(uuid.uuid4())]118 )119 120 return json.dumps(results, indent=2)121 