How Macro's Email Service Integrates with Gmail API and Synchronizes Emails: A Production-Grade Rust Implementation
Macro's email service uses a Pub/Sub-driven pipeline with Redis-cached authentication, SQS message processing, and cursor-based history synchronization to maintain bidirectional Gmail integration.
The Macro email service, implemented in the macro-inc/macro repository, provides robust, near-real-time synchronization between user inboxes and Macro's PostgreSQL data store. This architecture handles authentication, change notifications, and conflict resolution through a carefully designed multi-stage pipeline.
Authentication and Public Key Management
Before any Gmail API calls, the service verifies JWTs using Google's public keys. These keys are cached in Redis to minimize latency and external dependencies.
In services/email_service/src/util/gmail/auth.rs, the get_google_public_keys function implements a cache-first strategy:
// services/email_service/src/util/gmail/auth.rs
pub async fn get_google_public_keys(
redis_client: Arc<RedisClient>,
gmail_client: Arc<GmailClient>,
) -> anyhow::Result<KeyMap> {
let cached = redis_client.get_google_public_keys().await.ok().flatten();
match cached {
Some(keys) => Ok(keys),
None => fetch_and_cache_google_public_keys(redis_client, gmail_client).await,
}
}
On a cache miss, fetch_and_cache_google_public_keys calls the Gmail endpoint, stores the result in Redis, and returns the key map. This pattern ensures that verification remains fast even under load.
Inbox Change Notification Pipeline
Gmail pushes Inbox Sync notifications to an Amazon SQS queue. The process_message function in services/email_service/src/pubsub/inbox_sync/process.rs serves as the entry point:
// services/email_service/src/pubsub/inbox_sync/process.rs
pub async fn process_message(
ctx: PubSubContext,
message: &aws_sdk_sqs::types::Message,
) -> Result<()> {
// Deserialize the SQS payload
let data = extract_inbox_sync_message(message)?;
// Delegate to a typed handler based on the operation enum
inner_process_message(&ctx, &data).await?;
cleanup_message(&ctx.sqs_worker, message).await?;
Ok(())
}
The message is parsed into an InboxSyncPubsubMessage. The InboxSyncOperation enum determines which handler executes:
gmail_message— Handles new message notifications and history changesupsert_message— Creates or updates message records in PostgreSQLdelete_message— Removes messages from Macro's storeupdate_labels— Synchronizes label assignments
Gmail History Synchronization Logic
The core synchronization logic resides in services/email_service/src/pubsub/inbox_sync/operations/gmail_message.rs. The gmail_message handler implements cursor-based incremental sync:
// services/email_service/src/pubsub/inbox_sync/operations/gmail_message.rs
pub async fn gmail_message(
ctx: &PubSubContext,
link: &Link,
payload: &GmailMessagePayload,
) -> Result<(), ProcessingError> {
// 1️⃣ Load the stored Gmail history cursor
let db_history_id = email_db_client::histories::fetch_history_id_for_link(...).await?;
// 2️⃣ Compare with the notification's history_id
if db_history_u64 >= payload.history_id { return Ok(()); }
// 3️⃣ Pull the list of new labels (cached in Redis) and sync them locally
let labels = ctx.email_api.list_labels(link.id).await?;
sync_labels(&ctx.db, link.id, &labels).await?;
// 4️⃣ Query Gmail for the delta of changes since the last cursor
let change_batch = ctx.email_api
.list_changes(link.id, &SyncCursor::gmail(&db_history_id))
.await
.map_err(handle_gmail_message_error)?;
// 5️⃣ Persist the new cursor
email_db_client::histories::upsert_gmail_history(&ctx.db, link.id,
change_batch.next_cursor.as_str()).await?;
// 6️⃣ Turn each change into a downstream Pub/Sub message
for ps_message in build_pubsub_messages(link.id, change_batch.changes) {
ctx.sqs_client.enqueue_gmail_inbox_sync_notification(ps_message).await?;
}
Ok(())
}
This handler performs six critical operations:
- Retrieves the last known Gmail history cursor from the database via
fetch_history_id_for_link - Deduplicates notifications by comparing history IDs — early returns prevent redundant processing
- Synchronizes labels using
email_api.list_labelswith Redis caching - Fetches incremental changes through
email_api.list_changes, which wraps Gmail'shistory.listendpoint - Updates the stored cursor with
upsert_gmail_historyfor the next sync cycle - Emits downstream messages via
build_pubsub_messagesandenqueue_gmail_inbox_sync_notification
Stale Cursor Recovery and Backfill
When Gmail returns an OutdatedCursor error, the service cannot proceed with incremental sync. The schedule_stale_cursor_backfill function creates a recovery job that:
- Re-processes the entire mailbox range
- Patches the stale cursor once backfill completes via
repair_stale_cursor
This guarantees eventual consistency even when Gmail's sync cursor expires or becomes invalid.
Sending Emails with Proper Threading
Outbound messages require correct threading headers for Gmail compatibility. The generate_email_threading_headers function in services/email_service/src/util/gmail/send.rs constructs In-Reply-To and References headers:
let (in_reply_to, references) = generate_email_threading_headers(
&db_pool,
Some(parent_message_id),
link.id,
).await;
email_api.send_message(
SendEmailRequest {
to: vec![recipient],
subject: "Re: …".into(),
in_reply_to,
references,
body: "Your reply".into(),
attachments: None,
},
).await?;
For messages with attachments, the service supports:
- Draft attachments — Fetched from S3 via
fetch_and_attach_draft_attachments - Forwarded attachments — Pulled directly from Gmail via
fetch_and_attach_forwarded_attachments
Key Implementation Files
| File | Purpose |
|---|---|
services/email_service/src/util/gmail/auth.rs |
Google public key retrieval with Redis caching |
services/email_service/src/util/gmail/send.rs |
Threading headers, attachment handling, message sending |
services/email_service/src/pubsub/inbox_sync/process.rs |
SQS message entry point and routing |
services/email_service/src/pubsub/inbox_sync/operations/gmail_message.rs |
Core sync logic: cursor management, change fetching, backfill scheduling |
services/email_service/src/pubsub/inbox_sync/operations/upsert_message.rs |
Database writes for new and updated messages |
services/email_service/src/pubsub/inbox_sync/operations/delete_message.rs |
Message deletion propagation |
services/email_service/src/pubsub/inbox_sync/operations/update_labels.rs |
Label synchronization |
services/email_service/src/outbound/email_api.rs |
GmailApi wrapper exposing list_labels, list_changes, send_message |
Summary
- Cache-first authentication — Google public keys are cached in Redis to minimize Gmail API calls for JWT verification
- SQS-driven processing — Gmail push notifications flow through
process_messageininbox_sync/process.rs - Cursor-based incremental sync — The
gmail_messagehandler inoperations/gmail_message.rsuses Gmail's history API for efficient delta synchronization - Automatic recovery — Stale cursors trigger backfill jobs to maintain data consistency
- Full threading support — Outbound replies include proper RFC 5322 headers for Gmail compatibility
Frequently Asked Questions
How does Macro handle Gmail authentication securely?
Macro caches Google public keys in Redis using the get_google_public_keys function in services/email_service/src/util/gmail/auth.rs. On a cache hit, verification uses the stored keys; on a miss, the service fetches fresh keys from Gmail, caches them, and proceeds. This eliminates redundant external calls while maintaining security.
What triggers email synchronization in Macro?
Gmail sends Inbox Sync push notifications to an Amazon SQS queue. The process_message function receives these notifications, validates them, and routes them to appropriate handlers based on the InboxSyncOperation enum — including gmail_message, upsert_message, delete_message, and update_labels.
How does Macro recover from expired or invalid Gmail sync cursors?
When email_api.list_changes returns an OutdatedCursor error, the schedule_stale_cursor_backfill function creates a recovery job that re-processes the full mailbox range. Once complete, repair_stale_cursor updates the stored cursor. This ensures eventual consistency without manual intervention.
What happens when a user sends a reply through Macro?
The service calls generate_email_threading_headers to build proper In-Reply-To and References headers, optionally attaches files from S3 or Gmail using fetch_and_attach_draft_attachments or fetch_and_attach_forwarded_attachments, and dispatches the message through GmailApi::send_message.
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 →