Implementing Event-Driven Asynchronous Communication in Agent Systems: A Complete Guide to In-Process Message Buses
Event-driven asynchronous communication in agent systems can be implemented using a lightweight, in-process message bus built on Python's asyncio primitives, enabling decoupled coordination without external infrastructure like Redis.
Modern multi-agent applications require robust communication patterns that preserve isolation while maintaining non-blocking performance. The bojieli/ai-agent-book repository demonstrates how to build such systems using purely in-process, event-driven architectures. This guide examines two production-ready implementations: a minimal publish-subscribe bus for parallel research tasks and a full-featured bus with trace persistence and data redaction for sensitive workflows.
Core Architecture: The Message Bus Pattern
The Envelope as the Atomic Unit
Every message traversing the bus is wrapped in an Envelope (in chapter10/parallel-web-research/message_bus.py) or AgentMessage (in chapter10/autonomous-phone-registration/bus.py). These dataclasses encapsulate:
- sender: Originating agent identifier
- target: Recipient ID or
BROADCAST = "*"for fan-out - type: Message type identifier for routing logic
- payload: Arbitrary data dictionary
- sequence number: Monotonically increasing counter for deterministic ordering
- timestamp: Human-readable and monotonic timing for debugging
The sequence counter and start time are initialized as module-level globals (lines 31-35 in message_bus.py), ensuring consistent ordering across all messages in a session.
Subscription Model
Consumers express interest via the Subscription class, which holds an asyncio.Queue for incoming envelopes. Registration follows this pattern:
# Subscribe to specific message types
sub = bus.subscribe(owner_id="agent_B", types=["task_assigned", "status_update"])
# Or subscribe to all message types
sub = bus.subscribe(owner_id="agent_B")
The Subscription.get() method provides a blocking-but-async interface that yields control when no messages are available, preventing event loop starvation.
Two Production Implementations
Minimal Bus: Parallel Web Research
The MessageBus in chapter10/parallel-web-research/message_bus.py provides core functionality with minimal overhead:
- Broadcast suppression: The sender never receives its own broadcast messages
- Type-based filtering: Subscriptions declare accepted types via predicate matching
- Verbose logging: Optional timestamped console output for development
The dispatch logic in MessageBus.publish iterates registered subscriptions and delivers matching envelopes:
# Agent A publishes a task assignment
await bus.send(
sender_id="agent_A",
target="agent_B", # point-to-point delivery
type="task_assigned",
payload={"url": "https://example.com"},
)
# Agent B consumes in an event loop
sub = bus.subscribe(owner_id="agent_B")
while True:
env = await sub.get() # yields control when empty
if env.type == "task_assigned":
# Process and respond
await bus.send(
sender_id="agent_B",
target="agent_A",
type="status_update",
payload={"status": "done"},
)
Full-Featured Bus: Autonomous Phone Registration
The MessageBus in chapter10/autonomous-phone-registration/bus.py extends the pattern with operational necessities:
Trace persistence: Every message is append-written to a JSONL file for post-mortem analysis
Payload redaction: Sensitive keys are masked in traces while preserved for runtime consumers
Direct receive API: Simpler bus.receive(recipient) method for single-consumer scenarios
# Initialize with trace capture
bus = MessageBus(trace_path="traces/run01.json")
# Send with automatic redaction
await bus.send(
sender="frontend",
recipient="phone_agent",
type="register",
sensitive_keys=("phone_number", "ssn"),
phone_number="+1-555-1234", # appears as "<redacted>" in trace
username="alice",
)
# Receiver gets unredacted payload
msg = await bus.receive("phone_agent")
print(msg.payload["phone_number"]) # "+1-555-1234"
The send method in this implementation handles redaction before appending to self.history and writing to disk, ensuring compliance without compromising runtime functionality.
Event Flow and Guarantees
The message bus implements exactly-once delivery semantics within process boundaries:
- Registration phase: Agent obtains Subscription with isolated Queue
- Publishing phase:
bus.send()constructs envelope, assigns sequence number, appends to history - Dispatch phase:
publish()routes to all matching subscriptions (all for broadcast, specific ID for point-to-point) - Consumption phase: Agent awaits queue, processes message, potentially publishes responses
This flow guarantees that:
- Broadcast messages reach all matching subscribers except the sender
- Point-to-point messages route to exactly one recipient
- Message ordering is consistent across all observers via global sequence counter
- No external dependencies (Redis, Kafka, RabbitMQ) are required
Alternative: Priority-Based Event Bus
The repository also contains chapter9/gaia-experience/AWorld/aworld/core/event/event_bus.py, which implements priority ordering via asyncio.PriorityQueue. This variant suits scenarios where message urgency varies:
- High priority: Critical alerts, shutdown signals
- Normal priority: Standard task assignments
- Low priority: Telemetry, logging
Priority inversion is prevented by the queue's heap-based ordering, with tie-breaking on sequence number for fairness.
Performance Characteristics
| Characteristic | Parallel-Web-Research | Autonomous-Phone-Registration |
|---|---|---|
| Memory footprint | Minimal (in-memory queues only) | Moderate (trace file buffering) |
| Latency | Sub-microsecond (direct method calls) | Sub-millisecond (with disk I/O) |
| Throughput | Thousands of messages/second | Hundreds of messages/second (I/O bound) |
| Persistence | None | JSONL trace files |
| Security | Basic | Payload redaction |
Both implementations scale to hundreds of agents within a single Python process, limited primarily by the GIL around queue operations rather than architectural constraints.
Summary
- In-process event buses eliminate infrastructure dependencies while preserving async semantics through
asyncio.Queueprimitives - Envelope/AgentMessage dataclasses provide standardized, timestamped, sequenced message containers
- Subscription objects isolate consumers with private queues, enabling safe blocking without starving the event loop
- Broadcast suppression prevents self-delivery loops in fan-out scenarios
- Trace persistence and redaction balance observability with data protection requirements
- Sequence counters provide deterministic, debuggable message ordering across all participants
Frequently Asked Questions
What is the difference between message_bus.py and bus.py in the repository?
The message_bus.py implementation (chapter10/parallel-web-research/) provides a minimal, high-performance bus optimized for educational clarity and low overhead. The bus.py implementation (chapter10/autonomous-phone-registration/) adds production features including JSONL trace persistence, sensitive key redaction, and a simplified receive() API. Both share core semantics but target different operational requirements.
How does the message bus handle backpressure when consumers are slower than producers?
The asyncio.Queue underlying each Subscription has an unbounded default in these implementations, so producers never block on put(). For memory-constrained deployments, you can subclass Subscription to pass maxsize=N to the Queue constructor, causing publish() to await space availability. The repository's GAIA event bus variant (chapter9/gaia-experience/) demonstrates this pattern with bounded queues.
Can multiple processes or machines use this message bus?
No—these implementations are strictly in-process, relying on Python's single-event-loop architecture. For cross-process communication, the repository patterns would need extension via Unix domain sockets, HTTP webhooks, or integration with external brokers. The wrapped dataclasses (Envelope, AgentMessage) serialize cleanly to JSON for such extensions.
Why use sequence numbers instead of native asyncio queue ordering?
Native queue ordering preserves insertion sequence, but sequence numbers provide visibility and determinism across system restarts when combined with trace files. They also enable merging multiple trace streams during post-mortem analysis, where wall-clock timestamps may collide or drift. The monotonic counter in message_bus.py (lines 31-35) initializes once at import time, ensuring consistent ordering for the process lifetime.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →