sqlite-graphrag 1.1.8

Persistent GraphRAG memory for Claude Code, Codex, Cursor, and 27 AI agents — one self-contained ~19 MiB Rust binary, zero daemon. Never re-explain your codebase again. Hybrid retrieval (FTS5 BM25 + cosine similarity + multi-hop graph traversal) surfaces the right memory in milliseconds. Embedding and entity enrichment run as parallel REST calls against your cloud LLM — no fragile headless subprocesses, no ONNX runtime, no model downloads. Soft-delete with full version history, transactional atomic writes, BLAKE3-tracked mutations. OAuth-only: raw API keys ABORT the spawn.
Documentation
// v1.0.97: modularised into queue.rs, scan.rs, postprocess.rs, extraction.rs.
// See ADR-0056 (closes the ADR-0046 "Known Tech Debt (v1.0.89+)" item).
// v1.1.8 Wave C1: further split into schemas/args/events/run modules (≤800 LOC).

//! Handler for the `enrich` CLI subcommand (GAP-14 + GAP-18).
//!
//! Enriches the knowledge graph by running LLM-powered analysis over memories
//! and entities that are missing key structural data. Operations are:
//!
//! - `memory-bindings`: memories without `memory_entities` rows get entity extraction
//! - `entity-descriptions`: entities with NULL/empty descriptions get LLM descriptions
//! - `body-enrich`: memories with short bodies get expanded by the LLM (GAP-18)
//! - `re-embed`: memories without a vector row get re-embedded without rewriting body
//!
//! Architecture mirrors `ingest_claude.rs`: SCAN → JUDGE (LLM) → PERSIST, with a
//! SQLite queue DB derived next to `--db` (GAP-SG-64) for resume/retry support.
// Workload: Subprocess I/O-bound (claude/codex API calls with network wait)
//!
//! # DRY note
//!
//! v1.0.97: `claude_runner.rs` now hosts the shared Claude invocation helpers
//! (`run_claude`, `parse_claude_output`, `spawn_with_memory_limit`). The queue
//! DB schema in `ingest_claude.rs` still duplicates `open_queue_db` here — a
//! future pass can unify them.

mod args;
mod drain_parallel;
mod drain_serial;
mod events;
mod extraction;
// extraction_providers is a child of extraction.rs (#[path])
// extraction_ops_* are submodules of extraction.rs (#[path])
mod postprocess;
mod predicates;
mod prompts;
mod queue;
mod queue_ops;
mod run;
mod scan;
mod scan_ec;
mod quality_sample;
mod scheduler;
mod schemas;
mod status;

// Surface schemas + defaults for sibling modules (`super::BINDINGS_SCHEMA`, etc.).
#[allow(unused_imports)]
pub(crate) use schemas::*;

pub(crate) const DEFAULT_RATE_LIMIT_WAIT: u64 = 60;
pub(crate) const DEFAULT_BODY_ENRICH_MIN_CHARS: usize = 500;
pub(crate) const DEFAULT_BODY_ENRICH_MAX_CHARS: usize = 2000;

pub use args::{EnrichArgs, EnrichMode, EnrichOperation, ReEmbedTarget};
pub use queue::{cleanup_queue_entry, DeadItem, DeadSummary, EnrichStatus, WaitingItem};
pub use run::run;

// External types re-exported so historical `use super::*` in child modules keeps working.
#[allow(unused_imports)]
pub(crate) use crate::commands::ingest_claude::find_claude_binary;
#[allow(unused_imports)]
pub(crate) use crate::constants::MAX_MEMORY_BODY_LEN;
#[allow(unused_imports)]
pub(crate) use crate::entity_type::EntityType;
#[allow(unused_imports)]
pub(crate) use crate::errors::AppError;
#[allow(unused_imports)]
pub(crate) use crate::paths::AppPaths;
#[allow(unused_imports)]
pub(crate) use crate::storage::connection::{ensure_db_ready, open_rw};
#[allow(unused_imports)]
pub(crate) use crate::storage::entities::{self as entities, NewEntity, NewRelationship};
#[allow(unused_imports)]
pub(crate) use crate::storage::memories;
#[allow(unused_imports)]
pub(crate) use rusqlite::Connection;
#[allow(unused_imports)]
pub(crate) use serde::{Deserialize, Serialize};
#[allow(unused_imports)]
pub(crate) use std::io::Write;
#[allow(unused_imports)]
pub(crate) use std::path::{Path, PathBuf};
#[allow(unused_imports)]
pub(crate) use std::time::Instant;
#[allow(unused_imports)]
pub(crate) use prompts::ENTITY_DESCRIPTION_PROMPT_PREFIX;

use queue::{enqueue_candidate_with_priority, open_queue_db, PRIORITY_HOT};

/// GAP-CLI-PRIO-02: enqueue entity-descriptions for a hot set of entity names
/// into the enrich sidecar queue with elevated priority.
pub fn enqueue_priority_entity_descriptions(
    paths: &crate::paths::AppPaths,
    namespace: &str,
    entity_names: &[String],
) -> Result<usize, AppError> {
    let _ = namespace; // queue keys are entity names; namespace is scoped by DB path
    let queue_path = crate::paths::sidecar_path(&paths.db, ".enrich-queue.sqlite");
    let queue = open_queue_db(&queue_path)?;
    let mut n = 0usize;
    for name in entity_names {
        enqueue_candidate_with_priority(
            &queue,
            name,
            "entity",
            "EntityDescriptions",
            PRIORITY_HOT,
        );
        n += 1;
    }
    Ok(n)
}


#[cfg(test)]
mod tests {
    use super::*;
    use super::events::{enrich_operation_cli_name, is_sqlite_interrupt, scan_operation_with_deadline};
    use rusqlite::ErrorCode;
    use std::time::{Duration, Instant};

    #[test]
    fn bindings_schema_is_valid_json() {
        let _: serde_json::Value =
            serde_json::from_str(BINDINGS_SCHEMA).expect("BINDINGS_SCHEMA must be valid JSON");
    }

    #[test]
    fn entity_description_schema_is_valid_json() {
        let _: serde_json::Value = serde_json::from_str(ENTITY_DESCRIPTION_SCHEMA)
            .expect("ENTITY_DESCRIPTION_SCHEMA must be valid JSON");
    }

    #[test]
    fn body_enrich_schema_is_valid_json() {
        let _: serde_json::Value = serde_json::from_str(BODY_ENRICH_SCHEMA)
            .expect("BODY_ENRICH_SCHEMA must be valid JSON");
    }

    // v1.1.06 — GAP-ENTITY-CONNECT-SCAN-CARTESIAN observability + interrupt

    #[test]
    fn enrich_operation_cli_name_pair_ops_are_kebab_case() {
        assert_eq!(
            enrich_operation_cli_name(&EnrichOperation::EntityConnect),
            "entity-connect"
        );
        assert_eq!(
            enrich_operation_cli_name(&EnrichOperation::CrossDomainBridges),
            "cross-domain-bridges"
        );
        assert_eq!(
            enrich_operation_cli_name(&EnrichOperation::EntityDescriptions),
            "entity-descriptions"
        );
    }

    #[test]
    fn is_sqlite_interrupt_detects_operation_interrupted() {
        let ffi_err = rusqlite::ffi::Error {
            code: ErrorCode::OperationInterrupted,
            extended_code: 9,
        };
        let err = rusqlite::Error::SqliteFailure(ffi_err, Some("interrupted".into()));
        assert!(is_sqlite_interrupt(&err));

        let busy = rusqlite::ffi::Error {
            code: ErrorCode::DatabaseBusy,
            extended_code: 5,
        };
        let busy_err = rusqlite::Error::SqliteFailure(busy, None);
        assert!(!is_sqlite_interrupt(&busy_err));
    }

    #[test]
    fn scan_deadline_already_elapsed_returns_timeout() {
        // Past deadline must fail fast without running SQL (exit path → Timeout).
        use clap::Parser;
        let cli = crate::cli::Cli::try_parse_from([
            "sqlite-graphrag",
            "enrich",
            "--operation",
            "entity-connect",
            "--mode",
            "openrouter",
            "--openrouter-model",
            "test/model",
            "--dry-run",
            "--limit",
            "1",
        ])
        .expect("parse enrich args");
        let Some(crate::cli::Commands::Enrich(args)) = cli.command else {
            panic!("expected Commands::Enrich");
        };
        let conn = Connection::open_in_memory().unwrap();
        let past = Instant::now() - Duration::from_secs(1);
        let err = scan_operation_with_deadline(&conn, "global", &args, Some(past))
            .expect_err("elapsed deadline must Timeout");
        match err {
            AppError::Timeout { .. } => {}
            other => panic!("expected Timeout, got {other:?}"),
        }
    }

    #[test]
    fn interrupt_handle_maps_long_query_to_sqlite_interrupt() {
        // Live SQLite: watchdog interrupt aborts a recursive CTE (same mechanism
        // as scan_operation_with_deadline). Confirms rusqlite InterruptHandle.
        let conn = Connection::open_in_memory().unwrap();
        let handle = conn.get_interrupt_handle();
        let stop = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
        let stop_w = std::sync::Arc::clone(&stop);
        let watchdog = std::thread::spawn(move || {
            std::thread::sleep(Duration::from_millis(30));
            if !stop_w.load(std::sync::atomic::Ordering::Relaxed) {
                handle.interrupt();
            }
        });
        let result = conn.query_row(
            "WITH RECURSIVE t(x) AS (
                 SELECT 1
                 UNION ALL
                 SELECT x + 1 FROM t WHERE x < 500000000
             )
             SELECT COUNT(*) FROM t",
            [],
            |r| r.get::<_, i64>(0),
        );
        stop.store(true, std::sync::atomic::Ordering::Relaxed);
        let _ = watchdog.join();
        let err = result.expect_err("recursive CTE must be interrupted");
        assert!(
            is_sqlite_interrupt(&err),
            "expected SQLITE_INTERRUPT, got {err:?}"
        );
    }
}