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.rsunpacks OOXML structures and extracts paragraph text - Markdown files:
markdown.rsstrips 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:
- The stale entry is removed from both properties and content indices
- A re-index may be triggered for related entities
- 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, buildingDocumentInfostructs 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.rswrites dual indices—properties for metadata filtering, content for full-text token matching - Query layer: Unified search in
simple_unified.rsexecutes 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →