Skip to content

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:

  1. When the Unit of Work commits, events are written to both the event store and an outbox table atomically.
  2. The outbox processor reads from the outbox table and publishes events to Redis Streams.
  3. 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": Use StreamSubscription (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.toml configuration.
  • 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
#
# ---------------------------------------------------------------

Next

Chapter 9: Transferring Funds →