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_service handles message storage and job enqueueing
  • Backfill workers: Pub/Sub tasks process batches of messages asynchronously
  • Data models: crates/models_soup and crates/entity_access define 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:

  1. Contact resolution: Maps each email address to an existing crm_contact record or creates a new one
  2. Timestamp aggregation: Computes MIN(internal_date_ts) and MAX(internal_date_ts) across all emails involving this contact
  3. UPSERT execution: Persists the aggregated values to first_interaction_ts and last_interaction_ts columns

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() → triggers enqueue_populate_crm_contacts for real-time events
  • backfill_user() → invokes populate_crm_for_user for historical imports
  • depopulate_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/MAX SQL operations in populate_crm_contact to track first and last interaction timestamps
  • User-level seeding via populate_crm_for_user guarantees complete backfill coverage for historical mailboxes
  • Company aggregates derive from linked contacts through the crm_company_access query 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:

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 →