How Macro’s CRM Automatically Aggregates Emails for Company and Contact Objects
Macro's CRM aggregates email activity by processing every inbound and outbound message through a Pub/Sub backfill pipeline that computes MIN/MAX interaction timestamps per contact and propagates those aggregates to linked company records.
Email-based CRMs live or die by their ability to surface relationship history without manual data entry. In the macro-inc/macro repository, the CRM achieves this through a multi-stage email aggregation pipeline that transforms raw mailbox data into queryable contact and company activity metrics. This article examines the source code implementation, from message ingestion to aggregate persistence.
Overview of the Email Aggregation Pipeline
The aggregation system spans three architectural layers:
- Ingestion layer:
services/email_servicehandles message storage and job enqueueing - Backfill workers: Pub/Sub tasks process batches of messages asynchronously
- Data models:
crates/models_soupandcrates/entity_accessdefine schemas and query patterns
The pipeline triggers on every email event—real-time synchronization or historical backfill—and maintains first interaction and last interaction timestamps for both contacts and their associated companies.
Message Ingestion and Job Enqueueing
When a message arrives, the upsert_message handler in services/email_service/src/pubsub/inbox_sync/operations/upsert_message.rs extracts participant data and constructs CrmContactRecipient structures for each unique email address in the message headers.
Immediately after message persistence, the system calls enqueue_populate_crm_contacts from services/email_service/src/pubsub/util.rs:
// Simplified representation of the enqueue flow
pub async fn enqueue_populate_crm_contacts(
message_id: &str,
owner_email: &str,
recipients: Vec<CrmContactRecipient>,
) -> Result<(), Error> {
// Creates Pub/Sub task for backfill workers
let task = PopulateCrmContactTask {
message_id: message_id.to_string(),
user_email: owner_email.to_string(),
recipients,
};
pubsub.publish("backfill.crm_contact", task).await
}
This decoupled design ensures that message ingestion remains fast while CRM aggregation happens asynchronously.
Per-Contact Aggregation in populate_crm_contact
The core aggregation logic resides in services/email_service/src/pubsub/backfill/populate_crm_contact.rs. This worker receives the message context and performs three critical operations:
- Contact resolution: Maps each email address to an existing
crm_contactrecord or creates a new one - Timestamp aggregation: Computes
MIN(internal_date_ts)andMAX(internal_date_ts)across all emails involving this contact - UPSERT execution: Persists the aggregated values to
first_interaction_tsandlast_interaction_tscolumns
The aggregation query follows this pattern:
// Conceptual UPSERT from populate_crm_contact.rs
INSERT INTO crm_contacts (
id, email, name, company_id,
first_interaction_ts, last_interaction_ts,
updated_at
)
VALUES ($1, $2, $3, $4, $5, $6, NOW())
ON CONFLICT (email, user_id) DO UPDATE SET
first_interaction_ts = LEAST(
crm_contacts.first_interaction_ts,
EXCLUDED.first_interaction_ts
),
last_interaction_ts = GREATEST(
crm_contacts.last_interaction_ts,
EXCLUDED.last_interaction_ts
),
name = COALESCE(EXCLUDED.name, crm_contacts.name),
company_id = COALESCE(EXCLUDED.company_id, crm_contacts.company_id),
updated_at = NOW();
The LEAST/GREATEST functions guarantee monotonic aggregation—earlier timestamps never overwrite later ones, and vice versa.
User-Level Seeding with populate_crm_for_user
Historical mailbox backfills require batch processing. The populate_crm_for_user worker in services/email_service/src/pubsub/backfill/populate_crm_for_user.rs ensures complete coverage:
// From populate_crm_for_user.rs
pub async fn populate_crm_for_user(user_email: &str) -> Result<(), Error> {
let participants = fetch_all_participants_for_user(user_email).await?;
for batch in participants.chunks(100) {
enqueue_populate_crm_contacts(
"backfill", // synthetic message_id for batch operations
user_email,
batch.to_vec(),
).await?;
}
Ok(())
}
This seeds the aggregation queue with every distinct sender and recipient ever encountered in the user's mailbox, eliminating gaps in CRM data.
Company-Level Aggregation
Contact aggregation feeds company-level metrics through foreign key relationships. The crm_company model—defined in crates/models_soup/src/crm_company.rs and crates/models_search/src/crm_company.rs—links multiple contacts to a single company record.
The crm_company_access query layer in crates/entity_access/src/outbound/pg_access_repo/queries/crm_company_access.rs rolls up contact aggregates:
-- Company aggregation derived from linked contacts
SELECT
c.id as company_id,
MIN(cc.first_interaction_ts) as company_first_contact,
MAX(cc.last_interaction_ts) as company_last_contact,
COUNT(DISTINCT cc.id) as associated_contacts
FROM crm_companies c
JOIN crm_contacts cc ON cc.company_id = c.id
WHERE c.user_id = $1
GROUP BY c.id;
This query enables the CRM to display "First contacted" and "Last activity" at the company level without storing redundant timestamps.
Database Schema for Aggregates
The migration crates/macro_db_client/migrations/20260520150000_crm_contacts_name.up.sql establishes the aggregation columns:
CREATE TABLE IF NOT EXISTS crm_contacts (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
user_id UUID NOT NULL REFERENCES users(id),
email TEXT NOT NULL,
name TEXT,
company_id UUID REFERENCES crm_companies(id),
-- Aggregation fields
first_interaction_ts BIGINT, -- MIN of all email timestamps
last_interaction_ts BIGINT, -- MAX of all email timestamps
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
UNIQUE(user_id, email)
);
CREATE INDEX idx_crm_contacts_company_id
ON crm_contacts(company_id)
WHERE company_id IS NOT NULL;
The BIGINT timestamp storage (Unix epoch milliseconds) ensures timezone-agnostic comparisons and efficient indexing.
Backfill Orchestration
The crates/models_email/src/email/service/backfill.rs module coordinates the entire pipeline. It exports high-level functions that wire together:
process_message()→ triggersenqueue_populate_crm_contactsfor real-time eventsbackfill_user()→ invokespopulate_crm_for_userfor historical importsdepopulate_message()→ reverses aggregates when messages are deleted
The backfill service in services/email_service/src/pubsub/backfill/process.rs subscribes to Pub/Sub topics and dispatches to the appropriate worker based on task type.
Key Implementation Files
| Component | Path | Responsibility |
|---|---|---|
| Message ingestion | services/email_service/src/pubsub/inbox_sync/operations/upsert_message.rs |
Extracts recipients and triggers CRM jobs |
| Job enqueueing | services/email_service/src/pubsub/util.rs |
Builds and publishes Pub/Sub tasks |
| Contact aggregation | services/email_service/src/pubsub/backfill/populate_crm_contact.rs |
Computes and persists MIN/MAX timestamps |
| User batch seeding | services/email_service/src/pubsub/backfill/populate_crm_for_user.rs |
Ensures complete contact coverage |
| Backfill coordinator | services/email_service/src/pubsub/backfill/process.rs |
Routes tasks to workers |
| Company access layer | crates/entity_access/src/outbound/pg_access_repo/queries/crm_company_access.rs |
Rolls up contact aggregates to companies |
| Schema definitions | crates/macro_db_client/migrations/20260520150000_crm_contacts_name.up.sql |
Creates aggregation columns |
Summary
- Macro's CRM email aggregation runs through a Pub/Sub backfill pipeline that processes every message asynchronously
- Contact-level aggregates use
MIN/MAXSQL operations inpopulate_crm_contactto track first and last interaction timestamps - User-level seeding via
populate_crm_for_userguarantees complete backfill coverage for historical mailboxes - Company aggregates derive from linked contacts through the
crm_company_accessquery layer, avoiding data duplication - Idempotent UPSERTs ensure safe reprocessing and monotonic timestamp updates
Frequently Asked Questions
How does Macro prevent duplicate CRM entries when reprocessing the same email?
The populate_crm_contact worker uses PostgreSQL ON CONFLICT clauses with LEAST and GREATEST functions. These guarantee that reprocessing a message never corrupts existing aggregates—earlier timestamps only reduce first_interaction_ts, and later timestamps only increase last_interaction_ts.
Can the aggregation pipeline handle deleted or revoked emails?
Yes. The depopulate_message flow (referenced in services/email_service/src/pubsub/backfill/process.rs) reverses the aggregation by recomputing MIN/MAX from remaining messages. This ensures that deleting an email from Macro correctly updates CRM activity metrics.
What distinguishes real-time aggregation from historical backfill?
Real-time paths invoke enqueue_populate_crm_contacts immediately within upsert_message, providing sub-second latency for new emails. Historical backfill uses populate_crm_for_user to enqueue batches of all historical participants, processing them through the same workers without blocking ingestion.
Where are company aggregates actually stored?
Company-level timestamps are computed dynamically through crm_company_access.rs queries rather than stored columns. This design maintains normalization—updating a contact's interaction date automatically reflects in company metrics without additional write operations.
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 →