David Prince
production: clean source snapshot — no history bloat
71b4454
Raw
History Blame Contribute Delete
3.33 kB
"""
Agent-to-Agent Messaging Bus
Async, in-process pub/sub. Each registered agent gets an inbox (asyncio.Queue).
Messages can be sent directly to one agent (`send`) or broadcast to every
subscriber of a topic (`publish`), e.g. the orchestrator publishing task
status updates that other agents or a UI layer can subscribe to.
"""
from __future__ import annotations
import asyncio
import time
import uuid
from dataclasses import dataclass, field
from typing import Any, Optional
@dataclass(slots=True)
class AgentMessage:
id: str
sender: str
recipient: Optional[str]
topic: Optional[str]
type: str
payload: dict[str, Any] = field(default_factory=dict)
created_at: float = field(default_factory=time.time)
correlation_id: Optional[str] = None
class MessageBus:
"""Central bus shared by the orchestrator and all agents in a run."""
def __init__(self) -> None:
self._inboxes: dict[str, "asyncio.Queue[AgentMessage]"] = {}
self._subscriptions: dict[str, set[str]] = {}
self._history: list[AgentMessage] = []
def register(self, agent_id: str) -> "asyncio.Queue[AgentMessage]":
self._inboxes.setdefault(agent_id, asyncio.Queue())
return self._inboxes[agent_id]
def unregister(self, agent_id: str) -> None:
self._inboxes.pop(agent_id, None)
for subs in self._subscriptions.values():
subs.discard(agent_id)
def subscribe(self, agent_id: str, topic: str) -> None:
self._subscriptions.setdefault(topic, set()).add(agent_id)
self._inboxes.setdefault(agent_id, asyncio.Queue())
async def send(
self,
sender: str,
recipient: str,
type: str,
payload: Optional[dict[str, Any]] = None,
correlation_id: Optional[str] = None,
) -> AgentMessage:
message = AgentMessage(
id=str(uuid.uuid4()), sender=sender, recipient=recipient, topic=None,
type=type, payload=payload or {}, correlation_id=correlation_id,
)
inbox = self._inboxes.setdefault(recipient, asyncio.Queue())
await inbox.put(message)
self._history.append(message)
return message
async def publish(
self, sender: str, topic: str, type: str, payload: Optional[dict[str, Any]] = None
) -> AgentMessage:
message = AgentMessage(
id=str(uuid.uuid4()), sender=sender, recipient=None, topic=topic,
type=type, payload=payload or {},
)
for agent_id in self._subscriptions.get(topic, set()):
inbox = self._inboxes.setdefault(agent_id, asyncio.Queue())
await inbox.put(message)
self._history.append(message)
return message
async def receive(self, agent_id: str, timeout: Optional[float] = None) -> Optional[AgentMessage]:
inbox = self._inboxes.setdefault(agent_id, asyncio.Queue())
try:
if timeout is None:
return await inbox.get()
return await asyncio.wait_for(inbox.get(), timeout=timeout)
except asyncio.TimeoutError:
return None
def history_for(self, agent_id: str, limit: int = 50) -> list[AgentMessage]:
relevant = [m for m in self._history if m.sender == agent_id or m.recipient == agent_id]
return relevant[-limit:]