use super::*;
pub(super) use super::quality_sample::*;
pub(super) use super::scan_ec::*;
#[cfg(test)]
#[path = "scan_tests_a.rs"]
mod tests_a;
#[cfg(test)]
#[path = "scan_tests_b.rs"]
mod tests_b;
const UNBOUND_MEMORY_PREDICATE: &str =
"NOT EXISTS (SELECT 1 FROM memory_entities me WHERE me.memory_id = m.id)";
const NULL_DESCRIPTION_PREDICATE: &str = "(description IS NULL OR description = '')";
const SHORT_BODY_PREDICATE: &str = "LENGTH(COALESCE(m.body,'')) < ?2";
#[allow(dead_code)] const GENERIC_DESCRIPTION_PREDICATE: &str = "(description LIKE '%ingested%' \
OR description LIKE '%imported%' OR description LIKE '%added%' \
OR length(description) < 30)";
const HIGH_WEIGHT_PREDICATE: &str = "r.weight >= 0.7";
const GENERIC_RELATION_PREDICATE: &str = "r.relation = 'applies_to'";
fn reembed_memory_predicate(dim: usize) -> String {
format!(
"NOT EXISTS (SELECT 1 FROM memory_embeddings me WHERE me.memory_id = m.id \
AND me.dim = {dim} AND LENGTH(me.embedding) > 0)"
)
}
fn reembed_entity_predicate(dim: usize) -> String {
format!(
"NOT EXISTS (SELECT 1 FROM entity_embeddings ev WHERE ev.entity_id = e.id \
AND ev.dim = {dim} AND LENGTH(ev.embedding) > 0)"
)
}
fn reembed_chunk_predicate(dim: usize) -> String {
format!(
"NOT EXISTS (SELECT 1 FROM chunk_embeddings ce WHERE ce.chunk_id = c.id \
AND ce.dim = {dim} AND LENGTH(ce.embedding) > 0)"
)
}
pub(super) fn scan_unbound_memories(
conn: &Connection,
namespace: &str,
limit: Option<usize>,
name_filter: &[String],
) -> Result<Vec<(i64, String, String)>, AppError> {
let limit_clause = limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
if name_filter.is_empty() {
let sql = format!(
"SELECT m.id, m.name, m.body
FROM memories m
WHERE m.namespace = ?1
AND m.deleted_at IS NULL
AND {UNBOUND_MEMORY_PREDICATE}
ORDER BY m.id
{limit_clause}"
);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params![namespace], |r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
))
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
} else {
let placeholders: Vec<String> = (2..=name_filter.len() + 1)
.map(|i| format!("?{i}"))
.collect();
let in_clause = placeholders.join(", ");
let sql = format!(
"SELECT m.id, m.name, m.body
FROM memories m
WHERE m.namespace = ?1
AND m.deleted_at IS NULL
AND m.name IN ({in_clause})
AND {UNBOUND_MEMORY_PREDICATE}
ORDER BY m.id
{limit_clause}"
);
let mut params_vec: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(1 + name_filter.len());
params_vec.push(&namespace);
for n in name_filter {
params_vec.push(n);
}
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(
rusqlite::params_from_iter(params_vec.iter().copied()),
|r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
))
},
)?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
}
pub(super) fn scan_bound_memories_for_augment(
conn: &Connection,
namespace: &str,
limit: Option<usize>,
name_filter: &[String],
) -> Result<Vec<String>, AppError> {
if name_filter.is_empty() {
return Err(AppError::Validation(
"augment-bindings requires an explicit subset: pass --names or \
--names-file (it refuses to re-scan the whole namespace)"
.into(),
));
}
let limit_clause = limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
let placeholders: Vec<String> = (2..=name_filter.len() + 1)
.map(|i| format!("?{i}"))
.collect();
let in_clause = placeholders.join(", ");
let sql = format!(
"SELECT m.name
FROM memories m
WHERE m.namespace = ?1
AND m.deleted_at IS NULL
AND m.name IN ({in_clause})
AND EXISTS (
SELECT 1 FROM memory_entities me WHERE me.memory_id = m.id
)
ORDER BY m.id
{limit_clause}"
);
let mut params_vec: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(1 + name_filter.len());
params_vec.push(&namespace);
for n in name_filter {
params_vec.push(n);
}
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(
rusqlite::params_from_iter(params_vec.iter().copied()),
|r| r.get::<_, String>(0),
)?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
pub(super) fn read_names_file(path: &Path) -> Result<Vec<String>, AppError> {
let content = std::fs::read_to_string(path).map_err(|e| {
AppError::Validation(format!("failed to read names file {}: {e}", path.display()))
})?;
let mut seen = std::collections::HashSet::new();
let mut out = Vec::new();
for line in content.lines() {
let trimmed = line.trim();
if trimmed.is_empty() || trimmed.starts_with('#') {
continue;
}
if seen.insert(trimmed.to_string()) {
out.push(trimmed.to_string());
}
}
Ok(out)
}
pub(super) fn resolve_name_filter(args: &EnrichArgs) -> Result<Vec<String>, AppError> {
let mut combined: Vec<String> = args.names.clone();
if let Some(p) = &args.names_file {
let from_file = read_names_file(p)?;
for n in from_file {
if !combined.contains(&n) {
combined.push(n);
}
}
}
Ok(combined)
}
pub(super) fn scan_entities_without_description(
conn: &Connection,
namespace: &str,
limit: Option<usize>,
name_filter: &[String],
) -> Result<Vec<(i64, String, String)>, AppError> {
let limit_clause = limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
if name_filter.is_empty() {
let sql = format!(
"SELECT id, name, type
FROM entities
WHERE namespace = ?1
AND {NULL_DESCRIPTION_PREDICATE}
ORDER BY id
{limit_clause}"
);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params![namespace], |r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
))
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
} else {
let placeholders: Vec<String> = (2..=name_filter.len() + 1)
.map(|i| format!("?{i}"))
.collect();
let in_clause = placeholders.join(", ");
let sql = format!(
"SELECT id, name, type
FROM entities
WHERE namespace = ?1
AND name IN ({in_clause})
AND {NULL_DESCRIPTION_PREDICATE}
ORDER BY id
{limit_clause}"
);
let mut params_vec: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(1 + name_filter.len());
params_vec.push(&namespace);
for n in name_filter {
params_vec.push(n);
}
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(
rusqlite::params_from_iter(params_vec.iter().copied()),
|r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
))
},
)?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
}
pub(super) fn scan_short_body_memories(
conn: &Connection,
namespace: &str,
min_chars: usize,
limit: Option<usize>,
name_filter: &[String],
) -> Result<Vec<(i64, String, String)>, AppError> {
let limit_clause = limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
if name_filter.is_empty() {
let sql = format!(
"SELECT m.id, m.name, m.body
FROM memories m
WHERE m.namespace = ?1
AND m.deleted_at IS NULL
AND {SHORT_BODY_PREDICATE}
ORDER BY m.id
{limit_clause}"
);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params![namespace, min_chars as i64], |r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
))
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
} else {
let placeholders: Vec<String> = (3..=name_filter.len() + 2)
.map(|i| format!("?{i}"))
.collect();
let in_clause = placeholders.join(", ");
let sql = format!(
"SELECT m.id, m.name, m.body
FROM memories m
WHERE m.namespace = ?1
AND m.deleted_at IS NULL
AND m.name IN ({in_clause})
AND {SHORT_BODY_PREDICATE}
ORDER BY m.id
{limit_clause}"
);
let mut params_vec: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(2 + name_filter.len());
let min_chars_i64 = min_chars as i64;
params_vec.push(&namespace);
params_vec.push(&min_chars_i64);
for n in name_filter {
params_vec.push(n);
}
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(
rusqlite::params_from_iter(params_vec.iter().copied()),
|r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
))
},
)?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
}
pub(super) fn scan_memories_without_embeddings(
conn: &Connection,
namespace: &str,
limit: Option<usize>,
name_filter: &[String],
) -> Result<Vec<(i64, String, String)>, AppError> {
let limit_clause = limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
let predicate = reembed_memory_predicate(crate::constants::embedding_dim());
if name_filter.is_empty() {
let sql = format!(
"SELECT m.id, m.name, COALESCE(m.body,'')
FROM memories m
WHERE m.namespace = ?1
AND m.deleted_at IS NULL
AND {predicate}
ORDER BY m.id
{limit_clause}"
);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params![namespace], |r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
))
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
} else {
let placeholders: Vec<String> = (2..=name_filter.len() + 1)
.map(|i| format!("?{i}"))
.collect();
let in_clause = placeholders.join(", ");
let sql = format!(
"SELECT m.id, m.name, COALESCE(m.body,'')
FROM memories m
WHERE m.namespace = ?1
AND m.deleted_at IS NULL
AND m.name IN ({in_clause})
AND {predicate}
ORDER BY m.id
{limit_clause}"
);
let mut params_vec: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(1 + name_filter.len());
params_vec.push(&namespace);
for n in name_filter {
params_vec.push(n);
}
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(
rusqlite::params_from_iter(params_vec.iter().copied()),
|r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
))
},
)?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
}
pub(super) fn scan_entities_missing_embeddings(
conn: &Connection,
namespace: &str,
limit: Option<usize>,
name_filter: &[String],
) -> Result<Vec<(i64, String)>, AppError> {
let limit_clause = limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
let predicate = reembed_entity_predicate(crate::constants::embedding_dim());
if name_filter.is_empty() {
let sql = format!(
"SELECT e.id, e.name
FROM entities e
WHERE e.namespace = ?1
AND {predicate}
ORDER BY e.id
{limit_clause}"
);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params![namespace], |r| {
Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?))
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
} else {
let placeholders: Vec<String> = (2..=name_filter.len() + 1)
.map(|i| format!("?{i}"))
.collect();
let in_clause = placeholders.join(", ");
let sql = format!(
"SELECT e.id, e.name
FROM entities e
WHERE e.namespace = ?1
AND e.name IN ({in_clause})
AND {predicate}
ORDER BY e.id
{limit_clause}"
);
let mut params_vec: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(1 + name_filter.len());
params_vec.push(&namespace);
for n in name_filter {
params_vec.push(n);
}
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(
rusqlite::params_from_iter(params_vec.iter().copied()),
|r| Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?)),
)?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
}
pub(super) fn scan_chunks_missing_embeddings(
conn: &Connection,
namespace: &str,
limit: Option<usize>,
name_filter: &[String],
) -> Result<Vec<i64>, AppError> {
let limit_clause = limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
let predicate = reembed_chunk_predicate(crate::constants::embedding_dim());
if name_filter.is_empty() {
let sql = format!(
"SELECT c.id
FROM memory_chunks c
LEFT JOIN memories m ON m.id = c.memory_id
WHERE (m.namespace = ?1 OR m.id IS NULL)
AND {predicate}
ORDER BY c.id
{limit_clause}"
);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params![namespace], |r| r.get::<_, i64>(0))?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
} else {
let placeholders: Vec<String> = (2..=name_filter.len() + 1)
.map(|i| format!("?{i}"))
.collect();
let in_clause = placeholders.join(", ");
let sql = format!(
"SELECT c.id
FROM memory_chunks c
LEFT JOIN memories m ON m.id = c.memory_id
WHERE (m.namespace = ?1 OR m.id IS NULL)
AND (m.name IN ({in_clause}) OR m.id IS NULL)
AND {predicate}
ORDER BY c.id
{limit_clause}"
);
let mut params_vec: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(1 + name_filter.len());
params_vec.push(&namespace);
for n in name_filter {
params_vec.push(n);
}
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(
rusqlite::params_from_iter(params_vec.iter().copied()),
|r| r.get::<_, i64>(0),
)?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
}
#[allow(clippy::type_complexity)]
pub(super) fn scan_weight_candidates(
conn: &Connection,
namespace: &str,
limit: Option<usize>,
) -> Result<Vec<(i64, String, String, String, f64)>, AppError> {
let limit_clause = limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
let sql = format!(
"SELECT r.id, e1.name, e2.name, r.relation, r.weight \
FROM relationships r \
JOIN entities e1 ON e1.id = r.source_id \
JOIN entities e2 ON e2.id = r.target_id \
WHERE {HIGH_WEIGHT_PREDICATE} AND e1.namespace = ?1 \
ORDER BY r.weight DESC {limit_clause}"
);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params![namespace], |r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
r.get::<_, String>(3)?,
r.get::<_, f64>(4)?,
))
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
pub(super) fn scan_generic_relations(
conn: &Connection,
namespace: &str,
limit: Option<usize>,
) -> Result<Vec<(i64, String, String, String)>, AppError> {
let limit_clause = limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
let sql = format!(
"SELECT r.id, e1.name, e2.name, r.relation \
FROM relationships r \
JOIN entities e1 ON e1.id = r.source_id \
JOIN entities e2 ON e2.id = r.target_id \
WHERE {GENERIC_RELATION_PREDICATE} AND e1.namespace = ?1 \
ORDER BY r.id {limit_clause}"
);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt
.query_map(rusqlite::params![namespace], |r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
r.get::<_, String>(3)?,
))
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
pub(super) fn scan_operation(
conn: &Connection,
namespace: &str,
args: &EnrichArgs,
) -> Result<Vec<String>, AppError> {
let name_filter = resolve_name_filter(args)?;
match args.operation() {
EnrichOperation::MemoryBindings => {
let rows = scan_unbound_memories(conn, namespace, args.limit, &name_filter)?;
Ok(rows.into_iter().map(|(_, name, _)| name).collect())
}
EnrichOperation::AugmentBindings => {
scan_bound_memories_for_augment(conn, namespace, args.limit, &name_filter)
}
EnrichOperation::EntityDescriptions => {
let rows =
scan_entities_without_description(conn, namespace, args.limit, &name_filter)?;
Ok(rows.into_iter().map(|(_, name, _)| name).collect())
}
EnrichOperation::BodyEnrich => {
let rows = scan_short_body_memories(
conn,
namespace,
args.min_output_chars,
args.limit,
&name_filter,
)?;
Ok(rows.into_iter().map(|(_, name, _)| name).collect())
}
EnrichOperation::ReEmbed => {
let mut keys: Vec<String> = Vec::new();
if matches!(args.target, ReEmbedTarget::Memories | ReEmbedTarget::All) {
let rows =
scan_memories_without_embeddings(conn, namespace, args.limit, &name_filter)?;
keys.extend(rows.into_iter().map(|(_, name, _)| name));
}
if matches!(args.target, ReEmbedTarget::Entities | ReEmbedTarget::All) {
let rows =
scan_entities_missing_embeddings(conn, namespace, args.limit, &name_filter)?;
keys.extend(rows.into_iter().map(|(_, name)| format!("entity:{name}")));
}
if matches!(args.target, ReEmbedTarget::Chunks | ReEmbedTarget::All) {
let ids =
scan_chunks_missing_embeddings(conn, namespace, args.limit, &name_filter)?;
keys.extend(ids.into_iter().map(|id| format!("chunk:{id}")));
}
Ok(keys)
}
EnrichOperation::WeightCalibrate => {
let rows = scan_weight_candidates(conn, namespace, args.limit)?;
Ok(rows
.into_iter()
.map(|(id, _, _, _, _)| id.to_string())
.collect())
}
EnrichOperation::RelationReclassify => {
let rows = scan_generic_relations(conn, namespace, args.limit)?;
Ok(rows
.into_iter()
.map(|(id, _, _, _)| id.to_string())
.collect())
}
EnrichOperation::EntityConnect | EnrichOperation::CrossDomainBridges => {
let pairs = scan_isolated_entity_pairs(conn, namespace, args.limit)?;
Ok(pairs
.into_iter()
.map(|(id1, _, id2, _)| format_pair_key(id1, id2))
.collect())
}
EnrichOperation::EntityTypeValidate => {
let rows = scan_entities_for_type_validation(conn, namespace, args.limit)?;
Ok(rows.into_iter().map(|(_, name, _)| name).collect())
}
EnrichOperation::DescriptionEnrich => {
let rows = scan_generic_descriptions(conn, namespace, args.limit)?;
Ok(rows.into_iter().map(|(_, name, _)| name).collect())
}
EnrichOperation::DomainClassify
| EnrichOperation::GraphAudit
| EnrichOperation::DeepResearchSynth
| EnrichOperation::BodyExtract => {
let limit_clause = args.limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
let sql = format!(
"SELECT name FROM memories WHERE namespace=?1 AND deleted_at IS NULL ORDER BY id {limit_clause}"
);
let mut stmt = conn.prepare(&sql)?;
let mut names = stmt
.query_map(rusqlite::params![namespace], |r| r.get::<_, String>(0))?
.collect::<Result<Vec<_>, _>>()?;
if !name_filter.is_empty() {
names.retain(|n| name_filter.iter().any(|f| f == n));
}
Ok(names)
}
}
}