spark_colony / core /runtime.py
diwash-barla1's picture
fix: complete implementation of core runtime module
b5cd855
Raw
History Blame Contribute Delete
9.87 kB
from datetime import datetime, timezone
import time
from typing import Any, Dict, List, Optional
import uuid
try:
import psutil
except ImportError:
psutil = None
from agents.browser_agent import BrowserAgent
from agents.commander import CommanderAgent
from agents.conversation_agent import ConversationAgent
from agents.base import DynamicWorkerAgent
from agents.fact_checker import FactCheckerAgent
from agents.judge import EvidenceJudgeAgent
from agents.memory_agent import MemoryAgent
from agents.researcher import ResearchAgent
from agents.vision_agent import VisionAgent
from agents.writer import WriterAgent
from browser.browser_pool import BrowserPoolManager
from core.config import ConfigEngine
from core.logging import logger
from core.notification import NotificationEngine
from database.db import DatabaseManager
from memory.blackboard import ColonyBlackboard
from models.key_rotator import APIKeyRotationEngine
from models.model_manager import ModelManager
from plugins.plugin_manager import PluginManager
from reasoning.consensus import ConsensusEngine
from reasoning.debate import DebateEngine
from reasoning.discussion import ColonyDiscussionEngine
from reasoning.negotiation import TaskNegotiationEngine
from reasoning.reputation import AgentReputationEngine
from schemas.enums import ActionType, AgentState, AppState, DecisionStage, MissionStatus, NotificationLevel, TimelineStepType
from schemas.models import ApprovalRequestModel, DynamicAgentInfo, MissionContextModel, TimelineEventModel
from telemetry.event_bus import EventBus
from telemetry.message_bus import MessageBus
from telemetry.monitor import SystemMonitor
from telemetry.resource_manager import ColonyResourceManager
from telemetry.websocket import WebSocketManager
from tools.tool_engine import ToolPermissionEngine
from utils.scheduler import TaskScheduler
from utils.version import VersionEngine
class ColonyRuntime:
"""V2 Colony Operating System Runtime Kernel managing dynamic agents, shared blackboard, resource telemetry, universal models, debates, consensus, human approvals, and notifications."""
def __init__(self, db: DatabaseManager, event_bus: EventBus, message_bus: MessageBus):
self.db = db
self.event_bus = event_bus
self.message_bus = message_bus
self.blackboard = ColonyBlackboard(db, event_bus)
self.resource_manager = ColonyResourceManager()
self.key_rotator = APIKeyRotationEngine()
self.model_manager = ModelManager(self.key_rotator, self.resource_manager)
self.tool_engine = ToolPermissionEngine()
self.discussion_engine = ColonyDiscussionEngine(db, event_bus)
self.consensus_engine = ConsensusEngine()
self.debate_engine = DebateEngine(db, event_bus, self.model_manager)
self.negotiation_engine = TaskNegotiationEngine(db, event_bus)
self.reputation_engine = AgentReputationEngine(db)
self.notification_engine = NotificationEngine(db, event_bus)
self.conversation_agent = ConversationAgent(db, message_bus, event_bus, self.model_manager)
self.dynamic_agents: Dict[str, DynamicWorkerAgent] = {}
self.active_contexts: Dict[str, MissionContextModel] = {}
async def request_human_approval(self, mission_id: str, agent_id: str, action_type: ActionType, prompt_message: str) -> ApprovalRequestModel:
req = ApprovalRequestModel(mission_id=mission_id, agent_id=agent_id, action_type=action_type, prompt_message=prompt_message)
await self.db.save_approval_request(req)
await self.db.update_mission_status(mission_id, MissionStatus.PAUSED)
await self.notification_engine.notify(NotificationLevel.ACTION_REQUIRED, f"Approval Required: {action_type.value}", prompt_message, mission_id)
return req
async def record_timeline_step(self, mission_id: str, agent_id: str, step_type: TimelineStepType, description: str, metadata: Optional[Dict[str, Any]] = None):
evt = TimelineEventModel(mission_id=mission_id, agent_id=agent_id, step_type=step_type, description=description, metadata=metadata or {})
await self.db.save_timeline_event(evt)
await self.event_bus.emit("TimelineStep", mission_id, agent_id, {"step_type": step_type.value, "description": description})
async def spawn_worker(
self, role: str, mission_id: str, capabilities: Optional[List[str]] = None, tools: Optional[List[str]] = None
) -> DynamicWorkerAgent:
worker_id = f"worker-{role.lower().replace(' ', '-')}-{uuid.uuid4().hex[:6]}"
worker_name = f"Dynamic {role} ({worker_id[-4:]})"
worker = DynamicWorkerAgent(
agent_id=worker_id,
name=worker_name,
role=f"Dynamic {role}",
mission_id=mission_id,
db=self.db,
message_bus=self.message_bus,
event_bus=self.event_bus,
capabilities=capabilities or [f"{role} Processing", "Dynamic Execution"],
tools=tools or ["Blackboard", "MemoryAccess"],
)
self.dynamic_agents[worker_id] = worker
await worker.set_state(AgentState.IDLE, current_task="Spawned & Waiting")
await self.event_bus.emit("AgentSpawned", mission_id, worker_name, {"agent_id": worker_id, "role": role})
return worker
async def cleanup_mission_workers(self, mission_id: str) -> int:
to_remove = [aid for aid, agent in self.dynamic_agents.items() if agent.assigned_mission_id == mission_id]
for aid in to_remove:
agent = self.dynamic_agents.pop(aid)
await agent.set_state(AgentState.SLEEPING, current_task="Despawned")
await self.event_bus.emit("AgentDespawned", mission_id, agent.name, {"agent_id": aid})
return len(to_remove)
async def update_mission_context(
self,
mission_id: str,
topic: str,
phase: DecisionStage,
progress: float,
owner: str,
priority: int = 5,
cost: float = 0.0,
health: str = "HEALTHY",
) -> MissionContextModel:
ctx = MissionContextModel(
mission_id=mission_id,
topic=topic,
priority=priority,
current_phase=phase,
progress_percent=progress,
owner_agent=owner,
resource_usage={"cpu": psutil.cpu_percent() if psutil else 0.0, "ram": psutil.virtual_memory().percent if psutil else 0.0},
estimated_cost=cost,
health_status=health,
)
self.active_contexts[mission_id] = ctx
await self.db.save_mission_context(ctx)
return ctx
class ColonyOS:
"""Core Operating System Kernel orchestrating autonomous agent sub-systems and V2 Colony Runtime."""
def __init__(self):
self.start_time = time.time()
self.app_state = AppState.BOOTING
self.config = ConfigEngine()
self.db = DatabaseManager()
self.ws_manager = WebSocketManager()
self.event_bus = EventBus(self.ws_manager)
self.message_bus = MessageBus(self.db, self.event_bus)
self.monitor = SystemMonitor(self.start_time)
self.scheduler = TaskScheduler()
self.plugin_manager = PluginManager(self.db)
self.browser_pool = BrowserPoolManager(max_pool_size=self.config.browser_pool_size)
# V2 Colony Runtime Engine
self.runtime = ColonyRuntime(self.db, self.event_bus, self.message_bus)
# Initialize Core Base Agents
self.commander = CommanderAgent(self.db, self.message_bus, self.event_bus)
self.researcher = ResearchAgent(self.db, self.message_bus, self.event_bus)
self.fact_checker = FactCheckerAgent(self.db, self.message_bus, self.event_bus)
self.evidence_judge = EvidenceJudgeAgent(self.db, self.message_bus, self.event_bus)
self.writer = WriterAgent(self.db, self.message_bus, self.event_bus)
self.vision = VisionAgent(self.db, self.message_bus, self.event_bus)
self.browser = BrowserAgent(self.db, self.message_bus, self.event_bus, self.browser_pool)
self.memory = MemoryAgent(self.db, self.message_bus, self.event_bus)
self.conversation_agent = self.runtime.conversation_agent
self.agent_registry = {
self.commander.agent_id: self.commander,
self.researcher.agent_id: self.researcher,
self.fact_checker.agent_id: self.fact_checker,
self.evidence_judge.agent_id: self.evidence_judge,
self.writer.agent_id: self.writer,
self.vision.agent_id: self.vision,
self.browser.agent_id: self.browser,
self.memory.agent_id: self.memory,
self.conversation_agent.agent_id: self.conversation_agent,
}
async def awaken(self) -> None:
"""Boot sequence for Spark Colony OS."""
logger.info("==================================================")
logger.info(" SPARK COLONY OS v5.0 - PRODUCTION KERNEL READY ")
logger.info("==================================================")
for agent in self.agent_registry.values():
await agent.set_state(AgentState.IDLE, current_task=None, confidence=100.0)
self.app_state = AppState.READY
await self.event_bus.emit("ServerStarted", "system", "OS_Kernel", {"state": self.app_state.value, "version": VersionEngine.get_version_info()})
logger.info("Colony Operating System initialized cleanly.")
async def shutdown(self) -> None:
"""Graceful shutdown sequence."""
logger.info("Deactivating Spark Colony OS...")
self.app_state = AppState.SHUTDOWN
for agent in self.agent_registry.values():
await agent.set_state(AgentState.SLEEPING, current_task="Standby", confidence=100.0)
logger.info("Spark Colony OS suspended gracefully.")
colony_os = ColonyOS()