Team Ai
Apppublic

Tribh/devops-copilot

sourceHugging Faceupdated 8mo agoView on Hugging Face
0likes
engine.py121 linesDownload Raw Back to core
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