spark_colony / agents /base.py
diwash-barla1's picture
refactor: decompose app into modular domain packages for v2.5
0f336cf
Raw
History Blame Contribute Delete
5.08 kB
from typing import Any, Dict, List, Optional
from core.logging import logger
from database.db import DatabaseManager
from memory.blackboard import ColonyBlackboard
from schemas.enums import AgentState
from schemas.models import AgentCapabilityProfile, AgentMessage, BlackboardEntry
from telemetry.event_bus import EventBus
from telemetry.message_bus import MessageBus
class BaseAgent:
"""Abstract Base Class enforcing strict agent rules, journal logging, capability profiles, and message bus integration."""
def __init__(
self,
agent_id: str,
name: str,
role: str,
db: DatabaseManager,
message_bus: MessageBus,
event_bus: EventBus,
capabilities: Optional[List[str]] = None,
tools: Optional[List[str]] = None,
is_dynamic: bool = False,
):
self.agent_id = agent_id
self.name = name
self.role = role
self.db = db
self.message_bus = message_bus
self.event_bus = event_bus
self.is_dynamic = is_dynamic
self.state: AgentState = AgentState.IDLE
self.current_task: Optional[str] = None
self.confidence: float = 100.0
self.enabled: bool = True
self.capability_profile = AgentCapabilityProfile(
capabilities=capabilities or ["General Reasoning", "Task Execution"],
tools=tools or ["Journal", "MessageBus", "MemoryVault"],
confidence=100.0,
)
async def publish_to_blackboard(self, blackboard: ColonyBlackboard, mission_id: str, topic: str, data: Dict[str, Any]) -> BlackboardEntry:
return await blackboard.publish(mission_id, self.agent_id, topic, data)
async def set_state(self, state: AgentState, current_task: Optional[str] = None, confidence: Optional[float] = None) -> None:
self.state = state
if current_task is not None:
self.current_task = current_task
if confidence is not None:
self.confidence = confidence
logger.info(f"Agent [{self.name}] State -> {self.state.value} | Task: {self.current_task}")
await self.db.upsert_agent(
agent_id=self.agent_id,
name=self.name,
role=self.role,
state=self.state,
current_task=self.current_task,
confidence=self.confidence,
enabled=self.enabled,
)
await self.event_bus.emit(
event_type="AgentStateChanged",
mission_id="system",
agent_name=self.name,
data={"state": self.state.value, "task": self.current_task, "confidence": self.confidence, "enabled": self.enabled},
)
async def write_journal(self, mission_id: str, entry: str) -> None:
logger.info(f"Journal [{self.name}]: {entry}")
await self.db.save_journal(self.agent_id, mission_id, entry)
await self.event_bus.emit(
event_type="JournalUpdated",
mission_id=mission_id,
agent_name=self.name,
data={"entry": entry},
)
async def record_memory(self, mission_id: str, content: str, tags: List[str]) -> str:
mem_id = await self.db.save_memory(self.agent_id, mission_id, content, tags)
await self.event_bus.emit(
event_type="MemoryUpdated",
mission_id=mission_id,
agent_name=self.name,
data={"memory_id": mem_id, "content": content[:100], "tags": tags},
)
return mem_id
async def send_message(self, recipient: str, mission_id: str, status_msg: str, summary: str, next_request: str, confidence: float) -> None:
msg = AgentMessage(
sender=self.name,
recipient=recipient,
mission_id=mission_id,
status=status_msg,
summary=summary,
next_request=next_request,
confidence=confidence,
)
await self.message_bus.publish(msg)
class DynamicWorkerAgent(BaseAgent):
"""Dynamically spawned worker agent assigned to temporary mission tasks."""
def __init__(
self,
agent_id: str,
name: str,
role: str,
mission_id: str,
db: DatabaseManager,
message_bus: MessageBus,
event_bus: EventBus,
capabilities: Optional[List[str]] = None,
tools: Optional[List[str]] = None,
):
super().__init__(agent_id, name, role, db, message_bus, event_bus, capabilities, tools, is_dynamic=True)
self.assigned_mission_id = mission_id
async def execute_task(self, task_description: str, task_fn) -> Any:
await self.set_state(AgentState.PLANNING, current_task=task_description)
await self.write_journal(self.assigned_mission_id, f"Dynamic Worker [{self.name}] starting task: {task_description}")
res = await task_fn()
await self.set_state(AgentState.COMPLETED, current_task="Task Completed")
await self.write_journal(self.assigned_mission_id, f"Dynamic Worker [{self.name}] completed task: {task_description}")
return res