# How Macro’s CRM Automatically Aggregates Emails for Company and Contact Objects

> Discover how Macro CRM automatically aggregates emails for company and contact objects via a Pub/Sub pipeline, streamlining your communication tracking and providing essential interaction timestamps.

- Repository: [Macro/macro](https://github.com/macro-inc/macro)
- Tags: how-to-guide
- Published: 2026-08-16

---

**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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/services/email_service/src/pubsub/util.rs):

```rust
// 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`](https://github.com/macro-inc/macro/blob/main/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:

```rust
// 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`](https://github.com/macro-inc/macro/blob/main/services/email_service/src/pubsub/backfill/populate_crm_for_user.rs) ensures complete coverage:

```rust
// 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`](https://github.com/macro-inc/macro/blob/main/crates/models_soup/src/crm_company.rs) and [`crates/models_search/src/crm_company.rs`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/crates/entity_access/src/outbound/pg_access_repo/queries/crm_company_access.rs) rolls up contact aggregates:

```sql
-- 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`](https://github.com/macro-inc/macro/blob/main/crates/macro_db_client/migrations/20260520150000_crm_contacts_name.up.sql) establishes the aggregation columns:

```sql
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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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.