OpenSearch Indexing Pipeline for the Search Processing Service: How Macro Indexes Documents
The OpenSearch indexing pipeline in Macro's search processing service consumes document events from Kafka, maps them to specific indexing actions, and executes a multi-stage flow that extracts content, enriches metadata with indexed properties, and bulk-upserts searchable chunks into OpenSearch with automatic rollback on failure to ensure idempotency.
The search processing service in the macro-inc/macro repository serves as the primary ingestion point for document-related search indexing. Written in Rust, this service bridges the internal event bus and the OpenSearch cluster, coordinating text extraction, metadata resolution, and chunk-based indexing. Understanding this OpenSearch indexing pipeline is critical for debugging search availability issues or extending support for new file types.
Event Ingestion and Action Dispatch
The pipeline begins at the Kafka consumer boundary. In services/search_processing_service/src/inbound/kafka_consumer/document.rs, the consumer reads a DocumentMacroEvent and invokes describe_document_event to translate the raw event into a strongly-typed DocumentIndexAction.
The possible actions include:
- Extract text – Triggered by
ContentUploadedevents requiring full-text indexing - Extract sync content – For synchronous extraction workflows
- Refresh name – When document metadata changes but content remains static
- Remove – Triggered by
Purgedevents to delete from the index - Ignore – For events that require no search-side processing
The process_document_event function (lines ~120-170) matches on this action enum and dispatches to the appropriate handler in the process::document module.
Document Processing and Text Extraction
For content extraction events, the pipeline enters the core indexing flow defined in services/search_processing_service/src/process/document/mod.rs. The process_extract_text_message function (lines 56-71) builds a SearchExtractorMessage and delegates to raw_document::update_search_with_raw_document.
Metadata Resolution
Before parsing content, get_document_info in services/search_processing_service/src/process/document/document_info.rs queries the MacroDB to retrieve essential metadata including owner identity, file type classification, and document version details. This context accompanies the document through all subsequent stages.
Content Parsing
The service supports pluggable parsers based on file type:
- Markdown – The
parse_markdown_legacyfunction inservices/search_processing_service/src/parsers/markdown.rshandles.mdfiles - PDF – When compiled with the
pdffeature,parse_pdf_pagesinservices/search_processing_service/src/parsers/pdf.rsextracts text page-by-page
The parsed content is then segmented into searchable chunks, each receiving a unique node_id. For every chunk, the system constructs an UpsertDocumentArgs structure containing the raw text, metadata, and parent document_id.
Indexed Properties Enrichment
Before sending data to OpenSearch, the pipeline denormalizes entity properties. In services/search_processing_service/src/process/document/raw_document.rs (lines 88-107), the attach_indexed_properties function calls entity_properties_get_query::get_entity_properties_for_index to fetch tags and custom fields. These properties are injected into every chunk, enabling faceted search and filtering without JOIN operations at query time.
Bulk Upsert and Error Handling
The opensearch_client.bulk_upsert_documents(&upserts, index_override) call (lines 48-63 in raw_document.rs) transmits all chunks to the OpenSearch cluster in a single request. The pipeline implements strict consistency guarantees: if any chunk reports an indexing error, the entire parent document is deleted via delete_document to prevent partial index states. This ensures the system remains idempotent—retrying a failed event will not produce duplicate or corrupted search entries.
Parent-Only Documents
Not all file types contain searchable text. For Canvas files or non-searchable binaries, the pipeline invokes upsert_parent_document (lines 111-185) to create a metadata-only entry. This parent document excludes content chunks but includes the document name, owner, and indexed properties, ensuring the file remains discoverable through title searches and property filters.
Document Lifecycle Operations
Beyond ingestion, the pipeline handles structural changes to existing documents.
In services/search_processing_service/src/process/document/mod.rs:
- Deletion – The
process_remove_messagefunction (lines 8-15) processesDocumentMacroEvent::Purgedby callingopensearch_client.delete_documentto remove all associated chunks and parent records - Renaming –
process_update_name_message(lines 23-30) handlesDocumentMacroEvent::Updatedwith name changes by fetching the latest title from the database and updating only the parent document metadata in OpenSearch, avoiding expensive re-indexing of unchanged content
Reliability and Context Management
All pipeline stages execute within a retry wrapper defined in services/search_processing_service/src/inbound/kafka_consumer/mod.rs. The retry_processing utility applies exponential back-off using PROCESSING_RETRY_BASE_DELAY and respects MAX_PROCESSING_ATTEMPTS before dead-lettering the event.
Shared resources are propagated via KafkaProcessingContext (defined in services/search_processing_service/src/inbound/kafka_consumer/context.rs). This struct carries the OpenSearch client, database connection pool, S3 client, and lexical client, ensuring each processing stage accesses initialized connections without recreating them per message.
Code Examples
The following examples demonstrate how to trigger the pipeline programmatically or invoke core functions directly for testing.
Simulating a document upload event through the Kafka consumer:
use macro_event_broker::MacroEvent;
use documents::domain::events::DocumentMacroEvent;
use services::search_processing_service::inbound::kafka_consumer::{
process_document_event, KafkaProcessingContext
};
use services::search_processing_service::inbound::kafka_consumer::context::KafkaProcessingContext;
// Context normally initialized at service startup
let ctx = KafkaProcessingContext::new(/* db, opensearch_client, s3_client */);
// Construct a ContentUploaded event for a Markdown file
let event = DocumentMacroEvent::new(
/* user_id, document_id, file_type = FileType::Md, ... */
);
// Execute the full pipeline
let outcome = process_document_event(&ctx, &event, 0, 123).await;
assert!(matches!(outcome, EventOutcome::Processed));
Directly invoking the raw document upsert for unit testing:
use opensearch_client::OpensearchClient;
use sqs_client::search::document::SearchExtractorMessage;
use services::search_processing_service::process::document::raw_document::update_search_with_raw_document;
use documents::domain::types::FileType;
let client = OpensearchClient::new(/* connection params */).await?;
let db = sqlx::Pool::<sqlx::Postgres>::connect(/* database url */).await?;
let s3 = s3_client::S3::new(/* aws config */);
let bucket = "macro-documents";
let msg = SearchExtractorMessage {
user_id: "user-123".into(),
document_id: "doc-456".into(),
file_type: FileType::Md,
document_version_id: Some("version-789".into()),
index_override: None,
};
update_search_with_raw_document(&client, &db, &s3, bucket, &msg).await?;
Summary
- The OpenSearch indexing pipeline begins with Kafka event ingestion in
document.rs, where events map to concreteDocumentIndexActionvariants - Content extraction flows through
process_extract_text_message, which coordinates parsing via type-specific modules likemarkdown.rsandpdf.rs - Indexed properties are denormalized and attached to every chunk via
attach_indexed_propertiesbefore bulk insertion - The system guarantees atomic updates by deleting the entire parent document if any chunk fails during
bulk_upsert_documents - Parent-only documents support non-text file types by indexing metadata without content chunks
- Exponential backoff and shared resource management via
KafkaProcessingContextensure reliable processing under load
Frequently Asked Questions
What happens when a PDF fails to parse during indexing?
If the parse_pdf_pages function in services/search_processing_service/src/parsers/pdf.rs returns an error, the error propagates up through update_search_with_raw_document, triggering the retry logic. After exhausting MAX_PROCESSING_ATTEMPTS, the event is dead-lettered. Because the bulk upsert occurs only after successful parsing, partial index corruption is avoided—no chunks are written until all content is successfully extracted.
How does the pipeline handle partial indexing failures?
The pipeline implements an all-or-nothing consistency model. In services/search_processing_service/src/process/document/raw_document.rs, if opensearch_client.bulk_upsert_documents reports errors for any chunk in the batch, the system immediately invokes delete_document to remove the parent document and all previously inserted chunks. This rollback mechanism ensures the OpenSearch cluster never contains incomplete document data that would pollute search results.
What is the difference between a full document and a parent-only document?
A full document contains both a parent record and multiple child chunks representing searchable text segments, created when parsers successfully extract content. A parent-only document, created via upsert_parent_document for file types like Canvas or encrypted binaries, contains only the metadata record with no child chunks. Both types receive indexed properties and support name-based search, but only full documents support content-body queries.
How are document renames propagated to OpenSearch without re-indexing content?
When a DocumentMacroEvent::Updated signals a name change, the process_update_name_message function queries the MacroDB for the current document title and constructs a targeted update request. This updates only the name field in the parent document metadata within OpenSearch, avoiding the expensive text extraction, parsing, and chunking stages required for full re-indexing.
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 →