Spaces:
Sleeping
Sleeping
File size: 9,869 Bytes
0f336cf b5cd855 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 | 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()
|