Spaces:
Sleeping
Sleeping
| 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() | |