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 includinguser_id,chat_id, andtextWebhookEvent– Encapsulates external webhook payloads withproviderandevent_type_namefieldsScheduledEvent– Represents cron-triggered or timer-based automationAgentResponseEvent– 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
Eventclass insrc/events/bus.py, enabling type-safe routing and extensibility through subclasses defined insrc/events/types.py. - Publish-Subscribe Pattern: Producers call
publish()without knowing consumers, while handlers register viasubscribe()orsubscribe_all(), eliminating direct dependencies. - Concurrent Dispatch: The
_dispatch()method executes all matching handlers usingasyncio.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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →