AI Processing Pipeline for the Document Cognition Service: A Technical Deep Dive

The Document Cognition Service processes raw documents through a nine-stage Rust-based pipeline that routes HTTP requests via Axum, retrieves document context, registers AI streams, invokes external LLMs through the mcp_client crate, and assembles structured responses containing generated text, citations, and attachment references.

The Document Cognition Service in the macro-inc/macro repository transforms unstructured documents into AI-enhanced representations through a sophisticated streaming architecture. When clients such as the Lexical SDK interact with the service via the /cognitionv2 endpoint, the system orchestrates document retrieval, LLM communication via Anthropic, and real-time response assembly. This guide traces the complete AI processing pipeline from HTTP request entry to final JSON response delivery.

Pipeline Architecture Overview

The service operates as an Axum-based microservice that wires together HTTP routing, request validation, and asynchronous AI processing. Each cognition request travels through a well-defined sequence of modules that handle document retrieval, stream management, and persistent storage of results.

Step-by-Step Processing Flow

HTTP Request Routing and Validation

Client requests enter through the Axum router defined in services/document_cognition_service/src/api/mod.rs. The system registers the POST /cognitionv2/:document_id endpoint, which extracts the document identifier from the URL path and validates the JSON payload before dispatching to the handler.

Structured Completion Handler

The router forwards validated requests to the structured completion handler in services/document_cognition_service/src/api/structured_completion.rs. This module extracts the document ID, validates the incoming payload structure, and prepares the AI request context. It serves as the central orchestrator that coordinates subsequent pipeline stages.

Document Context Retrieval

Before invoking the LLM, the handler retrieves the document's content and chat history via services/document_cognition_service/src/service/get_chat.rs. This module loads the raw document data from the document storage service and assembles the conversation context required for the AI processing.

AI Stream Registration and Management

The service creates a new AI stream through services/document_cognition_service/src/service/ai_stream_registry.rs. This module generates a unique stream ID, stores metadata including the user identifier, document reference, and target model, and registers an associated chat object via services/document_cognition_service/src/service/chat_renamer.rs. The registry maintains the lifecycle of streaming jobs throughout the cognition process.

LLM Invocation and Token Streaming

The actual LLM communication occurs in services/document_cognition_service/src/service/notification.rs, which utilizes the mcp_client crate located in crates/mcp_client to forward structured completion requests to external providers such as Anthropic. This module handles HTTP retries, error mapping, and manages the WebSocket notification infrastructure. As the LLM streams response tokens, the service aggregates them into manageable chunks.

Stream Aggregation and Persistence

Incoming tokens flow into the Stream model defined in services/document_cognition_service/src/model/stream.rs. This model aggregates raw tokens into coherent chunks and persists the streaming data to PostgreSQL via the SQLx-based client in crates/macro_db_client/src/dcs. The dual approach of in-memory aggregation and database persistence ensures durability while maintaining low-latency delivery.

Citations and Attachments Processing

When the LLM emits citations or references to document attachments, the handler captures these relationships in services/document_cognition_service/src/model/response/attachments.rs. This module manages entity mentions and attachment metadata, enabling clients to retrieve referenced files through presigned URL endpoints. The system maintains referential integrity between generated content and source document elements.

Response Construction and Delivery

Once the token stream completes, the service constructs the final response using the CognitionV2ResponseData struct defined in services/document_cognition_service/src/model/response/mod.rs. This structure encapsulates the generated text, citation arrays, and attachment identifiers. The handler serializes this data to JSON for the HTTP response while triggering WebSocket notifications via the notification service to update connected clients in real time.

Database Persistence and Storage

The AI processing pipeline relies on crates/macro_db_client/src/dcs for durable storage of streaming data, citations, and entity relationships. This SQLx-based data access layer persists AI stream metadata, token chunks, and attachment references to PostgreSQL, ensuring that cognition results remain available for subsequent queries and historical analysis.

Client Integration Examples

The Lexical SDK interacts with the Document Cognition Service using the following Rust pattern:

let client = lexical_client::LexicalClient::new("https://document-cognition.macro.com");
let result = client
    .parse_cognition_v2("doc_12345")
    .await
    .expect("failed to run cognition");

println!("Generated text: {}", result.generated_text);
println!("Citations: {:?}", result.citations);

For direct HTTP access, clients can invoke the endpoint using curl:

curl -X POST \
     -H "Authorization: Bearer $TOKEN" \
     -H "Content-Type: application/json" \
     https://document-cognition.macro.com/cognitionv2/doc_12345 \
     -d '{"prompt":"Summarize this document"}'

The response body conforms to the CognitionV2ResponseData schema defined in the Rust source, containing the full generated content and associated metadata.

Summary

  • The Document Cognition Service implements a Rust-based microservice architecture using Axum for HTTP handling and Tokio for async processing.
  • Request processing flows through src/api/mod.rs to the structured completion handler in src/api/structured_completion.rs, which coordinates the entire pipeline.
  • Document context is retrieved via src/service/get_chat.rs before AI processing begins.
  • Stream management occurs through src/service/ai_stream_registry.rs, which tracks LLM jobs and registers chat contexts via src/service/chat_renamer.rs.
  • LLM communication leverages the mcp_client crate with retry logic implemented in src/service/notification.rs.
  • Response assembly uses src/model/response/mod.rs to structure the final JSON output containing text, citations, and attachment references.
  • Persistence is handled by crates/macro_db_client/src/dcs using PostgreSQL for durable storage of streams and metadata.

Frequently Asked Questions

What LLM provider does the Document Cognition Service use?

The service currently integrates with Anthropic for language model processing, communicating through the mcp_client crate located in crates/mcp_client. The notification service in src/service/notification.rs handles the HTTP transport, retry logic, and error mapping for these external calls.

How does the service handle real-time streaming of AI responses?

The system creates a dedicated AI stream via src/service/ai_stream_registry.rs that maintains a unique stream ID and metadata. As tokens arrive from the LLM, they aggregate in the Stream model defined in src/model/stream.rs, which simultaneously writes chunks to PostgreSQL via crates/macro_db_client/src/dcs while forwarding data to the client.

Where does the Document Cognition Service store processing results?

Streaming data, citations, and attachment references persist to PostgreSQL through the SQLx-based client in crates/macro_db_client/src/dcs. The src/model/stream.rs module manages the in-memory representation and database synchronization, ensuring that AI-generated content and metadata remain durable and queryable.

How do clients retrieve document attachments referenced by the AI?

When the LLM generates citations or attachment references, the system stores these relationships in src/model/response/attachments.rs. Clients subsequently retrieve the actual files through presigned URL endpoints, enabling secure, temporary access to referenced document components without exposing underlying storage credentials.

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 →