Event Bus Architecture in Claude-Code-Telegram: Decoupled Message Routing Explained

The Claude-Code-Telegram repository implements an asynchronous, typed event bus that decouples event producers from consumers through a publish-subscribe pattern with concurrent handler execution and error isolation.

The event bus architecture in RichardAtCT/claude-code-telegram serves as the backbone for routing messages between Telegram updates, webhook endpoints, and agent execution logic. By leveraging Python's asyncio primitives and strict type hierarchies, the system enables clean separation of concerns while maintaining high throughput and fault tolerance.

Core Components of the Event Bus Architecture

Event Base Class and Typed Subclasses

All events in the system inherit from a base Event class defined in src/events/bus.py. This base class provides every event with a unique UUID, UTC timestamp, and source identifier, ensuring traceability across the distributed pipeline.


# From src/events/bus.py

@dataclass
class Event:
    id: UUID = field(default_factory=uuid4)
    timestamp: datetime = field(default_factory=datetime.utcnow)
    source: str = "unknown"

Concrete event types in src/events/types.py extend this base to create strongly-typed payloads for specific domains:

  • UserMessageEvent – Carries Telegram message data including user_id, chat_id, and text
  • WebhookEvent – Encapsulates external webhook payloads with provider and event_type_name fields
  • ScheduledEvent – Represents cron-triggered or timer-based automation
  • AgentResponseEvent – Contains output from Claude API calls ready for delivery back to users

The EventBus Central Hub

The EventBus class in src/events/bus.py acts as the central message broker. It maintains two internal registries: _handlers (a mapping of event types to handler lists) and _global_handlers (handlers receiving all events regardless of type).


# Core registries in EventBus.__init__

self._handlers: Dict[Type[Event], List[EventHandler]] = {}
self._global_handlers: List[EventHandler] = []
self._queue: asyncio.Queue[Event] = asyncio.Queue()

How the Event Bus Facilitates Decoupled Routing

Publish-Subscribe Pattern Implementation

The event bus architecture enables decoupling through explicit subscription APIs. Handlers register for specific event types using subscribe(), or capture all events using subscribe_all(). This eliminates direct dependencies between message producers and consumers.

In src/events/bus.py, the subscription logic appends handlers to the appropriate registry:

def subscribe(self, event_type: Type[Event], handler: EventHandler) -> None:
    """Register a handler for a specific event type."""
    if event_type not in self._handlers:
        self._handlers[event_type] = []
    self._handlers[event_type].append(handler)
    logger.info(f"Handler subscribed to {event_type.__name__}")

Producers remain completely unaware of downstream handlers. They simply instantiate an event and call publish(), which enqueues the event into the internal asyncio.Queue:

async def publish(self, event: Event) -> None:
    """Enqueue an event for processing."""
    await self._queue.put(event)
    logger.debug(f"Event {event.id} published")

Concurrent Handler Execution

The event bus architecture maximizes throughput by executing all matching handlers concurrently. When the background _process_events loop retrieves an event from the queue, it invokes _dispatch(), which collects both type-specific and global handlers, then executes them using asyncio.gather().

From src/events/bus.py, the dispatch logic demonstrates this parallelism:

async def _dispatch(self, event: Event) -> None:
    """Route event to all registered handlers concurrently."""
    handlers = []
    
    # Collect type-specific handlers (including subclass matches)

    for event_type, handler_list in self._handlers.items():
        if isinstance(event, event_type):
            handlers.extend(handler_list)
    
    # Add global handlers

    handlers.extend(self._global_handlers)
    
    # Execute all handlers concurrently

    if handlers:
        await asyncio.gather(*[self._safe_call(h, event) for h in handlers])

This design ensures that a slow handler—such as one performing a long-running Claude API call—cannot block other critical handlers like audit logging or metrics collection.

Error Isolation and Resilience

Decoupled systems require fault tolerance. The event bus architecture isolates handler failures through the _safe_call wrapper, which catches exceptions, logs them via structlog, and prevents propagation to other handlers or the main event loop.

The error isolation implementation in src/events/bus.py:

async def _safe_call(self, handler: EventHandler, event: Event) -> None:
    """Execute handler with error isolation."""
    try:
        await handler(event)
    except Exception as e:
        logger.error(
            "Handler failed for event",
            event_id=str(event.id),
            handler=handler.__name__,
            error=str(e)
        )
        # Exception is not re-raised, preserving event loop stability

This resilience pattern ensures that a buggy handler cannot crash the entire messaging pipeline, maintaining system availability even under partial failure conditions.

Practical Implementation Examples

Subscribing Handlers to Event Types

The following example demonstrates registering specific handlers for Telegram messages and global handlers for monitoring:

from src.events.bus import EventBus
from src.events.types import UserMessageEvent, WebhookEvent

bus = EventBus()

# Specific handler for user messages

async def process_message(event: UserMessageEvent) -> None:
    await claude_client.send_message(event.text, event.chat_id)

# Specific handler for webhooks

async def process_webhook(event: WebhookEvent) -> None:
    if event.provider == "github":
        await handle_github_event(event.payload)

# Global audit handler

async def audit_all_events(event) -> None:
    await database.log_event(event.id, event.source, event.timestamp)

# Register subscriptions

bus.subscribe(UserMessageEvent, process_message)
bus.subscribe(WebhookEvent, process_webhook)
bus.subscribe_all(audit_all_events)

# Start the bus

await bus.start()

Publishing Events from Producers

Producers remain decoupled from handlers, publishing events without knowledge of downstream processing:

from src.events.types import UserMessageEvent
from src.events.bus import bus  # Singleton instance

from pathlib import Path

async def on_telegram_update(update):
    """Called by python-telegram-bot when a message arrives."""
    msg = update.effective_message
    
    # Create typed event

    event = UserMessageEvent(
        user_id=msg.from_user.id,
        chat_id=msg.chat.id,
        text=msg.text or "",
        working_directory=Path.cwd()
    )
    
    # Publish without knowing which handlers will process it

    await bus.publish(event)

Adding New Event Types Without Modifying the Bus

The architecture supports extension through inheritance without touching core routing logic:

from src.events.bus import Event
from dataclasses import dataclass

# Define new event type

@dataclass
class ScheduledBackupEvent(Event):
    backup_type: str
    target_path: str

# Handler automatically eligible for dispatch

async def backup_handler(event: ScheduledBackupEvent):
    await perform_backup(event.target_path)

# Subscribe without modifying EventBus source code

bus.subscribe(ScheduledBackupEvent, backup_handler)

Summary

The event bus architecture in RichardAtCT/claude-code-telegram provides a robust foundation for decoupled message routing through these key mechanisms:

  • Typed Event Hierarchy: All events inherit from a base Event class in src/events/bus.py, enabling type-safe routing and extensibility through subclasses defined in src/events/types.py.
  • Publish-Subscribe Pattern: Producers call publish() without knowing consumers, while handlers register via subscribe() or subscribe_all(), eliminating direct dependencies.
  • Concurrent Dispatch: The _dispatch() method executes all matching handlers using asyncio.gather(), ensuring high throughput and preventing slow operations from blocking the pipeline.
  • Fault Isolation: The _safe_call() wrapper catches and logs exceptions without propagating them, keeping the event loop stable even when individual handlers fail.

This architecture allows the Telegram bot, webhook endpoints, and background schedulers to operate independently while maintaining clean, testable boundaries between components.

Frequently Asked Questions

How does the event bus handle different event types without hardcoding logic?

The event bus uses Python's isinstance() checks within the _dispatch() method to match events against registered handler types. When a handler subscribes to a specific event class, the bus stores that mapping in _handlers. During dispatch, it iterates through all registered types and checks if the incoming event is an instance of that type, including subclass matches. This allows new event types to be added in src/events/types.py without modifying the core routing logic in src/events/bus.py.

What happens if one event handler crashes or raises an exception?

Individual handler failures are isolated through the _safe_call() method, which wraps each handler execution in a try-except block. If a handler raises an exception, the error is logged with the event ID and handler name via structlog, but the exception is not re-raised. This ensures that asyncio.gather() continues executing remaining handlers, and the background _process_events() loop remains stable to process subsequent events from the queue.

Can multiple handlers process the same event simultaneously?

Yes, the event bus architecture supports concurrent execution of all matching handlers. When an event is dispatched, the _dispatch() method collects both type-specific and global handlers into a single list, then executes them concurrently using asyncio.gather(). This design ensures that independent operations—such as sending a Claude API request, writing to an audit log, and updating metrics—can happen in parallel rather than sequentially, maximizing throughput and preventing slow handlers from blocking faster ones.

How do producers remain decoupled from consumers in this architecture?

Producers interact with the event bus exclusively through the publish() method, which accepts an Event instance and places it onto an internal asyncio.Queue. The producer has no knowledge of which handlers are registered, how many exist, or what business logic they implement. Conversely, handlers register interest in specific event types via subscribe() or subscribe_all(), receiving only the events they care about without knowing their origin. This publish-subscribe pattern, implemented in src/events/bus.py, creates a clean boundary that allows components to evolve independently.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →