Skip to main content

Event Bus Architecture

This page provides a comprehensive deep dive into OmniDaemon’s event bus architecture, focusing on the Redis Streams implementation - the current production-ready backend.

What is the Event Bus?

The event bus is the message broker that delivers events from publishers to subscribers (agents). Think of it as the “nervous system” of your AI agent infrastructure. Key Responsibilities:
  • 📨 Receive events from publishers
  • 🚀 Deliver events to appropriate subscribers
  • 💾 Persist messages for durability
  • 🔄 Handle retries when processing fails
  • ⚖️ Load balance across multiple agent instances
  • 💀 Manage DLQ for failed messages
  • 📊 Track metrics for observability

Redis Streams Implementation

Overview

Redis Streams is a data structure that acts as an append-only log with powerful consumption semantics. It’s perfect for event-driven architectures. Why Redis Streams?
  • Durable - Messages persisted to disk
  • Reliable - At-least-once delivery guaranteed
  • Fast - Sub-millisecond latency
  • Scalable - Handles millions of messages/second
  • Consumer Groups - Built-in load balancing
  • Reclaiming - Automatic failure recovery
  • Simple - No separate broker to manage

Architecture Diagram


Core Components

1. Redis Stream

A Redis Stream is an append-only log where each entry has: Structure:
Message ID Format:
Stream Naming:
Configuration:
Default Values:
  • default_maxlen: 10,000 messages (approximate)
  • reclaim_interval: 30 seconds
  • default_reclaim_idle_ms: 180,000 ms (3 minutes)
  • default_dlq_retry_limit: 3 retries

2. Consumer Groups

A consumer group is a set of consumers that process messages from the same stream. Redis ensures each message is delivered to only ONE consumer in the group. Creation:
Group Naming Convention:
Why Groups?
  • Load Balancing - Messages distributed across consumers
  • Fault Tolerance - If consumer dies, another picks up
  • At-Least-Once - Messages not lost even if consumer crashes
  • Position Tracking - Group remembers last processed message
Multiple Consumers in Group:

3. Message Publishing

Publishing Process:
  1. Event Created:
  2. Serialized to JSON:
  3. Added to Stream:
XADD Parameters:
  • maxlen: Maximum stream length (older messages trimmed)
  • approximate: If True, trimming is faster but less precise (~)
Return Value:

4. Message Consumption

Consumption Loop:
XREADGROUP Parameters:
  • groupname: Consumer group name
  • consumername: Individual consumer name (unique within group)
  • streams: Dict of {stream_name: start_id}
    • ">" = Only NEW messages (not pending)
    • "0" = All messages from beginning
  • count: Max messages per read (batch size)
  • block: Milliseconds to block waiting (0 = don’t block)
Return Format:

5. Pending Entries List (PEL)

The Pending Entries List tracks messages that have been delivered to consumers but not yet acknowledged. Why PEL?
  • Failure Recovery - Know which messages weren’t processed
  • Reclaiming - Reassign stuck messages to other consumers
  • Monitoring - See what’s in-flight
PEL Entry:
Checking PEL:

6. Message Reclaiming

If a consumer crashes or takes too long, its messages are reclaimed by another consumer. Reclaim Loop:
Reclaim Configuration:
  • reclaim_idle_ms: How long before reclaiming (default: 180,000 ms = 3 min)
  • reclaim_interval: How often to check (default: 30 seconds)
  • dlq_retry_limit: Max retries before DLQ (default: 3)
In-Flight Tracking:

7. Dead Letter Queue (DLQ)

Messages that fail repeatedly go to the Dead Letter Queue for manual inspection. DLQ Stream:
DLQ Entry Structure:
Sending to DLQ:
Inspecting DLQ:

8. Acknowledgment (XACK)

Acknowledging a message tells Redis it was successfully processed and can be removed from the PEL. When to Acknowledge:
XACK Behavior:
  • ✅ Removes message from PEL
  • ✅ Message no longer redelivered
  • ✅ Group’s “last delivered ID” advances

9. Metrics Tracking

All events are tracked in a dedicated metrics stream. Metrics Stream:
Metric Events:
Emitting Metrics:

Redis Commands Used

XADD - Add Message to Stream

Parameters:
  • Stream name
  • * = Auto-generate message ID
  • Field-value pairs
Returns: Message ID (e.g., 1678901234567-0)

XREADGROUP - Read from Group

Parameters:
  • GROUP group consumer
  • COUNT - Max messages
  • BLOCK - Timeout (ms)
  • STREAMS - Stream and start ID (> = new only)
Returns: List of messages

XACK - Acknowledge Message

Parameters:
  • Stream name
  • Group name
  • Message ID(s)
Returns: Count of acknowledged messages

XPENDING - Get Pending Messages

Parameters:
  • Stream name
  • Group name
  • Start ID (- = beginning)
  • End ID (+ = end)
  • Count
Returns: List of pending messages with idle time

XCLAIM - Reclaim Message

Parameters:
  • Stream name
  • Group name
  • New consumer name
  • Min idle time (ms)
  • Message ID(s)
Returns: Claimed messages

XGROUP CREATE - Create Consumer Group

Parameters:
  • CREATE stream group start-id
  • MKSTREAM - Create stream if doesn’t exist
Returns: OK

XGROUP DESTROY - Delete Consumer Group

Returns: OK

XINFO GROUPS - List Consumer Groups

Returns: List of groups with stats

Performance Characteristics

Throughput

Single Redis Instance:
  • Publishing: ~100,000 messages/second
  • Consuming: ~50,000 messages/second per consumer
  • With Persistence: ~50,000 messages/second
Redis Cluster:
  • Publishing: >1,000,000 messages/second
  • Consuming: Scales linearly with consumers

Latency

End-to-End Latency:
  • Publish to Deliver: <10ms (typical)
  • Publish to Process: <50ms (typical)
  • With Network: <100ms (typical)
Breakdown:
  • XADD: <1ms
  • XREADGROUP: <1ms
  • Network: 1-10ms
  • Callback execution: Variable (depends on your agent)

Memory Usage

Per Message:
  • ~1 KB average (depends on payload size)
Stream Overhead:
  • Minimal (<10% of message size)
Pending Entries:
  • Stored in memory
  • ~100 bytes per pending message
Example:

Persistence

AOF (Append-Only File):
  • Every write logged to disk
  • Slower but safer (no data loss)
  • Recommended for production
RDB (Snapshot):
  • Periodic snapshots
  • Faster but risk of data loss
  • OK for development
Configuration:

Monitoring Event Bus

Via CLI

Via SDK

Via Redis CLI


Advanced Configuration

Custom Consumer Count

Effect:
  • 5 consumers in the group
  • Load distributed across all 5
  • Higher throughput

Custom Reclaim Settings

Use Case: Quick reclaim for fast-failing tasks

Custom Retry Limit

Use Case: Transient errors (network issues, rate limits)

Custom Stream Max Length

Use Case: High-volume topics needing more history

Best Practices

1. Choose Appropriate reclaim_idle_ms

2. Set Appropriate max_retries

3. Monitor DLQ

4. Configure Persistence

5. Scale with Consumer Count


Troubleshooting

Messages Not Being Delivered

Messages Stuck in Pending

High DLQ Count

Slow Processing


Further Reading


Summary

Key Components:
  • Redis Stream - Append-only log for messages
  • Consumer Group - Load balancing and fault tolerance
  • Pending Entries List (PEL) - Track unacknowledged messages
  • Message Reclaiming - Automatic failure recovery
  • Dead Letter Queue (DLQ) - Failed message storage
Key Redis Commands:
  • XADD - Publish message
  • XREADGROUP - Consume message
  • XACK - Acknowledge message
  • XCLAIM - Reclaim message
  • XPENDING - Check pending messages
Configuration:
  • default_maxlen: 10,000 messages
  • reclaim_interval: 30 seconds
  • default_reclaim_idle_ms: 180,000 ms (3 min)
  • default_dlq_retry_limit: 3 retries
Performance:
  • Throughput: ~100K msgs/sec (single instance)
  • Latency: <10ms (typical)
  • Memory: ~1KB per message
Redis Streams provides a robust, scalable foundation for OmniDaemon’s event-driven architecture! 🚀