How the Macro Connection Gateway Handles WebSockets and Real-Time Messaging
The Macro connection gateway upgrades HTTP requests to persistent WebSocket connections and broadcasts real-time messages to clients via a Redis Pub/Sub architecture, enabling live document updates and scheduled-action notifications across the platform.
The connection gateway is a critical service in the macro-inc/macro repository that manages persistent client connections and real-time message delivery. Written in Rust using the Axum framework, this service enables live collaboration features by bridging WebSocket clients with backend services through a scalable Pub/Sub infrastructure.
WebSocket Upgrade and Connection Management
The gateway exposes a single HTTP endpoint at /connection-gateway that upgrades incoming requests to WebSocket connections. When a client connects, the handler in services/connection_gateway/src/router.rs calls axum::extract::ws::upgrade to switch the protocol from HTTP to WebSocket.
The upgraded socket is wrapped in a Connection struct that stores the client’s identity—including the user ID and session token—and maintains a sender channel for outbound messages. This design keeps the connection state encapsulated while allowing asynchronous message delivery to the client.
Connection Registry Architecture
All active WebSocket connections are tracked in a thread-safe registry defined in services/connection_gateway/src/registry.rs. The implementation uses a HashMap<ConnectionId, Sender> protected by an RwLock, enabling concurrent reads across multiple tasks while ensuring safe mutation during connection establishment and teardown.
When a new WebSocket is created, the gateway inserts its sender handle into this map using a unique connection identifier. When the socket closes—either through client disconnection or server shutdown—the entry is removed from the registry to prevent memory leaks and message delivery attempts to dead connections.
Redis Pub/Sub Message Broadcasting
Real-time message distribution relies on Redis Pub/Sub channels. Each gateway instance subscribes to a tenant-specific channel named connection-gateway via the subscriber implementation in services/connection_gateway/src/redis_subscriber.rs.
When any service publishes a payload—such as a new comment, a scheduled-action update, or a typing indicator—the gateway receives the message from Redis, looks up all matching target sockets in the registry, and forwards the JSON payload through each socket’s sender channel. This decouples message producers from direct WebSocket management.
Publishing Messages from Services
Downstream services do not communicate directly with Redis. Instead, they use the connection_gateway_client crate located at crates/connection_gateway_client/src/lib.rs. This client provides a thin wrapper around an authenticated HTTP POST call to the gateway’s /publish endpoint.
The ConnectionGatewayClient::publish method marshals the request into a JSON payload and transmits it over HTTP. The gateway validates an internal API secret, then republishes the message onto the Redis channel for broadcasting to connected clients. This abstraction ensures consistent authentication and message formatting across all services.
Real-World Use Cases
Live Comment and Anchor Updates
When a user creates, edits, or deletes a comment or anchor in a document, the document-storage service publishes a ConnectionMessage via the client crate. As implemented in services/document_storage_service/src/api/annotations/create_comment.rs, these events trigger immediate updates to every client that has the document open, ensuring collaborative editing remains synchronized.
Scheduled-Action State Changes
The scheduled-action service sends live status updates—such as "action started" or "action completed"—through the gateway. The implementation in services/scheduled_action/src/outbound/conn_gateway_live_updates.rs demonstrates how long-running background jobs communicate progress to frontend components in real time.
End-to-End WebSocket Testing
The sync-service test suite validates the entire messaging pipeline. Located at services/sync-service/tests/e2e.test.ts, these tests spin up real WebSocket clients, register them with the gateway, and assert that messages published by backend services are received within milliseconds.
Horizontal Scaling and Resilience
The gateway architecture supports horizontal scaling through stateless design principles. The /publish HTTP endpoint is stateless; any gateway instance can handle publication requests because Redis serves as the single source of truth for message distribution.
Multiple gateway instances can run behind a load balancer, with all instances subscribing to the same Redis channels. This guarantees that every connected client receives every broadcast exactly once, regardless of which gateway node manages their specific WebSocket connection.
For graceful shutdown, the Drop implementation for the Connection struct automatically removes the client from the registry and unsubscribes the Redis listener when a gateway pod stops. This prevents dangling references and ensures clean disconnection without message loss.
Code Examples
The following examples demonstrate how to publish messages from a service and how clients connect to receive real-time updates.
// Publishing a message from a backend service
use connection_gateway_client::client::ConnectionGatewayClient;
use connection_gateway_client::message::ConnectionMessage;
let client = ConnectionGatewayClient::new(
std::env::var("INTERNAL_API_SECRET_KEY")?,
std::env::var("CONNECTION_GATEWAY_URL")?,
);
let msg = ConnectionMessage {
tenant_id: "tenant-123".into(),
target: ConnectionMessageTarget::User("user-456".into()),
payload: serde_json::json!({ "type": "comment_created", "comment_id": "c789" }),
};
client.publish(msg).await?;
// Connecting from a web client
const ws = new WebSocket(
`${import.meta.env.VITE_CONNECTION_GATEWAY_WS_URL}/connection-gateway`
);
ws.addEventListener('open', () => {
console.log('WebSocket connection established');
});
ws.addEventListener('message', (event) => {
const data = JSON.parse(event.data);
// Handle real-time updates (e.g., new comment, typing indicator)
console.log('Received:', data);
});
ws.addEventListener('close', () => {
console.log('WebSocket closed');
});
Summary
- The gateway upgrades HTTP connections to WebSockets at
/connection-gatewayusing Axum’s upgrade mechanism inservices/connection_gateway/src/router.rs. - Active connections are stored in a thread-safe
HashMap<ConnectionId, Sender>withinservices/connection_gateway/src/registry.rs. - Messages are broadcast via Redis Pub/Sub through
services/connection_gateway/src/redis_subscriber.rs, enabling horizontal scaling across multiple gateway instances. - Services publish messages using the
connection_gateway_clientcrate rather than connecting directly to Redis, ensuring consistent authentication via the/publishendpoint. - Real-time features—including live document comments and scheduled-action updates—rely on this architecture to deliver sub-second updates to connected clients.
- Graceful shutdown is handled through Rust’s
Droptrait, which cleans up registry entries and Redis subscriptions automatically.
Frequently Asked Questions
How does the connection gateway scale horizontally?
The gateway achieves horizontal scalability by maintaining stateless HTTP endpoints for message publishing and delegating message distribution to Redis. Multiple gateway instances can run behind a load balancer, with each instance subscribing to the same Redis Pub/Sub channels. Since all connection state is local to each instance (stored in the registry) while broadcasts originate from Redis, the system scales linearly with the number of clients and messages.
What Redis channel does the gateway use for broadcasting?
According to the source code in services/connection_gateway/src/redis_subscriber.rs, the gateway subscribes to a channel named connection-gateway for each tenant. All backend services publish messages to this channel, and every active gateway instance receives these messages to forward to their connected WebSocket clients.
How do backend services publish messages without talking directly to Redis?
Services utilize the connection_gateway_client crate defined in crates/connection_gateway_client/src/lib.rs. This client provides a ConnectionGatewayClient::publish method that sends an authenticated HTTP POST request to the gateway’s /publish endpoint. The gateway then validates the internal API secret and republishes the payload to Redis, abstracting the messaging infrastructure from service developers.
How is the WebSocket connection upgraded from HTTP?
The upgrade process occurs in services/connection_gateway/src/router.rs using axum::extract::ws::upgrade. When a client requests /connection-gateway, the Axum server switches the protocol to WebSocket, wraps the socket in a Connection struct containing the user identity and a sender channel, and registers the connection in the thread-safe registry for message delivery.
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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →