Skip to content

Chapter 21: The Event Store as a Database

We have been working with events through aggregates, projections, and handlers for the entire tutorial. But we have never looked at the event store directly. In this chapter we will explore it as a database, reading raw events, viewing statistics, searching by type, and understanding the stream naming conventions.

Reading Events from a Stream

View events for a specific account:

$ protean events read "fidelis::account-acc-001" --domain=fidelis
 Position  Global Pos  Type              Time                Data Keys
 0         1           AccountOpened      2025-01-10 09:15   account_id, account_number, holder_name, opening_deposit
 1         5           DepositMade        2025-01-10 10:30   account_id, amount, source_type, reference
 2         12          DepositMade        2025-02-01 14:00   account_id, amount, source_type, reference
 3         18          WithdrawalMade     2025-03-01 09:00   account_id, amount, reference

Showing 4 event(s) from position 0

Add --data to see full payloads:

$ protean events read "fidelis::account-acc-001" --data --domain=fidelis

Use --from and --limit for pagination:

$ protean events read "fidelis::account-acc-001" --from=10 --limit=5 --domain=fidelis

Reading Category Streams

Omit the instance ID to read across all accounts:

$ protean events read "fidelis::account" --limit=10 --domain=fidelis

This returns events from all account instances, ordered by global_position.

Stream Naming Conventions

Stream Pattern Example
Instance stream {domain}::{category}-{id} fidelis::account-acc-001
Category stream {domain}::{category} fidelis::account
Command stream {domain}::{category}:command-{id} fidelis::account:command-acc-001
Snapshot stream {domain}::{category}:snapshot-{id} fidelis::account:snapshot-acc-001
Fact event stream {domain}::{category}-fact-{id} fidelis::account-fact-acc-001

The stream category is derived from the aggregate's class name (lowercased, underscored).

Domain-Wide Statistics

$ protean events stats --domain=fidelis
 Aggregate   Stream Category      ES?  Instances  Events  Latest Type         Latest Time
 Account     fidelis::account     Yes     1,247  245,891  DepositMade         2025-06-16 15:30
 Transfer    fidelis::transfer    Yes       312    1,248  TransferCompleted   2025-06-16 14:55

Total: 247,139 event(s) across 1,559 aggregate instance(s)

This gives you a high-level view of the entire event store: how many aggregates, how many events, and the latest activity.

Searching by Event Type

Find all events of a specific type:

$ protean events search --type=DepositMade --domain=fidelis
 Position  Global Pos  Type         Stream                        Time
 1         5           DepositMade  fidelis::account-acc-001      2025-01-10 10:30
 2         12          DepositMade  fidelis::account-acc-001      2025-02-01 14:00
 ...

Found 89,234 event(s) matching type 'DepositMade' (showing first 20)

Searches support partial matching and are case-insensitive:

$ protean events search --type=deposit --domain=fidelis

Event Store Positions

Two position numbers appear in event listings:

  • Position: The event's index within its specific stream (0-indexed). This is the aggregate's version number.
  • Global Position: A monotonically increasing counter across the entire event store. This establishes a total ordering of all events, regardless of which aggregate they belong to.

Global position is critical for projection rebuilding (events must be replayed in global order) and for subscription position tracking.

Memory vs. MessageDB

Throughout this tutorial we used the memory event store, great for development and testing, but not persistent across restarts.

For production, use MessageDB. A PostgreSQL-based event store:

# domain.toml (production)
[event_store]
provider = "message_db"
database_uri = "${MESSAGE_DB_URL|postgresql://message_store@localhost:5432/message_store}"

MessageDB provides:

  • Persistent storage with PostgreSQL durability
  • Optimistic concurrency control
  • Efficient category reads and stream queries
  • SQL access for ad-hoc analysis

The domain code does not change, only the configuration.

What We Built

  • protean events read for reading raw events from streams.
  • protean events stats for domain-wide statistics.
  • protean events search for finding events by type.
  • Understanding of stream naming conventions.
  • Understanding of position vs. global position.
  • Memory vs. MessageDB event store adapters.

The event store is not just an implementation detail. It is a database of facts about your business. Learning to query it directly is a useful debugging and analysis tool.

Full Source

"""Chapter 21: Event Store as a Database

Demonstrates the domain setup for using the event store as the primary source
of truth.  Protean provides CLI commands to inspect, query, and manage event
store data directly:

    protean events list --domain fidelis.domain --stream <stream-name>
    protean events read --domain fidelis.domain --stream <stream-name>
    protean events stats --domain fidelis.domain

This chapter is entirely about CLI workflows for inspecting the event store;
the Python source only needs to set up the domain so the CLI can operate on it.
"""

from protean import Domain, apply, handle, invariant
from protean.exceptions import ValidationError
from protean.fields import Float, Identifier, String
from protean.utils.globals import current_domain

domain = Domain("fidelis")


@domain.event(part_of="Account")
class AccountOpened:
    account_id: Identifier(required=True)
    account_number: String(required=True)
    holder_name: String(required=True)
    opening_deposit: Float(required=True)


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


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


@domain.aggregate(event_sourced=True)
class Account:
    account_number: String(max_length=20, required=True)
    holder_name: String(max_length=100, required=True)
    balance: Float(default=0.0)
    status: String(max_length=20, default="ACTIVE")

    @invariant.post
    def balance_must_not_be_negative(self):
        if self.balance is not None and self.balance < 0:
            raise ValidationError(
                {"balance": ["Insufficient funds: balance cannot be negative"]}
            )

    @classmethod
    def open(cls, account_number: str, holder_name: str, opening_deposit: float):
        account = cls._create_new()
        account.raise_(
            AccountOpened(
                account_id=str(account.id),
                account_number=account_number,
                holder_name=holder_name,
                opening_deposit=opening_deposit,
            )
        )
        return account

    def deposit(self, amount: float, reference: str | None = None) -> None:
        if amount <= 0:
            raise ValidationError({"amount": ["Deposit amount must be positive"]})
        self.raise_(
            DepositMade(
                account_id=str(self.id),
                amount=amount,
                reference=reference,
            )
        )

    def withdraw(self, amount: float, reference: str | None = None) -> None:
        if amount <= 0:
            raise ValidationError({"amount": ["Withdrawal amount must be positive"]})
        self.raise_(
            WithdrawalMade(
                account_id=str(self.id),
                amount=amount,
                reference=reference,
            )
        )

    @apply
    def on_account_opened(self, event: AccountOpened):
        self.id = event.account_id
        self.account_number = event.account_number
        self.holder_name = event.holder_name
        self.balance = event.opening_deposit
        self.status = "ACTIVE"

    @apply
    def on_deposit_made(self, event: DepositMade):
        self.balance += event.amount

    @apply
    def on_withdrawal_made(self, event: WithdrawalMade):
        self.balance -= event.amount


@domain.command(part_of=Account)
class OpenAccount:
    account_number: String(required=True)
    holder_name: String(required=True)
    opening_deposit: Float(required=True)


@domain.command(part_of=Account)
class MakeDeposit:
    account_id: Identifier(required=True)
    amount: Float(required=True)
    reference: String()


@domain.command(part_of=Account)
class MakeWithdrawal:
    account_id: Identifier(required=True)
    amount: Float(required=True)
    reference: String()


@domain.command_handler(part_of=Account)
class AccountCommandHandler:
    @handle(OpenAccount)
    def handle_open_account(self, command: OpenAccount):
        account = Account.open(
            account_number=command.account_number,
            holder_name=command.holder_name,
            opening_deposit=command.opening_deposit,
        )
        current_domain.repository_for(Account).add(account)
        return str(account.id)

    @handle(MakeDeposit)
    def handle_make_deposit(self, command: MakeDeposit):
        repo = current_domain.repository_for(Account)
        account = repo.get(command.account_id)
        account.deposit(command.amount, reference=command.reference)
        repo.add(account)

    @handle(MakeWithdrawal)
    def handle_make_withdrawal(self, command: MakeWithdrawal):
        repo = current_domain.repository_for(Account)
        account = repo.get(command.account_id)
        account.withdraw(command.amount, reference=command.reference)
        repo.add(account)

Next

Chapter 22: The Full Picture →