Rust SDK quickstart

Build an incremental Rust pipeline that chunks Markdown, batches and caches local embeddings, and keeps a Postgres/pgvector table in sync.

Time
~15 minutes
Language
Rust 1.89+
Requires
Postgres with pgvector
Version
v 1.0.14

This tutorial builds the Rust version of the text-embedding pipeline: read Markdown files, split them into chunks, embed the chunks, and declare the rows that should exist in Postgres.

The important word is declare. Your code describes the current target state; CocoIndex works out which rows to insert, update, retain, or delete on each run.

Create the project

Create a binary crate with a directory for input files:

bash
cargo new cocoindex-rust-quickstart
cd cocoindex-rust-quickstart
mkdir markdown_files

Until the Rust SDK is published separately, depend on the CocoIndex workspace from GitHub. Add these dependencies to Cargo.toml:

Cargo.toml
[dependencies]
cocoindex = { git = "https://github.com/cocoindex-io/cocoindex", features = [
    "postgres",
    "text",
    "fastembed",
] }
dotenvy = "0.15"
serde = { version = "1", features = ["derive"] }
tokio = { version = "1", features = ["full"] }

You also need a Postgres database where the pgvector extension can be created. Set its connection URL:

bash
export POSTGRES_URL=postgres://cocoindex:cocoindex@localhost/cocoindex

Add one or more .md files under markdown_files/.

Define resources and the row type

Start src/main.rs with the imports, constants, and two typed context keys:

main.rs
use std::path::PathBuf;

use cocoindex::connectors::postgres;
use cocoindex::ops::sentence_transformers::SentenceTransformerEmbedder;
use cocoindex::ops::text::{RecursiveChunkConfig, RecursiveSplitter};
use cocoindex::prelude::*;

const EMBED_MODEL: &str = "sentence-transformers/all-MiniLM-L6-v2";
const PG_SCHEMA: &str = "coco_examples";
const TABLE: &str = "doc_embeddings";

cocoindex::context_key!(static DB: postgres::Database);
cocoindex::context_key!(
    static EMBEDDER: SentenceTransformerEmbedder,
    detect_change
);

#[derive(Clone, Serialize, Deserialize, SchemaFields)]
struct DocEmbedding {
    id: i64,
    filename: String,
    chunk_start: i32,
    chunk_end: i32,
    text: String,
    #[coco(vector)]
    embedding: Vec<f32>,
}

context_key! gives each provided resource a stable name and type, using the Rust identifier as the key name. Change detection is disabled for the database connection. It is enabled for the embedder, whose type defines the model identity that invalidates dependent memoized work. Runtime handles such as connection pools do not need to be serializable.

SchemaFields derives the database columns from DocEmbedding. The vector dimension is intentionally absent because it will come from the loaded model.

Process one file

Add a processing function that splits one file and embeds its chunks:

main.rs
#[cocoindex::function]
async fn process_file(ctx: &Ctx, file: FileEntry) -> Result<Vec<DocEmbedding>> {
    let filename = file.key();
    let text = file.content_str()?;
    let chunks = RecursiveSplitter::new()?.split_with(
        &text,
        RecursiveChunkConfig {
            chunk_size: 2_000,
            min_chunk_size: None,
            chunk_overlap: Some(500),
            language: Some("markdown".to_string()),
        },
    );

    let texts: Vec<String> = chunks
        .iter()
        .map(|chunk| chunk.text(&text).to_string())
        .collect();
    let embedder = ctx.get_key(&EMBEDDER)?.clone();
    let embedding_ctx = ctx.clone();
    let embeddings = ctx
        .map(texts.clone(), move |chunk_text| {
            let embedder = embedder.clone();
            let ctx = embedding_ctx.clone();
            async move { embedder.embed(&ctx, chunk_text).await }
        })
        .await?;

    let mut id_gen = IdGenerator::new();
    let mut rows = Vec::with_capacity(texts.len());
    for ((chunk, chunk_text), embedding) in chunks.iter().zip(texts).zip(embeddings) {
        let id = i64::try_from(id_gen.next_id(ctx, &chunk_text).await?)
            .map_err(|_| Error::engine("generated id does not fit in BIGINT"))?;
        rows.push(DocEmbedding {
            id,
            filename: filename.clone(),
            chunk_start: chunk.start.char_offset as i32,
            chunk_end: chunk.end.char_offset as i32,
            text: chunk_text,
            embedding,
        });
    }
    Ok(rows)
}

SentenceTransformerEmbedder::embed is item-shaped at the call site, but its implementation automatically groups concurrent cache misses into batches of up to 64. Repeated texts are served from CocoIndex’s memo store.

Declare the table and rows

Now add the app’s main processing function:

main.rs
async fn app_main(ctx: Ctx, sourcedir: PathBuf) -> Result<()> {
    let vector_dim = ctx.get_key(&EMBEDDER)?.dimension();
    let schema = postgres::TableSchema::from_row::<DocEmbedding>(["id"])?
        .with_vector_dim("embedding", vector_dim)?;
    let table = postgres::mount_table_target(
        &ctx,
        &DB,
        TABLE,
        schema,
        Some(PG_SCHEMA),
    )
    .await?;
    table.declare_vector_index(
        &ctx,
        "embedding",
        postgres::VectorIndexOptions {
            method: "hnsw",
            ..Default::default()
        },
    )?;

    let files = walk_items(&sourcedir, &["**/*.md"])?;
    let rows_by_file = mount_each!(files, |file| process_file(ctx, file)).await?;

    for rows in rows_by_file {
        for row in rows {
            table.declare_row(&ctx, &row)?;
        }
    }
    Ok(())
}

mount_each! gives every relative file path its own stable processing component. It also fingerprints process_file and its arguments, so an unchanged file can skip the entire component on the next run. If a file or chunk disappears, the rows its component used to own are removed during table reconciliation.

Build the environment and run

Finish src/main.rs by loading the resources and providing them to an environment:

main.rs
fn database_url() -> String {
    std::env::var("POSTGRES_URL")
        .unwrap_or_else(|_| "postgres://cocoindex:cocoindex@localhost/cocoindex".to_string())
}

#[tokio::main]
async fn main() -> Result<()> {
    dotenvy::dotenv().ok();
    let sourcedir = std::env::args()
        .nth(1)
        .map(PathBuf::from)
        .unwrap_or_else(|| PathBuf::from("markdown_files"));

    let database = postgres::Database::connect(&database_url()).await?;
    let embedder = SentenceTransformerEmbedder::load(EMBED_MODEL).await?;
    let app = Environment::builder()
        .db_path(".cocoindex_db")
        .provide_key(&DB, database)
        .provide_key(&EMBEDDER, embedder)
        .build()
        .await?
        .app("RustTextEmbeddingQuickstart")
        .await?;

    let stats = app.run(move |ctx| app_main(ctx, sourcedir)).await?;
    println!("{stats}");
    Ok(())
}

Run the pipeline:

bash
cargo run

Run it again without changing an input: the target rows and embeddings are skipped. Then edit, add, or delete a Markdown file and rerun; CocoIndex updates only the affected component and reconciles its rows.

Next steps