Skip to content

Event Stores

The Event Store port provides persistence for domain events and commands in event-sourced systems. It serves a dual role: storing the event stream that forms the source of truth for event-sourced aggregates, and acting as the internal messaging backbone within a Protean-based application.

Overview

An event store is fundamentally an append-only log. Events are written to named streams and read back in order. Protean's BaseEventStore interface provides:

  • Stream writes: Append events and commands to named streams
  • Stream reads: Read messages from streams by position
  • Aggregate loading: Reconstitute event-sourced aggregates by replaying events
  • Temporal queries: Load an aggregate at a specific version or point in time
  • Snapshots: Create and restore aggregate snapshots for performance
  • Causation tracing: Traverse causal chains to understand how events triggered other events

Available Event Stores

Memory

The memory event store is the default. It stores events in Python data structures and requires no external services. Ideal for development, testing, and prototyping.

  • No external dependencies
  • All data is lost on process restart
  • Full interface compliance, same API as production event stores

Message DB

Message DB is a PostgreSQL-based event store that provides durable event storage with SQL-based stream operations.

  • Requires: PostgreSQL with the Message DB extension
  • Persistent, durable storage
  • Production-ready with proven reliability

Configuration

Event stores are configured in the [event_store] section of your domain configuration:

# Default: in-memory event store
[event_store]
provider = "memory"

For production, point it at Message DB instead:

[event_store]
provider = "message_db"
database_uri = "${MESSAGE_DB_URL|postgresql://message_store@localhost:5433/message_store}"

${MESSAGE_DB_URL|...} reads the URI from the MESSAGE_DB_URL environment variable, and uses the value after | when it is not set.

Configuration Options

Option Default Description
provider "memory" Event store provider (memory or message_db)
database_uri — Connection string (required for Message DB)

Core Operations

The examples below build on one another. They use a small banking domain with an event-sourced Account.

Writing Events

Events are written to streams by the framework as part of aggregate persistence. You do not typically call the event store directly:

from datetime import UTC, datetime

from protean import Domain, apply, handle
from protean.fields import Float, Identifier

domain = Domain(name="Banking")


@domain.event(part_of="Account")
class Opened:
    account_id: Identifier(required=True)


@domain.event(part_of="Account")
class Deposited:
    account_id: Identifier(required=True)
    amount: Float(required=True)


@domain.aggregate(event_sourced=True)
class Account:
    balance: Float(default=0.0)

    @classmethod
    def open(cls):
        account = cls._create_new()
        account.raise_(Opened(account_id=account.id))
        return account

    @apply
    def opened(self, event: Opened):
        self.id = event.account_id
        self.balance = 0.0

    @apply
    def deposited(self, event: Deposited):
        self.balance += event.amount

    def deposit(self, amount):
        self.raise_(Deposited(account_id=self.id, amount=amount))


@domain.command(part_of=Account)
class Deposit:
    account_id: Identifier(identifier=True)
    amount: Float(required=True)


@domain.command_handler(part_of=Account)
class AccountCommandHandler:
    @handle(Deposit)
    def deposit(self, command: Deposit):
        repo = domain.repository_for(Account)
        account = repo.get(command.account_id)
        account.deposit(command.amount)
        repo.add(account)


domain.init(traverse=False)

with domain.domain_context():
    account = Account.open()
    account.deposit(100.0)
    account.deposit(50.0)

    # Persisting the aggregate writes its three events to the event store
    domain.repository_for(Account).add(account)

When the aggregate is persisted, Protean writes the raised events to the event store automatically.

Reading Streams

The event store adapter is at domain.event_store.store. Each aggregate instance has its own stream, named after the aggregate's stream category and its identifier. For the account above, the stream is banking::account-<id>.

# An aggregate's stream is "<stream category>-<identifier>"
stream = f"{Account.meta_.stream_category}-{account.id}"

with domain.domain_context():
    store = domain.event_store.store

    # Read from the beginning of a stream
    messages = store.read(stream)

    # Read from a specific position
    later = store.read(stream, position=1)

    # Read the last message in a stream
    last = store.read_last_message(stream)

Positions count from 0, so position=1 skips the first event.

Temporal Queries

Load an event-sourced aggregate at a specific version or point in time:

with domain.domain_context():
    # Versions count from 0: version 0 is the state after the first event
    opened = domain.event_store.store.load_aggregate(Account, account.id, at_version=0)

    # Load as of a point in time
    current = domain.event_store.store.load_aggregate(
        Account, account.id, as_of=datetime.now(UTC)
    )

See Temporal Queries for the full guide.

Snapshots

Snapshots cache aggregate state to avoid replaying long event streams:

with domain.domain_context():
    # Create a snapshot for one aggregate
    created = domain.event_store.store.create_snapshot(Account, account.id)

    # Create snapshots for all instances of an aggregate type
    count = domain.event_store.store.create_snapshots(Account)

create_snapshot returns True when it wrote a snapshot. create_snapshots returns the number of aggregates it snapshotted.

Causation Tracing

Trace the causal chain of events to understand how one message led to another. trace_causation and trace_effects take a message id. build_causation_tree takes a correlation id:

with domain.domain_context():
    domain.process(Deposit(account_id=account.id, amount=25.0), asynchronous=False)
    event = domain.event_store.store.read_last_message(stream)

    # Walk up from an event to the command that caused it
    chain = domain.event_store.store.trace_causation(event.metadata.headers.id)

    # Find everything the command caused
    effects = domain.event_store.store.trace_effects(chain[0].metadata.headers.id)

    # Build the full tree for the event's correlation id
    tree = domain.event_store.store.build_causation_tree(
        event.metadata.domain.correlation_id
    )

Here chain holds the Deposit command and then the Deposited event, and effects holds the Deposited event.

See Message Tracing for the full guide.

Monitoring

Use the protean events CLI to inspect event store contents:

# Read from a stream
protean events read "banking::account-<id>" --domain=my_domain

# View aggregate history
protean events history --aggregate=Account --id=<id> --domain=my_domain

# Trace a causal chain
protean events trace "<correlation-id>" --domain=my_domain

See protean events for the full CLI reference.