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