How Macro Delivers Real‑Time Notifications via WebSocket: Notification Service Architecture Explained

Macro streams notifications to users in real time by combining a Kafka consumer with a WebSocket gateway adapter that batches messages through an HTTP API to the Connection Gateway service.

The Macro platform's notification service achieves low‑latency delivery by decoupling message ingestion from transport. Incoming notifications flow through Apache Kafka, get processed by a dedicated consumer, and reach active users via WebSocket connections maintained by a separate gateway service. This article breaks down the complete data flow using source code from the macro-inc/macro repository.

The Three‑Stage WebSocket Delivery Pipeline

The real‑time notification mechanism follows a clear producer‑adapter‑gateway pattern:

1. Kafka Consumption with NotificationTopicConsumer

The NotificationTopicConsumer in crates/notification/src/outbound/notification_consumer.rs listens to the macro.notifications topic. It runs as an ungrouped consumer starting at the latest offset, ensuring only new notifications trigger WebSocket pushes.

When a WebSocketDeliveryRequested event arrives, the consumer deserializes each UserNotificationRow into the target notification type T and yields a NotificationTopicEvent::WebSocketDeliveryRequested value. This design isolates the transport concern from the domain logic—the consumer knows nothing about WebSockets, only that a real‑time send was requested.

2. RealtimeSender Port Implementation

The domain layer defines the RealtimeSender trait in crates/notification/src/domain/ports.rs. The WebSocketGatewayAdapter<W> in crates/notification/src/outbound/websocket.rs implements this port for any type W satisfying WebSocketGatewayOps.

In production, the concrete type is ConnectionGatewayClient. This adapter translates domain calls into gateway operations, keeping the notification service agnostic to the underlying transport protocol.

3. Connection Gateway Batch Delivery

The ConnectionGatewayClient performs the actual WebSocket broadcast via HTTP. Its batch_send_to_entities method constructs this JSON payload:

{
  "message_type": "notification",
  "message": <serialized notification>,
  "entities": [
    { "type": "user", "id": "user-abc-123" },
    { "type": "user", "id": "user-def-456" }
  ]
}

This POSTs to {gateway_url}/message/batch_send. The Connection Gateway—running as a separate service—upgrades the request to WebSocket frames and broadcasts to all active connections for those users.

The gateway responds with MessageReceipt objects listing which users actually received the message. The client filters for delivery_count > 0 and returns the successful set:

// From crates/notification/src/outbound/websocket.rs
impl RealtimeSender for WebSocketGatewayAdapter<ConnectionGatewayClient> {
    async fn send_to_users<T: Serialize>(
        &self,
        user_ids: &[MacroUserIdStr],
        payload: &T,
    ) -> Result<Vec<MacroUserIdStr>, NotificationError> {
        let receipts = self.gateway.batch_send_to_entities(
            "notification",
            payload,
            &user_ids.iter().map(|id| EntityRef::user(id.clone())).collect(),
        ).await?;

        Ok(receipts
            .into_iter()
            .filter(|r| r.delivery_count > 0)
            .map(|r| r.user_id)
            .collect())
    }
}

End‑to‑End Code Example: Sending Notifications

Here is how a backend service produces and routes a notification through the full pipeline:

use notification::outbound::websocket::{
    ConnectionGatewayClient, WebSocketGatewayAdapter,
};
use notification::domain::ports::RealtimeSender;
use serde::Serialize;

#[derive(Serialize)]
struct NewMessageNotification {
    title: String,
    body: String,
    document_id: String,
}

// Initialize the gateway client with internal service credentials
let gateway_client = ConnectionGatewayClient::new(
    std::env::var("INTERNAL_AUTH_KEY").expect("INTERNAL_AUTH_KEY not set"),
    std::env::var("CONNECTION_GATEWAY_URL").expect("GATEWAY_URL not set"),
);

// Wrap in the RealtimeSender adapter
let realtime_sender = WebSocketGatewayAdapter {
    gateway: gateway_client,
};

// Define target recipients
let recipients = vec![
    MacroUserIdStr::parse_from_str("user-550e8400").unwrap(),
    MacroUserIdStr::parse_from_str("user-6ba7b810").unwrap(),
];

// Dispatch—returns only users who were actively connected
let delivered_to = realtime_sender
    .send_to_users(&recipients, &NewMessageNotification {
        title: "Comment received".into(),
        body: "Jane replied to your annotation".into(),
        document_id: "doc-0194f8a2".into(),
    })
    .await
    .expect("notification delivery failed");

println!("Real-time delivery confirmed for {} users", delivered_to.len());

Code Example: Consuming Kafka Events

The notification service itself runs this consumer loop to bridge Kafka to WebSocket:

use notification::outbound::notification_consumer::{
    NotificationTopicConsumer, NotificationTopicEvent,
};
use notification::outbound::websocket::{
    ConnectionGatewayClient, WebSocketGatewayAdapter,
};
use notification::domain::ports::NotificationRealtimePublisher;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // Kafka consumer reads from macro.notifications
    let consumer = NotificationTopicConsumer::<SystemNotification>::from_env(
        &std::env::var("KAFKA_BROKERS")?,
    )?;

    // Same adapter implements NotificationRealtimePublisher
    let publisher = WebSocketGatewayAdapter {
        gateway: ConnectionGatewayClient::new(
            std::env::var("INTERNAL_AUTH_KEY")?,
            std::env::var("CONNECTION_GATEWAY_URL")?,
        ),
    };

    loop {
        match consumer.recv().await? {
            NotificationTopicEvent::WebSocketDeliveryRequested(metadata) => {
                let user_ids = metadata.notifications
                    .iter()
                    .map(|n| n.user_id.clone())
                    .collect::<Vec<_>>();

                match publisher.send_to_users(&user_ids, &metadata).await {
                    Ok(delivered) => tracing::info!(
                        "WebSocket delivered to {}/{} users",
                        delivered.len(),
                        user_ids.len()
                    ),
                    Err(e) => tracing::error!("WebSocket send failed: {}", e),
                }
            }
            _ => {} // Handle other event variants (email, push, etc.)
        }
    }
}

Why This Architecture Scales

  • Kafka as buffer: Notification spikes get absorbed by the topic, preventing backpressure on the WebSocket gateway.
  • Ungrouped consumer: Each notification service instance reads independently—no partition rebalancing delays.
  • Latest offset start: Avoids replaying historical notifications on restart, keeping startup fast.
  • Batch HTTP API: Reduces connection churn versus per-message WebSocket calls to the gateway.
  • Delivery receipts: Enables accurate metrics and retry logic for failed pushes.

Summary

  • NotificationTopicConsumer in notification_consumer.rs ingests Kafka events for real‑time delivery.
  • WebSocketGatewayAdapter implements the RealtimeSender port to hide transport details.
  • ConnectionGatewayClient batches notifications via HTTP POST /message/batch_send to the Connection Gateway.
  • The Connection Gateway upgrades HTTP to WebSocket frames and broadcasts to active user connections.
  • Delivery receipts filter the result set to only users who were actually online and received the message.

Frequently Asked Questions

How does Macro guarantee message order for WebSocket notifications?

The Kafka topic macro.notifications provides ordered delivery per partition. The NotificationTopicConsumer processes events sequentially from its assigned partitions, preserving order for any single user's notifications as long as they route to the same partition key (typically user_id).

What happens when a user is offline?

The ConnectionGatewayClient returns delivery_count: 0 for users without active WebSocket connections. The notification service logs this and may trigger alternative delivery paths—email or push notifications—via separate Kafka topics and consumers, though these fallbacks exist outside the real‑time WebSocket flow.

Can the notification service handle millions of concurrent connections?

Not directly. The WebSocket connections themselves live in the separate Connection Gateway service, which can scale horizontally independent of the notification service. The notification service only maintains HTTP connections to the gateway, not to end users. This separation allows each component to scale according to its specific load pattern.

Why use HTTP rather than native WebSocket from the notification service?

The batch HTTP API reduces complexity and connection state. The notification service avoids managing thousands of persistent WebSocket connections; it simply fires-and-forgets to the gateway, which specializes in connection management. This also enables the gateway to implement features like connection coalescing, geographic routing, and protocol translation without changes to the notification service.

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 →