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

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, the consumer parses incoming events and constructs a DocumentInfo struct:

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 — parses thread metadata and sender information
  • Chats: chat.rs — handles message content and channel context
  • Channels: 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 model in services/search_processing_service/src/process/document/raw_document.rs defines this structure:

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, the service extracts text from PDF binaries:

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 unpacks OOXML structures and extracts paragraph text
  • Markdown files: 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. This module implements dual-index persistence through the models_opensearch client.

The index_property Function

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:

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

// 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 module implements the query side of the architecture. This unified search endpoint executes parallel OpenSearch queries across multiple entity indices:

// 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 writes dual indices—properties for metadata filtering, content for full-text token matching
  • Query layer: Unified search in 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 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.

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 →