# How Macro's Search Service Indexes Documents and Attachments Using OpenSearch: A Technical Deep Dive

> Discover how Macro's search service indexes documents and attachments using OpenSearch. Learn about its three-stage pipeline from Kafka ingestion to dual-index writing for full searchability.

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

---

**Macro's search service uses a three-stage pipeline—ingestion via Kafka, content extraction from PDFs/DOCX files, and dual-index writing to OpenSearch—to make documents and attachments fully searchable.**

The `macro-inc/macro` repository implements a sophisticated search infrastructure that transforms raw document binaries into queryable OpenSearch indices. This article examines the indexing pipeline from event ingestion through content extraction to final index persistence, with specific reference to the Rust implementation found in the search processing service.

## Overview of the Indexing Pipeline

The search processing service operates as an event-driven system. When users create or update documents, emails, chats, or channels, the service receives **Kafka events**, extracts **searchable plain text**, and writes structured data to **Amazon OpenSearch**.

The pipeline splits into two parallel index structures:

- **Properties index**: Stores metadata (title, entity type, attachment IDs, timestamps)
- **Content index**: Stores tokenized search terms for full-text relevance scoring

## Stage 1: Kafka Event Ingestion

The inbound processing layer listens to `document.created`, `document.updated`, and similar events across multiple entity types. Each entity has a dedicated consumer module in `services/search_processing_service/src/inbound/kafka_consumer/`.

### Document Consumer Implementation

In [`services/search_processing_service/src/inbound/kafka_consumer/document.rs`](https://github.com/macro-inc/macro/blob/main/services/search_processing_service/src/inbound/kafka_consumer/document.rs), the consumer parses incoming events and constructs a `DocumentInfo` struct:

```rust
let doc_info = DocumentInfo {
    id: event.document_id,
    title: event.title,
    body: event.body,
    attachment_ids: event.attachment_ids,  // References to S3-stored files
    binary_blob: fetch_blob(event.s3_key).await?,
};

```

The same pattern exists for:
- **Emails**: [`email.rs`](https://github.com/macro-inc/macro/blob/main/email.rs) — parses thread metadata and sender information
- **Chats**: [`chat.rs`](https://github.com/macro-inc/macro/blob/main/chat.rs) — handles message content and channel context
- **Channels**: [`channel.rs`](https://github.com/macro-inc/macro/blob/main/channel.rs) — indexes channel names and descriptions

### Attachment Handling During Ingestion

When documents contain attachments, the `attachment_ids` field captures references to separately stored files. The ingestion step records the **S3 storage key** for each attachment but does not immediately fetch binary content. The [`raw_document.rs`](https://github.com/macro-inc/macro/blob/main/raw_document.rs) model in [`services/search_processing_service/src/process/document/raw_document.rs`](https://github.com/macro-inc/macro/blob/main/services/search_processing_service/src/process/document/raw_document.rs) defines this structure:

```rust
pub struct RawDocument {
    pub id: Uuid,
    pub title: String,
    pub attachment_ids: Vec<Uuid>,
    pub attachment_s3_keys: Vec<String>,  // For later content extraction
}

```

## Stage 2: Content Extraction and Parsing

Binary attachments require transformation into searchable plain text. The `services/search_processing_service/src/parsers/` directory contains format-specific extractors.

### PDF Text Extraction

In [`services/search_processing_service/src/parsers/pdf.rs`](https://github.com/macro-inc/macro/blob/main/services/search_processing_service/src/parsers/pdf.rs), the service extracts text from PDF binaries:

```rust
pub fn extract_text(binary: &[u8]) -> Result<String, ParseError> {
    // Uses pdf-extract or similar Rust crate
    let document = pdf::load_from_memory(binary)?;
    let mut text = String::new();
    for page in document.pages() {
        text.push_str(&page.text());
    }
    Ok(text)
}

```

### DOCX and Markdown Parsing

- **DOCX files**: [`docx.rs`](https://github.com/macro-inc/macro/blob/main/docx.rs) unpacks OOXML structures and extracts paragraph text
- **Markdown files**: [`markdown.rs`](https://github.com/macro-inc/macro/blob/main/markdown.rs) strips formatting syntax, preserving structure through heading hierarchy

### Generic Attachment Fallback

For unsupported formats, the service stores filename and extension metadata without content extraction. Search queries can still match on attachment filenames even when body text is unavailable.

## Stage 3: OpenSearch Index Writing

The core indexing logic resides in [`services/search_processing_service/src/outbound/property_search_indexer.rs`](https://github.com/macro-inc/macro/blob/main/services/search_processing_service/src/outbound/property_search_indexer.rs). This module implements **dual-index persistence** through the `models_opensearch` client.

### The `index_property` Function

```rust
pub async fn index_property(
    client: &OpenSearchClient,
    entity_type: OpenSearchEntityType,
    entity_id: Uuid,
    document: PropertyDocument,
) -> Result<(), IndexError> {
    // Write to properties index (metadata)
    let property_doc = json!({
        "entity_type": entity_type.as_str(),
        "entity_id": entity_id,
        "title": document.title,
        "attachment_ids": document.attachment_ids,
        "created_at": document.created_at,
    });
    
    client.index(
        &format!("{}_properties", entity_type.index_prefix()),
        entity_id.to_string(),
        property_doc,
    ).await?;
    
    // Write to content index (tokenized search terms)
    let content_tokens: Vec<&str> = document.content.split_whitespace().collect();
    let content_doc = json!({
        "entity_id": entity_id,
        "tokens": content_tokens,
        "title_tokens": document.title.split_whitespace().collect::<Vec<_>>(),
    });
    
    client.index(
        &format!("{}_content", entity_type.index_prefix()),
        entity_id.to_string(),
        content_doc,
    ).await?;
    
    Ok(())
}

```

### Entity-Specific Index Mapping

The `models_opensearch` crate defines searchable entity types in [`crates/models_opensearch/src/lib.rs`](https://github.com/macro-inc/macro/blob/main/crates/models_opensearch/src/lib.rs):

| Entity Type | Index Prefix | Use Case |
|-------------|--------------|----------|
| `Documents` | `documents` | User-created documents with attachments |
| `Emails` | `emails` | Threaded email conversations |
| `Chats` | `chats` | Real-time message history |
| `Channels` | `channels` | Named discussion spaces |
| `CallRecords` | `call_records` | Voice/video meeting transcripts |
| `Projects` | `projects` | Grouped document collections |

### Complete Indexing Example

```rust
// Full pipeline: document with attachments
let plain_text = pdf::extract_text(&doc_info.binary_blob)?;

property_search_indexer::index_property(
    ctx.opensearch_client,
    OpenSearchEntityType::Documents,
    doc_info.id,
    PropertyDocument {
        title: doc_info.title,
        content: plain_text,
        attachment_ids: doc_info.attachment_ids,
        // Optional: attachment content extracted separately and concatenated
        attachment_content: fetch_and_extract_attachments(&doc_info.attachment_s3_keys).await?,
    },
).await?;

```

## Search Query Execution

The [`crates/search_service/src/api/search/simple/simple_unified.rs`](https://github.com/macro-inc/macro/blob/main/crates/search_service/src/api/search/simple/simple_unified.rs) module implements the query side of the architecture. This **unified search endpoint** executes parallel OpenSearch queries across multiple entity indices:

```rust
// Simplified query construction
let queries = vec![
    build_entity_query(OpenSearchEntityType::Documents, &search_terms),
    build_entity_query(OpenSearchEntityType::Emails, &search_terms),
    build_entity_query(OpenSearchEntityType::Chats, &search_terms),
];

let results = join_all(queries.into_iter().map(|q| client.search(q))).await;
// Merge, rank, paginate, and return unified response

```

## Stale Data Handling and Index Consistency

The search service in `crates/search_service/src/api/search/` validates database state before returning results. When OpenSearch references a **soft-deleted or hard-deleted** database record:

1. The stale entry is removed from both properties and content indices
2. A re-index may be triggered for related entities
3. The query continues with remaining valid results

This **eventual consistency pattern** tolerates transient desynchronization between the primary database and search indices.

## Summary

- **Ingestion layer**: Kafka consumers in `inbound/kafka_consumer/` parse creation and update events, building `DocumentInfo` structs with attachment references
- **Extraction layer**: Format-specific parsers in `parsers/` (PDF, DOCX, Markdown) convert binary content to searchable plain text
- **Persistence layer**: [`property_search_indexer.rs`](https://github.com/macro-inc/macro/blob/main/property_search_indexer.rs) writes dual indices—properties for metadata filtering, content for full-text token matching
- **Query layer**: Unified search in [`simple_unified.rs`](https://github.com/macro-inc/macro/blob/main/simple_unified.rs) executes parallel multi-index queries with result merging and ranking

## Frequently Asked Questions

### How does Macro handle search indexing for documents with multiple attachments?

The `DocumentInfo` struct captures `attachment_ids` as a vector during Kafka ingestion. The indexer stores these IDs in the OpenSearch properties document under `attachment_ids`. For full-text search, attachment content can be extracted separately and concatenated into the main `content` field, or indexed independently depending on configuration. The search API supports filtering queries by `has:attachments` using the presence of this field.

### What happens when a document is updated or deleted in Macro?

Update events trigger the same three-stage pipeline, with OpenSearch performing upsert operations on existing documents by `entity_id`. For deletions, the search service's query-time validation detects missing database records and **purges stale index entries** automatically. Soft deletes update a `deleted_at` timestamp in the properties index; hard deletes remove both properties and content documents entirely.

### Why does Macro use separate properties and content indices?

The **dual-index design** optimizes for different query patterns. The properties index supports structured filtering (entity type, date ranges, attachment presence) with minimal storage overhead. The content index stores whitespace-tokenized terms for **relevance-scored full-text matching**. This separation allows independent scaling: metadata queries hit smaller, cached indices while content queries traverse larger token stores with optimized analyzers.

### Which OpenSearch client does Macro use for Rust integration?

The `models_opensearch` crate wraps the official `opensearch` Rust client with domain-specific methods for entity type routing, index naming, and error handling. The [`property_search_indexer.rs`](https://github.com/macro-inc/macro/blob/main/property_search_indexer.rs) module depends on this abstraction rather than calling OpenSearch APIs directly, enabling centralized retry logic, connection pooling, and observability instrumentation across all search write operations.