How Macro's Search Service is Implemented Using OpenSearch and Event-Driven Processing
Macro's search architecture consists of two cooperating Rust services: a Search Service that translates HTTP queries into OpenSearch DSL calls, and a Search-Processing Service that consumes Kafka events to extract searchable text and maintain indices.
The macro-inc/macro repository implements a production-grade search system that separates read and write concerns into distinct layers. This design enables low-latency query handling through the Search Service while ensuring reliable, asynchronous indexing via the Search-Processing Service. Both components share a lightweight OpenSearch client wrapper to communicate with the cluster.
Architecture Overview
The implementation splits responsibilities across two main components:
| Component | Responsibility | Key Location |
|---|---|---|
| Search Service | HTTP API façade for client queries | crates/search_service/ |
| Search-Processing Service | Event-driven indexing and text extraction | services/search_processing_service/ |
Both services inject an Arc<OpensearchClient>—defined in the shared opensearch_client crate—to handle REST communication with the OpenSearch cluster. Configuration management is centralized via the macro_env_var crate and loaded through Doppler at runtime.
Search Service HTTP Layer
The Search Service acts as the query entry point, exposing REST endpoints built with Axum. The library entry point in crates/search_service/src/lib.rs re-exports the search_router and SearchHandlerState, which encapsulates the database pool and OpenSearch client reference.
Router Structure and State Management
The router construction in crates/search_service/src/api/mod.rs mounts several specialized sub-routers:
- Unified search (
api/search/unified.rs): Handles generic multi-entity queries across document, chat, and email indices. - Entity-specific searches (
api/search/simple/*.rs): Dedicated endpoints for projects, documents, and other domain objects. - Administrative endpoints (
api/internal/backfill.rsandextract_sync.rs): Operator-facing routes for re-indexing operations.
Each handler receives SearchHandlerState, which contains an Arc<OpensearchClient> and a Postgres pool connection. The unified search handler forwards incoming queries to the client via client.search(...), delegating DSL construction to the shared client layer.
OpenAPI Documentation
The service auto-generates Swagger documentation via api/swagger.rs, exposing the SearchApiDoc struct at the /api-docs endpoint for client SDK generation.
Search-Processing Service Pipeline
Running as a separate binary in services/search_processing_service/src/main.rs, this service maintains index freshness by consuming domain events from Kafka. It orchestrates text extraction and property-level indexing without blocking the main search API.
Kafka Event Consumption
The service initializes a Kafka consumer under inbound/kafka_consumer/*.rs that subscribes to topics including document.events, email.events, and chat.events. Upon receiving an event, the consumer delegates to the domain service (domain/service.rs) which coordinates the extraction and indexing workflow.
Content Extraction
Raw binary content passes through specialized parsers located in services/search_processing_service/src/parsers/. The parser modules—pdf.rs, docx.rs, and markdown.rs—convert proprietary formats into plain text and produce an IndexedDocument structure containing searchable tokens and metadata.
Property-Based Indexing
The core write path resides in services/search_processing_service/src/outbound/property_search_indexer.rs. The DirectPropertyBackfillIndexer implements the PropertyBackfillIndexer trait, providing a reindex method that:
- Retrieves current entity data from Postgres using
sqlx::PgPool. - Enriches the document via
process_entity_property_update, which handles custom field indexing. - Upserts the document into OpenSearch via
client.index_document(...).await.
A back-fill utility in api/internal/backfill.rs allows operators to trigger full re-indexing of any EntityType by iterating database rows and feeding them through the same indexer logic.
OpenSearch Integration and Client Abstraction
The opensearch_client::OpensearchClient wrapper abstracts HTTP methods and provides typed helpers for both services:
search(query: &SearchQuery) -> Result<SearchResponse, OpensearchError>: Used by the Search Service for DSL execution.index_document(index: &str, id: &str, body: &serde_json::Value) -> Result<(), OpensearchError>: Used by the processing service for document upserts.
Index names derive from the EntityType enum in models_properties, ensuring a one-to-one mapping between Macro's domain entities and OpenSearch indices. This guarantees that documents, chats, and emails route to isolated, type-specific indices.
End-to-End Implementation Examples
The following example demonstrates initializing the Search Service with shared state:
use macro::search_service::{search_router, SearchHandlerState};
use axum::{Router, routing::get};
use std::sync::Arc;
let opensearch = OpensearchClient::new_from_env()?; // Loads vars via macro_env_var
let pg_pool = PgPool::connect(&std::env::var("DATABASE_URL")?).await?;
let state = SearchHandlerState::new(pg_pool, Arc::new(opensearch));
let app = Router::new()
.nest("/search", search_router())
.with_state(state);
// Endpoint handles: GET /search?query=invoice+2024
To trigger a manual back-fill of all documents, operators use the indexer directly:
#[tokio::main]
async fn main() -> Result<()> {
let client = Arc::new(OpensearchClient::new_from_env()?);
let db = PgPool::connect(&std::env::var("DATABASE_URL")?).await?;
let indexer = DirectPropertyBackfillIndexer::new(db, client);
// Re-index all entities of type Document
indexer.reindex_all(EntityType::Document).await?;
Ok(())
}
Summary
- Macro implements a dual-service architecture separating query handling (Search Service) from indexing (Search-Processing Service) to optimize for latency and throughput.
- The Search Service in
crates/search_service/exposes Axum-based HTTP endpoints that translate client queries into OpenSearch DSL via the sharedOpensearchClient. - The Search-Processing Service consumes Kafka events from
document.eventsand similar topics, extracting text using parsers inparsers/pdf.rs,docx.rs, andmarkdown.rs. - Indexing logic resides in
property_search_indexer.rs, implementing thePropertyBackfillIndexertrait withprocess_entity_property_updatefor enriched document upserts. - Both services share an
Arc<OpensearchClient>configured through themacro_env_varcrate, ensuring consistent TLS and credential management across the search pipeline.
Frequently Asked Questions
How does Macro handle real-time document updates in OpenSearch?
When a user creates or updates a document, the application service writes to Postgres and publishes an event to the document.events Kafka topic. The Search-Processing Service consumes this event, extracts searchable text via the appropriate parser, and calls DirectPropertyBackfillIndexer to upsert the document into the corresponding OpenSearch index immediately.
What file formats can Macro's search processing extract text from?
According to the source code in services/search_processing_service/src/parsers/, Macro supports text extraction from PDFs (pdf.rs), Microsoft Word documents (docx.rs), and Markdown files (markdown.rs). Each parser normalizes content into an IndexedDocument structure before indexing.
How is the OpenSearch client configured across Macro's services?
Both the Search Service and Search-Processing Service initialize the OpensearchClient using OpensearchClient::new_from_env(), which reads endpoint URLs, credentials, and TLS settings from environment variables managed by the macro_env_var crate and Doppler. The client is wrapped in an Arc for thread-safe sharing across async handlers.
What is the purpose of the back-fill endpoint in Macro's search service?
The back-fill endpoint defined in crates/search_service/src/api/internal/backfill.rs provides administrative capabilities to re-index existing database records without requiring new Kafka events. Operators can trigger a full re-index of any EntityType (such as projects or documents) by iterating over Postgres rows and feeding them through the PropertyBackfillIndexer logic.
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 →