mod args;
mod drain_parallel;
mod drain_serial;
mod events;
mod extraction;
mod postprocess;
mod predicates;
mod prompts;
mod quality_sample;
mod queue;
mod queue_ops;
mod reembed;
mod run;
mod scan;
mod scan_ec;
mod scheduler;
mod schemas;
mod status;
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(crate) const DEFAULT_ENRICH_MAX_ATTEMPTS: u32 = 8;
pub(crate) const DEFAULT_ENRICH_STALE_CLAIM_SECS: u64 = 1800;
pub(crate) const DEFAULT_ENRICH_RATE_LIMIT_BUFFER_SECS: u64 = 300;
pub(crate) const DEFAULT_ENRICH_CIRCUIT_BREAKER_THRESHOLD: u32 = 5;
pub(crate) const DEFAULT_ENRICH_PRESERVE_THRESHOLD: f64 = 0.7;
pub(crate) const DEFAULT_ENRICH_GROUNDING_THRESHOLD: f64 = 0.30;
pub use args::{EnrichArgs, EnrichMode, EnrichOperation, ReEmbedTarget};
pub use queue::{cleanup_queue_entry, DeadItem, DeadSummary, EnrichStatus, WaitingItem};
pub use run::run;
use crate::errors::AppError;
use queue::{enqueue_candidate_with_priority, open_queue_db, PRIORITY_HOT};
pub fn enqueue_priority_entity_descriptions(
paths: &crate::paths::AppPaths,
namespace: &str,
entity_names: &[String],
) -> Result<usize, AppError> {
let _ = namespace; 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::events::{
enrich_operation_cli_name, is_sqlite_interrupt, scan_operation_with_deadline,
};
use super::schemas::{BINDINGS_SCHEMA, BODY_ENRICH_SCHEMA, ENTITY_DESCRIPTION_SCHEMA};
use super::*;
use crate::errors::AppError;
use rusqlite::{Connection, ErrorCode};
use std::time::{Duration, Instant};
#[test]
fn every_response_schema_is_valid_json_and_strict_mode_clean() {
let schemas: [(&str, &str); 12] = [
("BINDINGS_SCHEMA", BINDINGS_SCHEMA),
("BODY_ENRICH_SCHEMA", BODY_ENRICH_SCHEMA),
("BODY_EXTRACT_SCHEMA", super::schemas::BODY_EXTRACT_SCHEMA),
(
"DEEP_RESEARCH_SYNTH_SCHEMA",
super::schemas::DEEP_RESEARCH_SYNTH_SCHEMA,
),
(
"DESCRIPTION_ENRICH_SCHEMA",
super::schemas::DESCRIPTION_ENRICH_SCHEMA,
),
(
"DOMAIN_CLASSIFY_SCHEMA",
super::schemas::DOMAIN_CLASSIFY_SCHEMA,
),
(
"ENTITY_CONNECT_SCHEMA",
super::schemas::ENTITY_CONNECT_SCHEMA,
),
("ENTITY_DESCRIPTION_SCHEMA", ENTITY_DESCRIPTION_SCHEMA),
(
"ENTITY_TYPE_VALIDATE_SCHEMA",
super::schemas::ENTITY_TYPE_VALIDATE_SCHEMA,
),
("GRAPH_AUDIT_SCHEMA", super::schemas::GRAPH_AUDIT_SCHEMA),
(
"RELATION_RECLASSIFY_SCHEMA",
super::schemas::RELATION_RECLASSIFY_SCHEMA,
),
(
"WEIGHT_CALIBRATE_SCHEMA",
super::schemas::WEIGHT_CALIBRATE_SCHEMA,
),
];
for (name, text) in schemas {
let parsed: serde_json::Value = serde_json::from_str(text)
.unwrap_or_else(|e| panic!("{name} must be valid JSON: {e}"));
assert_eq!(
parsed.get("additionalProperties"),
Some(&serde_json::Value::Bool(false)),
"{name} must set additionalProperties to false; strict mode refuses anything else"
);
let properties = parsed
.get("properties")
.and_then(|p| p.as_object())
.unwrap_or_else(|| panic!("{name} must declare an object of properties"));
let required: Vec<&str> = parsed
.get("required")
.and_then(|r| r.as_array())
.unwrap_or_else(|| panic!("{name} must declare a required array"))
.iter()
.filter_map(|v| v.as_str())
.collect();
let missing: Vec<&String> = properties
.keys()
.filter(|k| !required.contains(&k.as_str()))
.collect();
assert!(
missing.is_empty(),
"{name} declares propertie(s) absent from `required`: {missing:?}. \
Under strict mode that is a REFUSED request, not an optional field; \
model an optional value as a nullable type inside `required` instead."
);
}
}
#[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() {
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() {
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:?}"
);
}
}