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()