Spaces:
Running
Running
| from abc import ABC, abstractmethod | |
| from dataclasses import dataclass, field | |
| from datetime import datetime | |
| from typing import Any, Dict, Optional | |
| from utils.helpers import generate_id | |
| class Message: | |
| """Inter-agent message envelope.""" | |
| sender: str | |
| recipient: str | |
| msg_type: str # e.g. "task_complete", "error", "data" | |
| payload: Dict[str, Any] = field(default_factory=dict) | |
| msg_id: str = field(default_factory=lambda: generate_id("msg")) | |
| timestamp: str = field(default_factory=lambda: datetime.utcnow().isoformat()) | |
| job_id: Optional[str] = None | |
| class BaseProtocol(ABC): | |
| """Abstract base for all inter-agent communication protocols.""" | |
| async def send(self, message: Message) -> None: | |
| """Publish a message.""" | |
| async def receive(self, recipient: str, timeout: float = 30.0) -> Optional[Message]: | |
| """Receive the next message for a given recipient.""" | |
| async def broadcast(self, message: Message, recipients: list) -> None: | |
| """Broadcast a message to multiple recipients.""" | |