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

> Explore the event-driven architecture of Macro's notification service. Discover its 6-stage pipeline, SQS queues, and Postgres LISTEN/NOTIFY for reliable, realtime communication at scale.

- Repository: [Macro/macro](https://github.com/macro-inc/macro)
- Tags: deep-dive
- Published: 2026-08-17

---

**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`](https://github.com/macro-inc/macro/blob/main/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.

```rust
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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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.

```rust
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`](https://github.com/macro-inc/macro/blob/main/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.

```rust
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`](https://github.com/macro-inc/macro/blob/main/mobile.rs), while the `EmailAdapter` manages digest aggregation in [`email.rs`](https://github.com/macro-inc/macro/blob/main/email.rs). Realtime publishers include `KafkaRealtimeSender`, `FanoutRealtimeSender`, and `WebSocketGatewayAdapter`, each implemented in [`kafka_notification_realtime.rs`](https://github.com/macro-inc/macro/blob/main/kafka_notification_realtime.rs), [`fanout_notification_realtime.rs`](https://github.com/macro-inc/macro/blob/main/fanout_notification_realtime.rs), and [`websocket.rs`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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.