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.rs and extract_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:

  1. Retrieves current entity data from Postgres using sqlx::PgPool.
  2. Enriches the document via process_entity_property_update, which handles custom field indexing.
  3. 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 shared OpensearchClient.
  • The Search-Processing Service consumes Kafka events from document.events and similar topics, extracting text using parsers in parsers/pdf.rs, docx.rs, and markdown.rs.
  • Indexing logic resides in property_search_indexer.rs, implementing the PropertyBackfillIndexer trait with process_entity_property_update for enriched document upserts.
  • Both services share an Arc<OpensearchClient> configured through the macro_env_var crate, 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:

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 →