Chapter 8: Going Async with the Server
Processing everything synchronously was fine for development. But in production, a slow compliance check should not block the deposit response. In this chapter we will configure Redis as the message broker, enable the outbox pattern for reliable delivery, and start the Protean server for asynchronous event processing.
The Outbox Pattern
When an aggregate raises events, they need to reach event handlers and projectors reliably. Protean uses the outbox pattern:
- When the Unit of Work commits, events are written to both the event store and an outbox table atomically.
- The outbox processor reads from the outbox table and publishes events to Redis Streams.
- StreamSubscriptions consume from Redis Streams and dispatch to handlers.
This guarantees at-least-once delivery. Events are never lost, even if Redis is temporarily unavailable.
Configuration
Create a domain.toml file in your project directory:
[brokers.default]
provider = "redis"
url = "${REDIS_URL|redis://localhost:6379/0}"
event_processing = "async"
command_processing = "async"
enable_outbox = true
[event_store]
provider = "memory"
[server]
default_subscription_type = "stream"
messages_per_tick = 100
[server.stream_subscription]
blocking_timeout_ms = 100
max_retries = 3
retry_delay_seconds = 1
enable_dlq = true
Key settings:
brokers.default.provider = "redis": Use Redis as the message broker.event_processing = "async": Events flow through the broker instead of being processed inline.enable_outbox = true: Reliable delivery via the outbox pattern.default_subscription_type = "stream": UseStreamSubscription(Redis Streams with consumer groups) for all handlers.enable_dlq = true: Failed messages go to a dead-letter queue instead of being lost.
Starting Docker Services
You need Redis running:
docker run -d --name fidelis-redis -p 6379:6379 redis:7-alpine
Or use Protean's Docker Compose (if available):
make up
Starting the Server
The Protean server is a long-running process that polls Redis Streams and dispatches messages to handlers:
$ protean server --domain=fidelis
Starting Protean Engine...
Registered subscriptions:
AccountCommandHandler -> fidelis::account:command (StreamSubscription)
AccountSummaryProjector -> fidelis::account (StreamSubscription)
ComplianceAlertHandler -> fidelis::account (StreamSubscription)
NotificationHandler -> fidelis::account (StreamSubscription)
Engine running. Press Ctrl+C to stop.
Each handler gets its own consumer group in Redis. This means:
- Each handler maintains its own read position
- Multiple instances of the same handler can run in parallel (horizontal scaling)
- A slow handler does not block other handlers
- Failed messages are retried automatically before moving to the DLQ
How StreamSubscription Works
┌─────────────┐
│ Event Store │
└──────┬──────┘
│ (events written)
┌──────▼──────┐
│ Outbox │
└──────┬──────┘
│ (outbox processor publishes)
┌──────▼──────┐
│Redis Stream │
│(fidelis:: │
│ account) │
└──┬───┬───┬──┘
│ │ │ (consumer groups)
┌───────────┘ │ └───────────┐
┌──────▼──────┐ ┌──────▼──────┐ ┌──────▼──────┐
│ Projector │ │ Compliance │ │Notification │
│ Consumer │ │ Consumer │ │ Consumer │
│ Group │ │ Group │ │ Group │
└─────────────┘ └─────────────┘ └─────────────┘
Each consumer group reads independently. If the compliance handler is slow, the projector and notification handler continue at full speed.
StreamSubscription vs. EventStoreSubscription
| StreamSubscription | EventStoreSubscription | |
|---|---|---|
| Backed by | Redis Streams | Event store directly |
| Delivery | At-least-once via consumer groups | At-least-once via position tracking |
| DLQ | Built-in | Not yet available |
| Retries | Configurable with backoff | None |
| Use for | Production handlers | Development, projections |
For production systems, StreamSubscription is the recommended choice. It provides consumer groups, automatic retries, dead-letter queues, and horizontal scaling.
Sending Commands Asynchronously
With async processing enabled, domain.process() publishes the command
to the command stream instead of executing it immediately:
# This returns immediately — the command is queued
domain.process(
MakeDeposit(account_id=account_id, amount=500.00, reference="paycheck")
)
The server picks up the command from the Redis stream and dispatches it to the command handler asynchronously.
What We Built
- Redis as the message broker with
domain.tomlconfiguration. - The outbox pattern for reliable event delivery.
- StreamSubscription with consumer groups, retries, and DLQ.
- The Protean server (
protean server) for async processing. - An understanding of how events flow from the aggregate through the outbox to Redis Streams to handlers.
The system is now truly asynchronous. In the next chapter, we will add account-to-account transfers. A multi-aggregate workflow that requires a process manager.
Full Source
"""Chapter 8: Going Async -- The Server
This chapter is about configuration and running the Protean server
for asynchronous event processing. The Python file contains the full
domain setup, while the server is started via the CLI:
protean server --domain fidelis.domain
The domain.toml configuration is shown below as a reference.
"""
from protean import Domain, apply, handle, invariant
from protean.core.projector import on
from protean.exceptions import ValidationError
from protean.fields import DateTime, Float, Identifier, Integer, 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.event(part_of="Account")
class AccountClosed:
account_id: Identifier(required=True)
reason: 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"]}
)
@invariant.post
def closed_account_must_have_zero_balance(self):
if self.status == "CLOSED" and self.balance != 0:
raise ValidationError(
{"status": ["Cannot close account with non-zero balance"]}
)
@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,
)
)
def close(self, reason: str | None = None) -> None:
self.raise_(
AccountClosed(
account_id=str(self.id),
reason=reason,
)
)
@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
@apply
def on_account_closed(self, event: AccountClosed):
self.status = "CLOSED"
@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(part_of=Account)
class CloseAccount:
account_id: Identifier(required=True)
reason: 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)
@handle(CloseAccount)
def handle_close_account(self, command: CloseAccount):
repo = current_domain.repository_for(Account)
account = repo.get(command.account_id)
account.close(reason=command.reason)
repo.add(account)
@domain.event_handler(part_of=Account)
class ComplianceAlertHandler:
@handle(DepositMade)
def on_large_deposit(self, event: DepositMade):
if event.amount >= 10000:
print(
f" [COMPLIANCE] Large deposit alert: "
f"${event.amount:.2f} into account {event.account_id}"
)
@domain.event_handler(part_of=Account)
class NotificationHandler:
@handle(AccountOpened)
def on_account_opened(self, event: AccountOpened):
self.id = event.account_id
print(
f" [NOTIFICATION] Welcome, {event.holder_name}! "
f"Your account {event.account_number} is now active."
)
@handle(WithdrawalMade)
def on_large_withdrawal(self, event: WithdrawalMade):
if event.amount >= 5000:
print(
f" [NOTIFICATION] Large withdrawal alert: "
f"${event.amount:.2f} from account {event.account_id}"
)
@domain.projection
class AccountSummary:
account_id: Identifier(identifier=True, required=True)
account_number: String(max_length=20, required=True)
holder_name: String(max_length=100, required=True)
balance: Float(default=0.0)
transaction_count: Integer(default=0)
last_transaction_at: DateTime()
@domain.projector(projector_for=AccountSummary, aggregates=[Account])
class AccountSummaryProjector:
@on(AccountOpened)
def on_account_opened(self, event: AccountOpened):
self.id = event.account_id
summary = AccountSummary(
account_id=event.account_id,
account_number=event.account_number,
holder_name=event.holder_name,
balance=event.opening_deposit,
transaction_count=1,
last_transaction_at=event._metadata.headers.time,
)
current_domain.repository_for(AccountSummary).add(summary)
@on(DepositMade)
def on_deposit_made(self, event: DepositMade):
repo = current_domain.repository_for(AccountSummary)
summary = repo.get(event.account_id)
summary.balance += event.amount
summary.transaction_count += 1
summary.last_transaction_at = event._metadata.headers.time
repo.add(summary)
@on(WithdrawalMade)
def on_withdrawal_made(self, event: WithdrawalMade):
repo = current_domain.repository_for(AccountSummary)
summary = repo.get(event.account_id)
summary.balance -= event.amount
summary.transaction_count += 1
summary.last_transaction_at = event._metadata.headers.time
repo.add(summary)
# ---------------------------------------------------------------
# domain.toml — place this file alongside your domain module
# ---------------------------------------------------------------
#
# [event_store]
# provider = "message_db"
# database_uri = "postgresql://message_store@localhost:5433/message_store"
#
# [broker]
# provider = "redis"
# redis_url = "redis://localhost:6379/0"
#
# [databases.default]
# provider = "postgresql"
# database_uri = "postgresql://postgres:postgres@localhost:5432/fidelis"
#
# [server]
# # Number of async workers for event/command processing
# workers = 4
#
# ---------------------------------------------------------------
# Start the server with:
#
# protean server --domain fidelis.domain
#
# Start the observatory dashboard with:
#
# protean observatory --domain fidelis.domain
#
# ---------------------------------------------------------------