Event-Driven Architecture of the Macro Notification Service: A 6-Stage Pipeline Deep Dive

The Macro notification service implements a decoupled, event-driven pipeline that separates notification creation from delivery using SQS queues, Postgres LISTEN/NOTIFY events, and horizontally scalable workers to ensure reliable, realtime communication at scale.

The notification service in the macro-inc/macro repository demonstrates production-grade event-driven architecture through a six-stage pipeline that isolates concerns between producers, brokers, and consumers. Written in Rust, the system leverages AWS SQS for at-least-once delivery guarantees alongside Postgres database triggers for realtime updates, allowing services like scheduled_action and teams to dispatch notifications without blocking on delivery logic.

Stage 1: Ingress Layer – Producing Notification Events

Services initiate notifications through the NotificationIngress trait, specifically the SqsNotificationIngress implementation defined in crates/notification/src/outbound/sqs_notification_ingress.rs. The ingress acts as the system's entry point, persisting notification metadata to the notification table before enqueueing a reference to the SQS queue.

The producer constructs requests using the SendNotificationRequestBuilder from notification::domain::models, specifying the target entity, sender, and notification metadata.

use notification::domain::models::SendNotificationRequestBuilder;
use notification::domain::service::NotificationIngress;
use std::sync::Arc;

let request = SendNotificationRequestBuilder::new()
    .notification_entity(EntityType::User.with_entity_string(user_id.to_string()))
    .sender_id(Some(actor_id.clone()))
    .notification(metadata)
    .build();

notification_ingress.send_notification(request).await?;

Stage 2: SQS Queue – The Event Broker

The SqsQueue implementation in crates/notification/src/outbound/queue.rs serves as the message broker, decoupling ingress from processing. It guarantees at-least-once delivery by accepting payloads from producers and making them available to worker consumers with configurable polling parameters (max_messages, wait_time_seconds).

Stage 3: Worker Pool – Processing Queued Events

The PushNotificationEventWorker in crates/notification/src/inbound/push_notification_event_worker.rs continuously polls the SQS queue, converting queued events into concrete delivery actions. Each worker instance resolves target device tokens, applies rate limiting through the RedisRateLimitAdapter, and invokes the appropriate outbound adapters.

use notification::inbound::push_notification_event_worker::PushNotificationEventWorker;

let worker = PushNotificationEventWorker::new(
    db.clone(),
    redis.clone(),
    queue.clone(),
    mobile_push_adapter,
    email_adapter,
);

while let Some(event) = queue.receive().await? {
    worker.handle(event).await?;
}

Stage 4: Realtime Propagation – Database Event Streaming

Parallel to the main delivery pipeline, the NotificationEventsListener in crates/notification/src/inbound/notification_events_listener.rs subscribes to Postgres LISTEN/NOTIFY channels on the user_notification table. When rows are inserted, updated, or deleted, the listener deserializes the payload into NotificationDatabaseEvent structures and broadcasts NotificationStatusPayload updates through realtime publishers.

use notification::inbound::notification_events_listener::NotificationEventsListener;
use notification::outbound::notification_events::PgNotificationEventsReceiver;
use notification::outbound::websocket::WebSocketGatewayAdapter;

let receiver = PgNotificationEventsReceiver::new(db.clone());
let realtime = WebSocketGatewayAdapter::new(gateway_client);
let mut listener = NotificationEventsListener::new(receiver, realtime);

listener.run().await; // Never returns

This architecture enables instantaneous UI updates without requiring clients to poll the database.

Stage 5: Delivery Adapters – Transport Implementation

The final delivery stage utilizes specialized adapters located in crates/notification/src/outbound/. The MobilePushAdapter handles APNs and FCM requests in mobile.rs, while the EmailAdapter manages digest aggregation in email.rs. Realtime publishers include KafkaRealtimeSender, FanoutRealtimeSender, and WebSocketGatewayAdapter, each implemented in kafka_notification_realtime.rs, fanout_notification_realtime.rs, and websocket.rs respectively.

Stage 6: Reliability Controls – Rate Limiting and Deduplication

Before emitting mobile pushes, workers consult the RedisRateLimitAdapter defined in crates/macro_cache_client/src/notification_rate_limit.rs to prevent spam. The SnsEndpointManagerAdapter maintains stable device token mappings, ensuring consistent delivery ordering and deduplication across retries.

Summary

  • Decoupled Architecture: The NotificationIngress trait isolates producers from delivery mechanics, allowing independent scaling of ingress and worker components.
  • Dual Event Streams: The system combines SQS queued events for reliable delivery with Postgres NOTIFY events for realtime UI synchronization.
  • Horizontal Scalability: Both SqsQueue consumers and NotificationEventsListener instances can scale horizontally without affecting the main service entry point in services/notification_service/src/main.rs.
  • Comprehensive Rate Limiting: Redis-backed rate limiting prevents notification flooding while maintaining per-user delivery guarantees.

Frequently Asked Questions

How does the Macro notification service handle realtime updates without polling?

The service uses a database event listener pattern where NotificationEventsListener subscribes to Postgres LISTEN/NOTIFY channels. When the user_notification table changes, triggers fire NOTIFY events that the listener captures, deserializes into NotificationDatabaseEvent, and forwards to WebSocket or Kafka realtime publishers. This eliminates the need for client-side polling while keeping UI state synchronized with the database.

What is the difference between the SQS queue and the Postgres event stream?

The SQS queue (crates/notification/src/outbound/queue.rs) handles the primary delivery workflow—persisting and processing push notifications, emails, and digests with at-least-once guarantees. The Postgres event stream (crates/notification/src/inbound/notification_events_listener.rs) handles status propagation—broadcasting creation, read, and deletion events to connected clients for immediate UI state synchronization.

How does the system prevent duplicate notification delivery?

Workers implement deduplication through the SnsEndpointManagerAdapter for device token stability and the RedisRateLimitAdapter for frequency control. Additionally, SQS visibility timeout mechanisms and the idempotent design of PushNotificationEventWorker ensure that retrying failed deliveries does not result in duplicate user notifications.

Where is the entry point for the notification service?

The service boots from services/notification_service/src/main.rs, which wires together the SqsNotificationIngress, queue configurations, worker pools, realtime publishers, and the NotificationEventsListener. This file initializes all adapters and starts the event-driven pipeline.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →