Skip to content

Chapter 18: Monitoring Subscription Health

After the DLQ incident, the team realizes they need proactive monitoring. They should know when a handler is falling behind or accumulating failures, not discover it from a customer complaint.

Checking Subscription Status

The protean subscriptions status command gives a dashboard of all subscriptions:

$ protean subscriptions status --domain=fidelis
                    Subscriptions - Fidelis
 Handler                     Type    Stream                Lag  Pending  DLQ  Status
 AccountCommandHandler       stream  fidelis::account:cmd    0        0    -  ok
 AccountSummaryProjector     stream  fidelis::account        2        0    -  ok
 ComplianceAlertHandler      stream  fidelis::account        0        0    3  lagging
 NotificationHandler         stream  fidelis::account        0        0    -  ok
 AccountReportProjector      stream  fidelis::account-fact   0        0    -  ok
 FundsTransferPM             stream  fidelis::transfer       0        0    -  ok

6 subscription(s), 5 ok, 1 lagging, total lag: 2

Key metrics:

  • Lag: How many messages the handler has not yet processed. High lag means the handler is falling behind.
  • Pending: Messages currently being processed (claimed but not acknowledged).
  • DLQ: Number of messages in the dead-letter queue.
  • Status: ok, lagging, or unknown.

For machine-readable output:

$ protean subscriptions status --domain=fidelis --json

The Observatory

For real-time monitoring, launch the Observatory, Protean's built-in observability dashboard:

$ protean observatory --domain=fidelis --port=9000
Observatory running at http://0.0.0.0:9000

The Observatory provides:

  • Live message traces: A real-time stream of handler.started, handler.completed, handler.failed, message.acked, message.dlq events via Server-Sent Events.
  • Subscription status: The same data as the CLI, auto-refreshing every 5 seconds.
  • DLQ management: Inspect, replay, and purge directly from the web interface.
  • Stream health: Queue depths per stream.

Observatory API Endpoints

Endpoint Description
GET / Dashboard
GET /stream SSE real-time trace stream
GET /api/health Health check
GET /api/subscriptions Subscription status
GET /api/traces Recent trace history
GET /api/streams Stream information
GET /api/outbox Outbox status
GET /metrics Prometheus metrics

Prometheus Metrics

The Observatory exposes Prometheus-compatible metrics at /metrics:

# HELP protean_subscription_lag Messages behind head position
# TYPE protean_subscription_lag gauge
protean_subscription_lag{domain="fidelis",handler="AccountSummaryProjector",stream="fidelis::account",type="stream"} 2

# HELP protean_subscription_dlq_depth Messages in dead-letter queue
# TYPE protean_subscription_dlq_depth gauge
protean_subscription_dlq_depth{domain="fidelis",handler="ComplianceAlertHandler",stream="fidelis::account",type="stream"} 3

# HELP protean_subscription_status Subscription status (1=ok, 0=error)
# TYPE protean_subscription_status gauge
protean_subscription_status{domain="fidelis",handler="AccountCommandHandler",stream="fidelis::account:command",type="stream"} 1

These metrics can be scraped by Grafana, Datadog, or any Prometheus- compatible monitoring tool. Set alerts on:

  • protean_subscription_lag > 100: Handler falling behind
  • protean_subscription_dlq_depth > 0: Failed messages accumulating

Trace Events

The engine emits structured trace events for every message processed:

  • handler.started: Handler began processing a message
  • handler.completed: Handler finished successfully
  • handler.failed: Handler threw an exception
  • message.acked: Message acknowledged (removed from pending)
  • message.nacked: Message negatively acknowledged (will retry)
  • message.dlq: Message moved to dead-letter queue
  • outbox.published: Outbox processor published a message to broker
  • outbox.failed: Outbox processor failed to publish

These events flow to the Observatory via Redis Pub/Sub in real-time. When nobody is listening, the emitter short-circuits with zero overhead.

What We Built

  • protean subscriptions status for quick health checks.
  • The Observatory for real-time monitoring and DLQ management.
  • Prometheus metrics for production alerting.
  • Understanding of trace events emitted by the engine.

With monitoring in place, we can detect problems early. In the next chapter, a bank acquisition triggers a massive migration, and we learn how to handle it without disrupting production.

Full Source

"""Chapter 18: Monitoring Health

Demonstrates the domain setup for health monitoring with the Observatory
dashboard and Prometheus metrics.  Protean provides CLI commands to inspect
subscription lag, handler throughput, and system health:

    protean observatory --domain fidelis.domain
    protean server --domain fidelis.domain --prometheus-port 9090

This chapter is primarily about operational tooling and CLI workflows; the
Python source only needs to set up the domain so the tools can operate on it.
"""

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.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)


@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.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)

Next

Chapter 19: The Great Migration, Run on Priority Lanes →