# How Macro's Search Service is Implemented Using OpenSearch and Event-Driven Processing

> Discover how Macro implements its search service using OpenSearch and event-driven processing. Learn about its two cooperating Rust services for efficient search and indexing.

- Repository: [Macro/macro](https://github.com/macro-inc/macro)
- Tags: architecture
- Published: 2026-08-18

---

**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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/crates/search_service/src/api/mod.rs) mounts several specialized sub-routers:

- **Unified search** ([`api/search/unified.rs`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/api/internal/backfill.rs) and [`extract_sync.rs`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/pdf.rs), [`docx.rs`](https://github.com/macro-inc/macro/blob/main/docx.rs), and [`markdown.rs`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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:

```rust
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:

```rust
#[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`](https://github.com/macro-inc/macro/blob/main/parsers/pdf.rs), [`docx.rs`](https://github.com/macro-inc/macro/blob/main/docx.rs), and [`markdown.rs`](https://github.com/macro-inc/macro/blob/main/markdown.rs).
- **Indexing logic** resides in [`property_search_indexer.rs`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/pdf.rs)), Microsoft Word documents ([`docx.rs`](https://github.com/macro-inc/macro/blob/main/docx.rs)), and Markdown files ([`markdown.rs`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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.