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 readfor reading raw events from streams.protean events statsfor domain-wide statistics.protean events searchfor 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)