Skip to content

Advanced Messaging Patterns

Overview

This tutorial covers advanced messaging patterns available in the FCC event bus system, including priority queues, dead-letter queues, circuit breakers, and event replay strategies.

Difficulty: Advanced Prerequisites: Familiarity with the FCC event bus (see Event Bus tutorial). Time: 45 minutes

Step 1: Basic Event Bus Review

from fcc.messaging.bus import EventBus
from fcc.messaging.events import Event, EventType

bus = EventBus()
received = []

def handler(event: Event) -> None:
    received.append(event)

bus.subscribe(EventType.SIMULATION_STARTED, handler)
bus.publish(Event(event_type=EventType.SIMULATION_STARTED))
print(f"Events received: {len(received)}")

Step 2: Priority Queue Pattern

Implement priority-based event processing:

import heapq
from dataclasses import dataclass, field

@dataclass(order=True)
class PrioritizedEvent:
    priority: int
    event: Event = field(compare=False)

queue = []
heapq.heappush(queue, PrioritizedEvent(2, Event(event_type=EventType.SIMULATION_STARTED)))
heapq.heappush(queue, PrioritizedEvent(0, Event(event_type=EventType.COMPLIANCE_CHECK_COMPLETED)))
heapq.heappush(queue, PrioritizedEvent(1, Event(event_type=EventType.WORKFLOW_COMPLETED)))

while queue:
    item = heapq.heappop(queue)
    print(f"Priority {item.priority}: {item.event.event_type.value}")

Step 3: Dead-Letter Queue Pattern

Capture failed events for retry:

class SimpleDeadLetterQueue:
    def __init__(self):
        self._failed = []

    def add(self, event, error):
        self._failed.append((event, str(error)))

    def retry(self, bus):
        for event, _ in self._failed:
            bus.publish(event)
        count = len(self._failed)
        self._failed.clear()
        return count

Step 4: Circuit Breaker Pattern

Prevent cascading failures:

class SimpleCircuitBreaker:
    def __init__(self, threshold=3):
        self._failures = 0
        self._threshold = threshold
        self._open = False

    def call(self, func, *args):
        if self._open:
            raise RuntimeError("Circuit open")
        try:
            return func(*args)
        except Exception:
            self._failures += 1
            if self._failures >= self._threshold:
                self._open = True
            raise

    def reset(self):
        self._failures = 0
        self._open = False

Step 5: Event Replay

Replay serialized events for debugging:

from fcc.messaging.serialization import EventSerializer

serializer = EventSerializer()
event = Event(event_type=EventType.SIMULATION_STARTED, data={"scenario": "test"})

# Serialize
json_str = serializer.serialize(event)
print(f"Serialized: {json_str[:80]}...")

# Deserialize and replay
restored = serializer.deserialize(json_str)
bus.publish(restored)
print(f"Replayed event type: {restored.event_type.value}")

Next Steps

  • Read Guidebook Chapter 23 for in-depth coverage.
  • Explore the ComplianceSubscriber for a real-world subscriber example.
  • Build a custom subscriber with DLQ and circuit breaker integration.