Subscription Manager¶
SubscriptionManager drives projections from an event store and an event bus. It replays historical events from the store ("catch-up"), transitions each subscription to live events from the bus without gaps, persists checkpoints so restarts resume where they left off, and adds retries, a circuit breaker, a dead-letter queue, health probes, and graceful shutdown around the whole thing.
These recipes assume you already have an event store, an event bus, and a checkpoint repository wired up, and that you want to run one or more projections against them in production. Each section is independent — jump to the task you have:
- Getting a projection running: Wire up a SubscriptionManager, Implement the projection protocol, Choose where a subscription starts
- Surviving failure: Retry transient failures with exponential backoff, Guard external calls with a circuit breaker, Route permanently failed events to the DLQ, React to errors with callbacks
- Operating it: Run multiple projections from one manager, Monitor subscription health, Expose Kubernetes readiness and liveness probes, Shut down gracefully
Import everything from eventsource.application.subscriptions:
from eventsource.application.subscriptions import (
SubscriptionConfig,
SubscriptionManager,
)
manager = SubscriptionManager(event_store, event_bus, checkpoint_repo)
await manager.subscribe(my_projection)
await manager.start()
Before you begin¶
You need Python 3.13 or later and three collaborators to construct a SubscriptionManager: an EventStore (historical events), an EventBus (live events), and a CheckpointRepository (position tracking). A DLQRepository is optional and only needed for the DLQ recipes.
Install the extras for the backends you plan to use — the core package depends only on pydantic and sqlalchemy:
# PostgreSQL store, checkpoints, and DLQ
pip install "eventsource-py[postgresql]"
# SQLite instead
pip install "eventsource-py[sqlite]"
# Redis event bus (also available: rabbitmq, kafka)
pip install "eventsource-py[redis]"
# Everything, including OpenTelemetry tracing
pip install "eventsource-py[all]"
For a database-backed deployment, apply the SQL schema first. The checkpoint recipes need the projection_checkpoints table and the DLQ recipes need dead_letter_queue; both ship in src/eventsource/adapters/sql/schemas/schemas/ (checkpoints.sql, dlq.sql, or all.sql / sqlite_all.sql for the full set).
The SQL repositories are dialect-parameterized -- the same classes serve PostgreSQL and SQLite -- and take a SQLAlchemy AsyncConnection or AsyncEngine:
from sqlalchemy.ext.asyncio import create_async_engine
from eventsource import SQLCheckpointRepository, SQLDLQRepository
engine = create_async_engine("postgresql+asyncpg://localhost/app")
checkpoint_repo = SQLCheckpointRepository(engine)
dlq_repo = SQLDLQRepository(engine)
If you just want to follow along without infrastructure, swap in the in-memory implementations — InMemoryCheckpointRepository() and InMemoryDLQRepository() take no arguments, and InMemoryEventStore / InMemoryEventBus need no services. They lose all state when the process exits, so use them for tests and local exploration only.
Everything below also assumes you are working inside an async context (asyncio.run(...) or an ASGI app), since every store, bus, and subscription API is async.
Wire up a SubscriptionManager¶
Construct the manager with the three required collaborators, register each projection with subscribe(), then call start():
import asyncio
from eventsource.application.subscriptions import SubscriptionConfig, SubscriptionManager
async def main() -> None:
manager = SubscriptionManager(
event_store=event_store,
event_bus=event_bus,
checkpoint_repo=checkpoint_repo,
)
await manager.subscribe(
OrderProjection(),
SubscriptionConfig(start_from="beginning", batch_size=500),
)
results = await manager.start()
failures = {name: err for name, err in results.items() if err is not None}
if failures:
raise RuntimeError(f"subscriptions failed to start: {failures}")
try:
await asyncio.Event().wait() # keep the process alive
finally:
await manager.stop()
asyncio.run(main())
Three things happen when start() runs, per subscription: the manager resolves the starting position (from config.start_from), replays historical events out of the event store in batches, and then transitions to live events from the bus. Checkpoints are written along the way, so a restart picks up where the last one left off.
What subscribe() gives you¶
subscribe(subscriber, config=None, name=None) returns the Subscription object, which is a live handle you can read for monitoring — subscription.name, .state, .last_processed_position, .events_processed, .events_failed, .lag, and .is_running:
subscription = await manager.subscribe(OrderProjection())
print(subscription.name) # "OrderProjection"
print(subscription.state) # SubscriptionState.STARTING
The name defaults to the subscriber's class name. Pass name= when you run two subscriptions from the same projection class, or when the class name is not a stable identifier — the name is the checkpoint key, so changing it makes the subscription replay from its configured start position again:
Names must be unique within a manager; registering the same name twice raises SubscriptionAlreadyExistsError. Omitting config gives you SubscriptionConfig() defaults — resume from checkpoint, batch size 100, retries and circuit breaker on.
Starting and stopping¶
start() returns a dict[str, Exception | None] mapping subscription name to None on success or the exception that stopped it. Subscriptions are isolated: one failing to start does not prevent the others, so always inspect the result rather than assuming success. Calling start() on an already-running manager logs a warning and returns an empty dict.
stop() drains in-flight events and saves checkpoints, with a timeout (default 30 seconds). Both start() and stop() accept subscription_names= to act on a subset:
await manager.start(subscription_names=["orders-read-model"])
await manager.stop(timeout=15.0, subscription_names=["orders-read-model"])
For a long-running daemon, prefer run_until_shutdown() over the manual start() / sleep / stop() dance — see Shut down gracefully.
Constructor options worth setting early¶
Beyond the three required arguments, the constructor takes:
| Argument | Default | Why you would set it |
|---|---|---|
shutdown_timeout |
30.0 |
Total seconds allowed for a graceful shutdown. |
drain_timeout |
10.0 |
Seconds to wait for in-flight events before forcing shutdown. |
dlq_repo |
None |
Enables the dead-letter queue for permanently failed events. |
error_handling_config |
ErrorHandlingConfig() |
Retry, DLQ, and error-classification behavior. |
health_check_config |
HealthCheckConfig() |
Lag and staleness thresholds used by the health probes. |
tracer / enable_tracing |
None / True |
Supply a custom OpenTelemetry tracer, or pass enable_tracing=False to turn spans off. |
A production wiring usually looks like this:
from eventsource.application.subscriptions import (
ErrorHandlingConfig,
HealthCheckConfig,
SubscriptionManager,
)
manager = SubscriptionManager(
event_store=event_store,
event_bus=event_bus,
checkpoint_repo=checkpoint_repo,
dlq_repo=dlq_repo,
shutdown_timeout=60.0,
drain_timeout=20.0,
error_handling_config=ErrorHandlingConfig(),
health_check_config=HealthCheckConfig(),
)
error_handling_config and health_check_config are manager-wide; SubscriptionConfig is per subscription. The rest of this guide fills in what to put in each.
Inspecting what is registered¶
The manager exposes read-only views of its registry, useful in tests and admin endpoints:
manager.subscription_count # int
manager.subscription_names # list[str]
manager.subscriptions # list[Subscription]
manager.get_subscription("orders-read-model") # Subscription | None
manager.get_all_statuses() # dict[str, SubscriptionStatus]
manager.is_running # bool
To remove a subscription at runtime, await manager.unsubscribe(name) — it stops the subscription first if it is running, and returns True if a subscription by that name existed.
Implement the projection protocol¶
Anything you pass to subscribe() must satisfy two members: subscribed_to(), returning the event classes you want, and handle(), processing one event. That is the whole contract — no base class required:
from eventsource.domain.event import DomainEvent
from eventsource.application.subscriptions import Subscriber
class OrderProjection:
def __init__(self, db) -> None:
self._db = db
def subscribed_to(self) -> list[type[DomainEvent]]:
return [OrderCreated, OrderShipped]
async def handle(self, event: DomainEvent) -> None:
if isinstance(event, OrderCreated):
await self._db.insert_order(event.aggregate_id, event.total)
elif isinstance(event, OrderShipped):
await self._db.mark_shipped(event.aggregate_id)
assert isinstance(OrderProjection(db), Subscriber) # runtime-checkable
Subscriber, SyncSubscriber, and BatchSubscriber are @runtime_checkable Protocols, so the isinstance() check above works and is worth asserting in a test — a typo in a method name otherwise surfaces only when the subscription starts.
If you prefer inheritance, BaseSubscriber is an ABC with the same two abstract methods plus can_handle(event) (defaults to type(event) in self.subscribed_to()) and a __repr__ that lists the subscribed types. FilteringSubscriber goes further: override should_handle(event) for predicate filtering and implement _process_event(event) instead of handle().
from eventsource.application.subscriptions import FilteringSubscriber
class TenantOrderProjection(FilteringSubscriber):
def __init__(self, tenant_id: UUID) -> None:
self._tenant_id = tenant_id
def subscribed_to(self) -> list[type[DomainEvent]]:
return [OrderCreated]
def should_handle(self, event: DomainEvent) -> bool:
return event.tenant_id == self._tenant_id
async def _process_event(self, event: DomainEvent) -> None:
await self._db.insert_order(event.aggregate_id, event.total)
Import all of these from eventsource.application.subscriptions. The manager's own type hint is the EventSubscriber ABC from eventsource.ports.handlers, but the runners only ever call subscribed_to() and handle(), so a plain duck-typed class works.
What subscribed_to() controls¶
The returned list does double duty, and the two paths differ:
- Live path: the live runner calls
subscribed_to()once at startup and registers a bus handler per event type. Types missing from the list are never delivered — and adding a type later requires restarting the subscription. - Catch-up path: the runner reads all events from the store and filters them in-process through an
EventFilterbuilt byEventFilter.from_config_and_subscriber(config, subscriber).SubscriptionConfig.event_typeswins if set; otherwise the filter falls back tosubscribed_to(). Returning an empty list means no filtering at all during catch-up, sohandle()sees every event in the store — declare your types explicitly unless that is what you want.
Because config.event_types overrides the subscriber during catch-up but the live subscription always uses subscribed_to(), keep the two consistent. Setting config.event_types to a narrower set than subscribed_to() means the projection sees fewer events while catching up than it does once live.
Handle events one at a time¶
During catch-up, handle() is called once per event in global-position order and awaited before the next event is delivered. The live runner does the same — it reads from the global feed rather than the bus (see ADR 0047), delivering one event at a time after duplicate and filter checks. Design for three things:
Raising is how you signal failure — but nothing retries your handler. RetryableOperation's exponential-backoff retry loop wraps the runner's own I/O: reading batches from the event store and saving checkpoints. Your handle()/handle_batch() call is never retried. An exception out of it records failure metrics and subscription.events_failed / last_error, then either propagates or is logged and swallowed — it is not automatically written to the DLQ. Retrying a flaky database write inside handle() is your job.
The handler circuit breaker observes every outcome, but does not retry either. SubscriptionConfig.circuit_breaker_enabled (default True) builds a handler_circuit_breaker distinct from the infra_circuit_breaker that guards read-batch/checkpoint-save (see circuit_breaker field docs for why they're separate instances). Every handle()/handle_batch() call feeds it — success or failure — regardless of continue_on_error or whether the event later goes to the DLQ. That needs no special case for DLQ'd events: a circuit breaker resets its consecutive-failure count to zero on the next success, so one bad event sandwiched between good ones never opens the circuit. Only a run of consecutive handler failures does — the signal that the handler itself, not one poisoned event, is broken. When open, further handle() calls raise CircuitBreakerOpenError without running your handler at all, until circuit_breaker_recovery_timeout elapses and one probe call is let through.
continue_on_error decides whether one bad event stops the subscription. With the default SubscriptionConfig(continue_on_error=True), both runners log a warning and move on; the catch-up runner then checkpoints past the failed event, so it will never be redelivered. Set continue_on_error=False to re-raise and stop the subscription at the first unhandled exception — the right choice when silently skipping an event would corrupt the read model. This applies equally to a CircuitBreakerOpenError raised while the handler breaker is open: it is just another handler-call exception as far as continue_on_error is concerned.
Handlers must be idempotent. Checkpoints are written per batch by default (CheckpointStrategy.EVERY_BATCH), so a crash redelivers every event processed since the last checkpoint write. Use upserts keyed by event.aggregate_id, or a processed-event-id table, rather than blind inserts. CheckpointStrategy.EVERY_EVENT narrows the redelivery window at the cost of a checkpoint write per event.
SubscriptionConfig.processing_timeout (default 30s) bounds one handler call. Both runners wrap the subscriber call in asyncio.timeout(), so a handle() that hangs no longer blocks its subscription indefinitely — it is abandoned at the deadline and raises TimeoutError, which is an ordinary handler failure from that point on: continue_on_error decides whether the subscription proceeds, the event takes the same DLQ path as any raising handler, and a consecutive run of them opens the handler circuit breaker.
Note that a handle_batch() call is one call, so a whole batch shares a single budget. If you deliver large batches and do slow per-event work, raise processing_timeout to match the batch, rather than assuming it is per event.
To route failures to the dead-letter queue, call the subscription's error handler yourself — see Route permanently failed events to the DLQ.
For an ergonomic alternative to the isinstance ladder, DeclarativeProjection discovers @handles(OrderCreated) methods, generates subscribed_to() from them, and dispatches per event type.
Handle events in batches with handle_batch()¶
When a projection can write in bulk — one multi-row insert instead of N — implement handle_batch(). The BatchSubscriber protocol (@runtime_checkable, import from eventsource.application.subscriptions) requires exactly two members: subscribed_to() and async def handle_batch(self, events: Sequence[DomainEvent]) -> None. Note it does not require handle().
The easier route is BatchAwareSubscriber, an ABC extending BaseSubscriber whose default handle_batch() simply loops over handle(), so you override it only where bulk writes pay off:
from collections.abc import Sequence
from eventsource.application.subscriptions import BatchAwareSubscriber
class AnalyticsProjection(BatchAwareSubscriber):
def subscribed_to(self) -> list[type[DomainEvent]]:
return [OrderCreated, OrderShipped]
async def handle(self, event: DomainEvent) -> None:
await self._db.record_metric(event)
async def handle_batch(self, events: Sequence[DomainEvent]) -> None:
await self._db.bulk_record_metrics(events)
Batch order is the order the events were read from the store, and the sequence may be empty — guard bulk calls that reject empty input.
Two utilities come with it. supports_batch_handling(obj) returns True when an object has a callable handle_batch attribute — a shape check only, it does not inspect the signature. And handle_batch_with_error_tracking(events), a method on BatchAwareSubscriber, processes a batch event by event via handle(), returning (success_count, failures) where failures is a list of (event, exception) pairs, so one bad event does not lose the rest of the batch (each failure is also logged at warning level):
success_count, failures = await projection.handle_batch_with_error_tracking(events)
for event, exc in failures:
logger.error("skipped %s: %s", event.event_id, exc)
Both runners call handle_batch() when your subscriber has one. Batch capability is detected once per subscriber with supports_batch_handling(), and handle_batch() takes precedence over handle() on both paths. Catch-up delivers each read batch as a unit; the live runner delivers each bounded feed read as a unit (live batch delivery). In both cases the whole batch settles — or a handle_batch() that raises falls back to per-event delivery — before the position advances, preserving the one-position-per-subscription invariant (ordered delivery). The live batch is a page the feed already returned, never an accumulator: nothing is held back waiting for more events, so when the feed is quiet a single event is still dispatched immediately, as a batch of one. There is no batching window, timer, or minimum batch size to configure. A subscriber implementing neither method is rejected at runner construction with a TypeError, rather than failing per event once it goes live.
Either way, SubscriptionConfig.batch_size sizes the read batch pulled from the event store; on catch-up that read batch is now also the unit handed to handle_batch(), but on live it still isn't a batch handed to your subscriber. If you implement both methods, keep handle() correct — it is what the live path will call — and note that any subscriber driven outside SubscriptionManager (a backfill script, a rebuild job, a test harness feeding it store.read_all() output) calls whichever method it chooses.
FilteringSubscriber subclasses BatchAwareSubscriber and overrides handle_batch() to apply should_handle() across the sequence before calling _process_event() per surviving event. If you override handle_batch() on a FilteringSubscriber for bulk writes, re-apply the filter yourself — the override replaces that logic.
Choose where a subscription starts¶
SubscriptionConfig.start_from decides the position a subscription resolves to when start() runs. It accepts four things (StartPosition = Literal["beginning", "end", "checkpoint"] | Position), resolved by StartFromResolver:
start_from |
Resolves to | Use it for |
|---|---|---|
"checkpoint" (default) |
The saved checkpoint for this subscription name, or the start of the feed if there is none | Normal long-running projections that must resume after a restart |
"beginning" |
The start of the feed | Rebuilding a read model from the full history, every time it starts |
"end" |
await event_store.current_position() |
Live-only consumers (notifications, cache invalidation) that must not replay history |
Position |
That exact opaque position token | Backfills and repairs from a known point |
from eventsource.application.subscriptions import SubscriptionConfig
await manager.subscribe(OrderProjection(), SubscriptionConfig()) # resume
await manager.subscribe(RebuildProjection(), SubscriptionConfig(start_from="beginning"))
await manager.subscribe(Notifier(), SubscriptionConfig(start_from="end"))
repair_position = await event_store.current_position()
await manager.subscribe(Repair(), SubscriptionConfig(start_from=repair_position), name="repair-from-known-point")
Positions are exclusive and opaque: the catch-up runner reads with from_position=<resolved position>, and the store returns only events strictly after it. start_from="beginning" therefore reads the entire feed, and a Position token delivers events from just after that point onward. Since a checkpoint records the last processed position, resuming never redelivers the checkpointed event itself. Position values are totally ordered only within the store that produced them — comparing or mixing tokens from different stores raises PositionForeignError (see ADR 0024).
SubscriptionConfig no longer accepts a bare int for start_from; pass one of the string literals or a Position obtained from the store or a prior checkpoint.
"checkpoint" when no checkpoint exists¶
The first ever start of a subscription has no row in the checkpoint repository. StartFromResolver logs "No checkpoint found, starting from beginning" and resolves to None — not a position value; positions are opaque adapter-owned tokens with no zero — so a fresh "checkpoint" subscription replays the whole store. That is usually what you want for a projection, but it means a mistyped subscription name silently triggers a full replay under a new checkpoint key instead of resuming the old one. Pin the name explicitly when it matters:
(CheckpointNotFoundError exists in eventsource.ports.exceptions for code that wants to treat a missing checkpoint as fatal; the resolver itself never raises it.)
What "end" actually skips¶
"end" resolves to the store's current max global position, and the transition coordinator's catch-up target is the same watermark — so catch-up finishes with zero events and the subscription goes straight to live bus delivery. Events written between resolution and the live subscription being registered are covered by the transition buffer, but anything written before the process started is never seen. A "end" subscription writes its first checkpoint only once it processes a live event.
Start position is resolved once¶
start_from is read at start() time only. After that the subscription tracks last_processed_position and checkpoints from there; stop() followed by start() re-resolves it. Two consequences:
- A subscription configured with
start_from="beginning"replays the entire store on every process restart, ignoring the checkpoints it wrote. Use it for one-shot rebuilds, not steady state. SubscriptionConfigis a frozen dataclass, so you cannot changestart_fromon a registered subscription.await manager.unsubscribe(name)and re-subscribe()with a new config instead.
To force a replay of a subscription that normally uses "checkpoint", either subscribe under a new name (a new checkpoint key) or clear the existing checkpoint while the manager is stopped:
Prebuilt configs¶
Two factory helpers in eventsource.application.subscriptions cover the common shapes:
from eventsource.application.subscriptions import create_catch_up_config, create_live_only_config
create_catch_up_config(batch_size=1000) # start_from="checkpoint", EVERY_BATCH checkpoints
create_live_only_config() # start_from="end", batch_size=100, EVERY_EVENT
create_catch_up_config(checkpoint_every_batch=False) switches to CheckpointStrategy.PERIODIC instead.