use std::collections::{BTreeSet, HashMap};
use crate::base::error::Result;
pub(crate) fn normalize_entity_name(name: &str) -> String {
name.to_lowercase()
}
pub(crate) fn extract_source_turn(metadata: &serde_json::Value) -> Option<i64> {
fn valid_turn(meta: &serde_json::Value, key: &str) -> Option<i64> {
meta.get(key)?.as_i64().filter(|t| *t >= 0)
}
valid_turn(metadata, "source_turn").or_else(|| valid_turn(metadata, "turn_id"))
}
pub(crate) fn decode_canonical_source_turn(
payload: &serde_json::Value,
) -> crate::base::error::Result<Option<Option<i64>>> {
match payload.get("source_turn") {
None => Ok(None),
Some(serde_json::Value::Null) => Ok(Some(None)),
Some(v) => match v.as_i64() {
Some(t) if t >= 0 => Ok(Some(Some(t))),
_ => Err(crate::base::error::YantrikDbError::InvalidInput(format!(
"replicated payload carries a malformed canonical source_turn \
({v}): must be JSON null or a nonnegative integer; rejecting \
the op rather than certifying a divergent column"
))),
},
}
}
pub(crate) const SOURCE_TURN_MARKER_KEY: &str = "source_turn_backfill_complete";
pub(crate) const SOURCE_TURN_EPOCH_KEY: &str = "source_turn_invalidation_epoch";
#[derive(Debug, Clone)]
pub(crate) struct MarkerSnapshot {
marker: Option<String>,
epoch: Option<String>,
}
fn meta_get(conn: &rusqlite::Connection, key: &str) -> rusqlite::Result<Option<String>> {
match conn.query_row(
"SELECT value FROM meta WHERE key = ?1",
rusqlite::params![key],
|r| r.get(0),
) {
Ok(v) => Ok(Some(v)),
Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
Err(e) => Err(e),
}
}
fn meta_put(
conn: &rusqlite::Connection,
key: &str,
value: &Option<String>,
) -> rusqlite::Result<()> {
match value {
Some(v) => {
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES (?1, ?2)",
rusqlite::params![key, v],
)?;
}
None => {
conn.execute("DELETE FROM meta WHERE key = ?1", rusqlite::params![key])?;
}
}
Ok(())
}
pub(crate) fn marker_snapshot(conn: &rusqlite::Connection) -> rusqlite::Result<MarkerSnapshot> {
Ok(MarkerSnapshot {
marker: meta_get(conn, SOURCE_TURN_MARKER_KEY)?,
epoch: meta_get(conn, SOURCE_TURN_EPOCH_KEY)?,
})
}
pub(crate) fn marker_restore(
conn: &rusqlite::Connection,
prior: &MarkerSnapshot,
) -> rusqlite::Result<()> {
meta_put(conn, SOURCE_TURN_MARKER_KEY, &prior.marker)?;
meta_put(conn, SOURCE_TURN_EPOCH_KEY, &prior.epoch)?;
Ok(())
}
pub(crate) fn repair_entity_norm(
conn: &rusqlite::Connection,
memory_rid: &str,
entity_name: &str,
) -> rusqlite::Result<()> {
let norm = normalize_entity_name(entity_name);
conn.execute(
"UPDATE memory_entities SET entity_name_norm = ?3 \
WHERE memory_rid = ?1 AND entity_name = ?2 \
AND (entity_name_norm IS NULL OR entity_name_norm != ?3)",
rusqlite::params![memory_rid, entity_name, norm],
)?;
Ok(())
}
#[derive(Debug, Clone, PartialEq, serde::Serialize)]
pub struct ThreadItem {
pub rid: String,
pub text: String,
pub created_at: f64,
pub source_turn: Option<i64>,
pub position: usize,
pub entities: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, serde::Serialize)]
pub struct ThreadItemV2 {
pub rid: String,
pub text: String,
pub created_at: f64,
pub source_turn: Option<i64>,
pub position: usize,
pub entities: Vec<String>,
pub routes: Vec<&'static str>,
pub phrases: Vec<String>,
pub topic_rids: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, serde::Serialize)]
pub struct ThreadRecallV2 {
pub items: Vec<ThreadItemV2>,
pub total: usize,
pub returned: usize,
pub omitted: usize,
}
#[derive(Debug, Clone, Default, PartialEq, serde::Serialize)]
pub struct ThreadQuery {
pub entities: Vec<String>,
pub phrases: Vec<String>,
pub topic_rids: Vec<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
pub struct MaintenanceProgress {
pub processed: usize,
pub remaining: usize,
pub complete: bool,
}
#[derive(Debug, Clone, PartialEq, serde::Serialize)]
pub struct ThreadRecall {
pub items: Vec<ThreadItem>,
pub total: usize,
pub omitted: usize,
}
pub(crate) const SOURCE_TURN_CURSOR_KEY: &str = "source_turn_repair_cursor";
pub(crate) fn source_turn_repair_batch<F>(
conn: &rusqlite::Connection,
decrypt: F,
batch: i64,
) -> Result<MaintenanceProgress>
where
F: Fn(&str) -> Result<String>,
{
let tx = conn.unchecked_transaction()?;
let epoch_start = meta_get(&tx, SOURCE_TURN_EPOCH_KEY)?.unwrap_or_else(|| "0".to_string());
let cursor: i64 = match meta_get(&tx, SOURCE_TURN_CURSOR_KEY)? {
Some(stored) => match stored.split_once(':') {
Some((epoch, rowid)) if epoch == epoch_start => rowid.parse::<i64>().unwrap_or(0),
_ => 0, },
None => 0,
};
let candidates: Vec<(i64, Option<String>, Option<i64>)> = {
let mut stmt = tx.prepare(
"SELECT rowid, metadata, source_turn FROM memories \
WHERE rowid > ?1 ORDER BY rowid LIMIT ?2",
)?;
let rows = stmt.query_map(rusqlite::params![cursor, batch], |r| {
Ok((r.get(0)?, r.get(1)?, r.get(2)?))
})?;
rows.collect::<std::result::Result<Vec<_>, _>>()?
};
let mut processed = 0usize;
let mut last = cursor;
for (rowid, stored_meta, stored_turn) in &candidates {
let expected: Option<i64> = match stored_meta.as_deref() {
None | Some("") => None,
Some(stored) => {
let plain = match decrypt(stored) {
Ok(p) => p,
Err(decrypt_err) => {
if serde_json::from_str::<serde_json::Value>(stored).is_ok() {
stored.to_string()
} else {
return Err(decrypt_err);
}
}
};
serde_json::from_str::<serde_json::Value>(&plain)
.ok()
.as_ref()
.and_then(extract_source_turn)
}
};
if expected != *stored_turn {
tx.execute(
"UPDATE memories SET source_turn = ?1 WHERE rowid = ?2",
rusqlite::params![expected, rowid],
)?;
}
processed += 1;
last = *rowid;
}
let remaining: i64 = tx.query_row(
"SELECT COUNT(*) FROM memories WHERE rowid > ?1",
rusqlite::params![last],
|r| r.get(0),
)?;
let complete = remaining == 0;
if complete {
meta_put(&tx, SOURCE_TURN_MARKER_KEY, &Some("1".to_string()))?;
meta_put(&tx, SOURCE_TURN_CURSOR_KEY, &None)?;
} else if processed > 0 {
let epoch_now = meta_get(&tx, SOURCE_TURN_EPOCH_KEY)?.unwrap_or_else(|| "0".to_string());
meta_put(
&tx,
SOURCE_TURN_CURSOR_KEY,
&Some(format!("{epoch_now}:{last}")),
)?;
}
tx.commit()?;
let remaining = usize::try_from(remaining).map_err(|_| {
crate::base::error::YantrikDbError::InvalidInput(format!(
"source_turn repair: remaining count {remaining} does not fit in usize"
))
})?;
Ok(MaintenanceProgress {
processed,
remaining,
complete,
})
}
impl crate::YantrikDB {
pub fn recall_thread(
&self,
namespace: &str,
entities: &[&str],
limit: usize,
) -> Result<ThreadRecall> {
let empty = || ThreadRecall {
items: Vec::new(),
total: 0,
omitted: 0,
};
if entities.is_empty() {
return Ok(empty());
}
let mut req_by_lower: HashMap<String, usize> = HashMap::new();
for (i, e) in entities.iter().enumerate() {
req_by_lower.entry(normalize_entity_name(e)).or_insert(i);
}
struct RowAgg {
text: String,
created_at: f64,
metadata: Option<String>,
matched: BTreeSet<usize>,
}
let mut by_rid: HashMap<String, RowAgg> = HashMap::new();
{
let conn = self.read_conn();
let mut norm_names: Vec<&String> = req_by_lower.keys().collect();
norm_names.sort();
let placeholders: String = (0..norm_names.len())
.map(|i| format!("?{}", i + 2))
.collect::<Vec<_>>()
.join(",");
let sql = format!(
"SELECT m.rid, m.text, m.created_at, m.metadata, me.entity_name_norm \
FROM memories m \
JOIN memory_entities me ON me.memory_rid = m.rid \
WHERE m.namespace = ?1 \
AND m.consolidation_status = 'active' \
AND (m.synthesis_state IS NULL OR m.synthesis_state = 'verified') \
AND me.entity_name_norm IN ({placeholders})"
);
let mut param_values: Vec<Box<dyn rusqlite::types::ToSql>> =
Vec::with_capacity(norm_names.len() + 1);
param_values.push(Box::new(namespace.to_string()));
for name in &norm_names {
param_values.push(Box::new((*name).clone()));
}
let params_ref: Vec<&dyn rusqlite::types::ToSql> =
param_values.iter().map(|p| p.as_ref()).collect();
let mut stmt = conn.prepare(&sql)?;
let rows = stmt.query_map(params_ref.as_slice(), |r| {
Ok((
r.get::<_, String>(0)?,
r.get::<_, String>(1)?,
r.get::<_, f64>(2)?,
r.get::<_, Option<String>>(3)?,
r.get::<_, String>(4)?,
))
})?;
for row in rows {
let (rid, text, created_at, metadata, entity_norm) = row?;
let idx = req_by_lower[&entity_norm];
by_rid
.entry(rid)
.or_insert_with(|| RowAgg {
text,
created_at,
metadata,
matched: BTreeSet::new(),
})
.matched
.insert(idx);
}
}
if self.status_read_policy() && !by_rid.is_empty() {
let rids: Vec<String> = by_rid.keys().cloned().collect();
let rid_refs: Vec<&str> = rids.iter().map(String::as_str).collect();
let superseded = self.superseded_rids_among(&rid_refs)?;
for rid in superseded {
by_rid.remove(&rid);
}
}
let mut eligible: Vec<(String, String, f64, Option<i64>, BTreeSet<usize>)> =
Vec::with_capacity(by_rid.len());
for (rid, agg) in by_rid {
let text = self.decrypt_text(&agg.text)?;
let source_turn = match agg.metadata.as_deref() {
None | Some("") => None,
Some(stored_meta) => {
let metadata = self.decrypt_text(stored_meta)?;
serde_json::from_str::<serde_json::Value>(&metadata)
.ok()
.as_ref()
.and_then(extract_source_turn)
}
};
eligible.push((rid, text, agg.created_at, source_turn, agg.matched));
}
eligible.sort_by(|a, b| {
a.2.total_cmp(&b.2)
.then_with(|| match (a.3, b.3) {
(Some(x), Some(y)) => x.cmp(&y),
(Some(_), None) => std::cmp::Ordering::Less,
(None, Some(_)) => std::cmp::Ordering::Greater,
(None, None) => std::cmp::Ordering::Equal,
})
.then_with(|| a.0.cmp(&b.0))
});
let total = eligible.len();
let mut items = Vec::with_capacity(total.min(limit));
for (pos0, (rid, text, created_at, source_turn, matched)) in
eligible.into_iter().enumerate()
{
if pos0 >= limit {
break; }
items.push(ThreadItem {
rid,
text,
created_at,
source_turn,
position: pos0 + 1,
entities: matched
.into_iter()
.map(|i| entities[i].to_string())
.collect(),
});
}
let omitted = total - items.len();
Ok(ThreadRecall {
items,
total,
omitted,
})
}
pub fn recall_thread_v2(
&self,
namespace: &str,
query: &ThreadQuery,
limit: usize,
) -> Result<ThreadRecallV2> {
use crate::base::error::YantrikDbError;
const MAX_ENTITIES: usize = 64;
const MAX_PHRASES: usize = 32;
const MAX_TOPIC_RIDS: usize = 16;
const MAX_ITEM_BYTES: usize = 512;
const MAX_LIMIT: usize = 10_000;
fn check_items(kind: &str, items: &[String], max_count: usize) -> Result<()> {
if items.len() > max_count {
return Err(YantrikDbError::InvalidInput(format!(
"recall_thread_v2: {} {kind} items exceed the cap of {max_count}",
items.len()
)));
}
for item in items {
if item.is_empty() || item.len() > MAX_ITEM_BYTES {
return Err(YantrikDbError::InvalidInput(format!(
"recall_thread_v2: every {kind} item must be 1..={MAX_ITEM_BYTES} bytes"
)));
}
}
Ok(())
}
check_items("entity", &query.entities, MAX_ENTITIES)?;
check_items("phrase", &query.phrases, MAX_PHRASES)?;
check_items("topic_rid", &query.topic_rids, MAX_TOPIC_RIDS)?;
if limit > MAX_LIMIT {
return Err(YantrikDbError::InvalidInput(format!(
"recall_thread_v2: limit {limit} exceeds the cap of {MAX_LIMIT}"
)));
}
let limit_i64 = i64::try_from(limit).map_err(|_| {
YantrikDbError::InvalidInput(format!(
"recall_thread_v2: limit {limit} does not fit in i64"
))
})?;
let empty = || ThreadRecallV2 {
items: Vec::new(),
total: 0,
returned: 0,
omitted: 0,
};
if query.entities.is_empty() && query.phrases.is_empty() && query.topic_rids.is_empty() {
return Ok(empty());
}
if !query.phrases.is_empty() && self.is_encrypted() {
return Err(YantrikDbError::CapabilityUnavailable {
capability: "phrase_thread_route".to_string(),
reason: "the memories_fts index on an encrypted store holds ciphertext, so \
a phrase MATCH would silently match nothing; use the entity or \
topic routes, or an unencrypted store"
.to_string(),
});
}
let mut req_by_norm: HashMap<String, usize> = HashMap::new();
for (i, e) in query.entities.iter().enumerate() {
req_by_norm.entry(normalize_entity_name(e)).or_insert(i);
}
let mut norm_names: Vec<String> = req_by_norm.keys().cloned().collect();
norm_names.sort();
fn dedup_first<'a>(items: &'a [String]) -> Vec<&'a str> {
let mut seen = std::collections::HashSet::new();
items
.iter()
.filter(|s| seen.insert(s.as_str()))
.map(String::as_str)
.collect()
}
let uniq_phrases: Vec<&str> = dedup_first(&query.phrases);
let uniq_topics: Vec<&str> = dedup_first(&query.topic_rids);
fn fts_literal(phrase: &str) -> String {
format!("\"{}\"", phrase.replace('"', "\"\""))
}
let fts_expr: Option<String> = if uniq_phrases.is_empty() {
None
} else {
Some(
uniq_phrases
.iter()
.map(|p| fts_literal(p))
.collect::<Vec<_>>()
.join(" OR "),
)
};
let mut params_v: Vec<Box<dyn rusqlite::types::ToSql>> =
vec![Box::new(namespace.to_string())];
let mut union_parts: Vec<String> = Vec::new();
if !norm_names.is_empty() {
let ph: Vec<String> = norm_names
.iter()
.map(|n| {
params_v.push(Box::new(n.clone()));
format!("?{}", params_v.len())
})
.collect();
union_parts.push(format!(
"SELECT DISTINCT me.memory_rid AS rid FROM memory_entities me \
WHERE me.entity_name_norm IN ({})",
ph.join(",")
));
}
if let Some(expr) = &fts_expr {
params_v.push(Box::new(expr.clone()));
union_parts.push(format!(
"SELECT DISTINCT fm.rid AS rid FROM memories fm \
JOIN memories_fts ON memories_fts.rowid = fm.rowid \
WHERE memories_fts MATCH ?{}",
params_v.len()
));
}
if !uniq_topics.is_empty() {
let ph: Vec<String> = uniq_topics
.iter()
.map(|t| {
params_v.push(Box::new((*t).to_string()));
format!("?{}", params_v.len())
})
.collect();
union_parts.push(format!(
"SELECT DISTINCT d.source_rid AS rid FROM synthesis_dependencies d \
WHERE d.namespace = ?1 AND d.is_direct = 1 \
AND d.synthesis_rid IN ({})",
ph.join(",")
));
}
let union_sql = union_parts.join(" UNION ");
let mut vis = String::from(
"m.namespace = ?1 AND m.consolidation_status = 'active' \
AND (m.synthesis_state IS NULL OR m.synthesis_state = 'verified')",
);
if self.status_read_policy() {
vis.push_str(
" AND NOT EXISTS (SELECT 1 FROM record_links rl \
WHERE rl.link_type = 'supersedes' AND rl.status = 'active' \
AND rl.selection_state = 'selected' AND rl.target_rid = m.rid)",
);
}
let params_ref: Vec<&dyn rusqlite::types::ToSql> =
params_v.iter().map(|p| p.as_ref()).collect();
let conn = self.read_conn();
let tx = conn.unchecked_transaction()?;
let topic_vis = format!(
"{vis} AND m.synthesis_state = 'verified' \
AND m.synthesis_axis IS NOT NULL \
AND m.synthesis_granularity IN ('atomic', 'rollup') \
AND m.synthesis_logical_key IS NOT NULL"
);
for topic_rid in &uniq_topics {
let visible: bool = tx.query_row(
&format!(
"SELECT EXISTS(SELECT 1 FROM memories m WHERE m.rid = ?2 AND {topic_vis})"
),
rusqlite::params![namespace, topic_rid],
|r| r.get(0),
)?;
if !visible {
return Err(YantrikDbError::InvalidThreadTopic {
topic_rid: (*topic_rid).to_string(),
});
}
}
let total: i64 = tx.query_row(
&format!(
"SELECT COUNT(*) FROM ({union_sql}) u \
JOIN memories m ON m.rid = u.rid WHERE {vis}"
),
params_ref.as_slice(),
|r| r.get(0),
)?;
let total_usize = usize::try_from(total).map_err(|_| {
YantrikDbError::InvalidInput(format!(
"recall_thread_v2: eligible count {total} does not fit in usize"
))
})?;
if total_usize == 0 {
return Ok(empty());
}
{
let marker = meta_get(&tx, SOURCE_TURN_MARKER_KEY)?;
if marker.as_deref() != Some("1") {
return Err(YantrikDbError::MaintenanceRequired {
operation: "maintain_source_turn_backfill".to_string(),
reason: "this store's source_turn columns are not known to mirror \
their metadata (backfill incomplete, or a raw SQL write \
staled them), so the strict (created_at, source_turn, rid) \
thread order cannot be guaranteed; call \
maintain_source_turn_backfill in batches until it reports \
complete"
.to_string(),
});
}
}
let page_sql = format!(
"SELECT rid, text, created_at, source_turn, pos FROM ( \
SELECT m.rid AS rid, m.text AS text, m.created_at AS created_at, \
m.source_turn AS source_turn, \
ROW_NUMBER() OVER (ORDER BY m.created_at ASC, \
(m.source_turn IS NULL) ASC, m.source_turn ASC, \
m.rid ASC) AS pos \
FROM ({union_sql}) u JOIN memories m ON m.rid = u.rid \
WHERE {vis} \
) ORDER BY pos ASC LIMIT ?{}",
params_v.len() + 1
);
struct PageRow {
rid: String,
stored_text: String,
created_at: f64,
source_turn: Option<i64>,
position: usize,
}
let mut page_params: Vec<&dyn rusqlite::types::ToSql> = params_ref.clone();
page_params.push(&limit_i64);
let page: Vec<PageRow> = {
let mut stmt = tx.prepare(&page_sql)?;
let rows = stmt.query_map(page_params.as_slice(), |r| {
Ok((
r.get::<_, String>(0)?,
r.get::<_, String>(1)?,
r.get::<_, f64>(2)?,
r.get::<_, Option<i64>>(3)?,
r.get::<_, i64>(4)?,
))
})?;
rows.map(|row| {
let (rid, stored_text, created_at, source_turn, pos) = row?;
let position = usize::try_from(pos).map_err(|_| {
YantrikDbError::InvalidInput(format!(
"recall_thread_v2: position {pos} does not fit in usize"
))
})?;
Ok(PageRow {
rid,
stored_text,
created_at,
source_turn,
position,
})
})
.collect::<Result<Vec<_>>>()?
};
let page_rids: Vec<&str> = page.iter().map(|p| p.rid.as_str()).collect();
let mut ent_matches: HashMap<String, BTreeSet<usize>> = HashMap::new();
let mut phrase_matches: HashMap<String, BTreeSet<usize>> = HashMap::new();
let mut topic_matches: HashMap<String, BTreeSet<usize>> = HashMap::new();
if !page_rids.is_empty() {
if !norm_names.is_empty() {
let mut p: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
let n_ph: Vec<String> = norm_names
.iter()
.map(|n| {
p.push(Box::new(n.clone()));
format!("?{}", p.len())
})
.collect();
let r_ph: Vec<String> = page_rids
.iter()
.map(|r| {
p.push(Box::new((*r).to_string()));
format!("?{}", p.len())
})
.collect();
let pr: Vec<&dyn rusqlite::types::ToSql> = p.iter().map(|b| b.as_ref()).collect();
let mut stmt = tx.prepare(&format!(
"SELECT me.memory_rid, me.entity_name_norm FROM memory_entities me \
WHERE me.entity_name_norm IN ({}) AND me.memory_rid IN ({})",
n_ph.join(","),
r_ph.join(",")
))?;
let rows = stmt.query_map(pr.as_slice(), |r| {
Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))
})?;
for row in rows {
let (rid, norm) = row?;
if let Some(idx) = req_by_norm.get(&norm) {
ent_matches.entry(rid).or_default().insert(*idx);
}
}
}
for (phrase_idx, phrase) in uniq_phrases.iter().enumerate() {
let literal = fts_literal(phrase);
let mut p: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
p.push(Box::new(literal));
let r_ph: Vec<String> = page_rids
.iter()
.map(|r| {
p.push(Box::new((*r).to_string()));
format!("?{}", p.len())
})
.collect();
let pr: Vec<&dyn rusqlite::types::ToSql> = p.iter().map(|b| b.as_ref()).collect();
let mut stmt = tx.prepare(&format!(
"SELECT fm.rid FROM memories fm \
JOIN memories_fts ON memories_fts.rowid = fm.rowid \
WHERE memories_fts MATCH ?1 AND fm.rid IN ({})",
r_ph.join(",")
))?;
let rows = stmt.query_map(pr.as_slice(), |r| r.get::<_, String>(0))?;
for row in rows {
phrase_matches.entry(row?).or_default().insert(phrase_idx);
}
}
if !uniq_topics.is_empty() {
let mut p: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
p.push(Box::new(namespace.to_string()));
let t_ph: Vec<String> = uniq_topics
.iter()
.map(|t| {
p.push(Box::new((*t).to_string()));
format!("?{}", p.len())
})
.collect();
let r_ph: Vec<String> = page_rids
.iter()
.map(|r| {
p.push(Box::new((*r).to_string()));
format!("?{}", p.len())
})
.collect();
let pr: Vec<&dyn rusqlite::types::ToSql> = p.iter().map(|b| b.as_ref()).collect();
let mut stmt = tx.prepare(&format!(
"SELECT d.synthesis_rid, d.source_rid FROM synthesis_dependencies d \
WHERE d.namespace = ?1 AND d.is_direct = 1 \
AND d.synthesis_rid IN ({}) AND d.source_rid IN ({})",
t_ph.join(","),
r_ph.join(",")
))?;
let rows = stmt.query_map(pr.as_slice(), |r| {
Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))
})?;
for row in rows {
let (synthesis_rid, source_rid) = row?;
if let Some(topic_idx) = uniq_topics.iter().position(|t| *t == synthesis_rid) {
topic_matches
.entry(source_rid)
.or_default()
.insert(topic_idx);
}
}
}
}
drop(tx);
drop(conn);
let mut items = Vec::with_capacity(page.len());
for row in page {
let text = self.decrypt_text(&row.stored_text)?;
let ent_idx = ent_matches.get(&row.rid);
let phrase_idx = phrase_matches.get(&row.rid);
let topic_idx = topic_matches.get(&row.rid);
let mut routes: Vec<&'static str> = Vec::new();
if ent_idx.is_some_and(|s| !s.is_empty()) {
routes.push("entity");
}
if phrase_idx.is_some_and(|s| !s.is_empty()) {
routes.push("phrase");
}
if topic_idx.is_some_and(|s| !s.is_empty()) {
routes.push("topic");
}
let entities: Vec<String> = ent_idx
.map(|s| s.iter().map(|i| query.entities[*i].clone()).collect())
.unwrap_or_default();
let phrases: Vec<String> = phrase_idx
.map(|s| s.iter().map(|i| uniq_phrases[*i].to_string()).collect())
.unwrap_or_default();
let topic_rids: Vec<String> = topic_idx
.map(|s| s.iter().map(|i| uniq_topics[*i].to_string()).collect())
.unwrap_or_default();
items.push(ThreadItemV2 {
rid: row.rid,
text,
created_at: row.created_at,
source_turn: row.source_turn,
position: row.position,
entities,
routes,
phrases,
topic_rids,
});
}
let returned = items.len();
let omitted = total_usize - returned;
Ok(ThreadRecallV2 {
items,
total: total_usize,
returned,
omitted,
})
}
pub fn maintain_source_turn_backfill(&self, batch: usize) -> Result<MaintenanceProgress> {
use crate::base::error::YantrikDbError;
const MAX_MAINTENANCE_BATCH: usize = 10_000;
if batch == 0 {
return Err(YantrikDbError::InvalidInput(
"maintain_source_turn_backfill: batch must be >= 1 (a zero batch can \
never make progress)"
.to_string(),
));
}
if batch > MAX_MAINTENANCE_BATCH {
return Err(YantrikDbError::InvalidInput(format!(
"maintain_source_turn_backfill: batch {batch} exceeds the cap of \
{MAX_MAINTENANCE_BATCH}"
)));
}
let batch_i64 = i64::try_from(batch).map_err(|_| {
YantrikDbError::InvalidInput(format!(
"maintain_source_turn_backfill: batch {batch} does not fit in i64"
))
})?;
let conn = self.conn();
source_turn_repair_batch(&conn, |stored| self.decrypt_text(stored), batch_i64)
}
}
#[cfg(test)]
mod thread_tests {
use crate::YantrikDB;
const NS: &str = "n";
const BASE_MICROS: i64 = 1_700_000_000_000_000;
fn vec_seed(seed: f32, dim: usize) -> Vec<f32> {
let raw: Vec<f32> = (0..dim).map(|i| (seed + i as f32) * 0.1).collect();
let norm: f32 = raw.iter().map(|x| x * x).sum::<f32>().sqrt();
raw.iter().map(|x| x / norm).collect()
}
#[allow(clippy::too_many_arguments)]
fn seed_row(
db: &YantrikDB,
rid: &str,
text: &str,
metadata: &serde_json::Value,
created_at_micros: i64,
entity_names: &[&str],
seed: f32,
) {
db.record_with_rid(
rid,
text,
"episodic",
0.5,
0.0,
604800.0,
metadata,
&vec_seed(seed, 8),
NS,
0.8,
"general",
"user",
None,
created_at_micros,
entity_names,
"test-model.v1",
None,
crate::provenance::WriteAdmission::Admitted,
)
.unwrap();
}
fn drain(db: &YantrikDB) {
for _ in 0..50 {
if db.apply_pending_ops_once(500).unwrap() == 0 {
return;
}
}
panic!("pending ops did not drain");
}
fn meta_empty() -> serde_json::Value {
serde_json::json!({})
}
#[test]
fn coverage_pin_returns_every_thread_member_in_order() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let mut alpha_rids = Vec::new();
for i in 0..120 {
let rid = format!("r{i:03}");
let micros = BASE_MICROS + (i as i64) * 1_000_000;
if i % 2 == 0 {
seed_row(
&db,
&rid,
&format!("Alpha update number {i}"),
&meta_empty(),
micros,
&["Alpha"],
i as f32,
);
alpha_rids.push(rid);
} else {
seed_row(
&db,
&rid,
&format!("Beta update number {i}"),
&meta_empty(),
micros,
&["Beta"],
i as f32,
);
}
}
drain(&db);
let out = db.recall_thread(NS, &["alpha"], 100).unwrap();
assert_eq!(out.total, 60, "eligible set is ALL 60 alpha rows");
assert_eq!(out.omitted, 0);
assert_eq!(out.items.len(), 60);
for (i, item) in out.items.iter().enumerate() {
assert_eq!(item.rid, alpha_rids[i], "chronological (insertion) order");
assert_eq!(item.position, i + 1, "positions 1..=60");
assert_eq!(item.entities, vec!["alpha".to_string()]);
assert!(item.text.contains("Alpha update"), "decrypted text");
if i > 0 {
assert!(
out.items[i - 1].created_at < item.created_at,
"created_at strictly ascending"
);
}
}
}
#[test]
fn turn_tie_break_within_equal_created_at() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let t = BASE_MICROS;
seed_row(
&db,
"m_turn7",
"Gamma event seven",
&serde_json::json!({"source_turn": 7}),
t,
&["Gamma"],
1.0,
);
seed_row(
&db,
"a_turn3",
"Gamma event three",
&serde_json::json!({"turn_id": 3}),
t,
&["Gamma"],
2.0,
);
seed_row(
&db,
"a_none",
"Gamma event with no turn",
&meta_empty(),
t,
&["Gamma"],
3.0,
);
seed_row(
&db,
"z_neg",
"Gamma event negative turn",
&serde_json::json!({"source_turn": -2}),
t,
&["Gamma"],
4.0,
);
seed_row(
&db,
"b_str",
"Gamma event string turn",
&serde_json::json!({"source_turn": "5"}),
t,
&["Gamma"],
5.0,
);
seed_row(
&db,
"later_turn0",
"Gamma event later",
&serde_json::json!({"source_turn": 0}),
t + 1_000_000,
&["Gamma"],
6.0,
);
drain(&db);
let out = db.recall_thread(NS, &["Gamma"], 10).unwrap();
let rids: Vec<&str> = out.items.iter().map(|i| i.rid.as_str()).collect();
assert_eq!(
rids,
vec![
"a_turn3",
"m_turn7",
"a_none",
"b_str",
"z_neg",
"later_turn0"
],
"turn asc, then NULLS LAST by rid, then created_at dominates"
);
let turns: Vec<Option<i64>> = out.items.iter().map(|i| i.source_turn).collect();
assert_eq!(
turns,
vec![Some(3), Some(7), None, None, None, Some(0)],
"invalid turns are None — never invented"
);
assert_eq!(
out.items.iter().map(|i| i.position).collect::<Vec<_>>(),
vec![1, 2, 3, 4, 5, 6]
);
}
#[test]
fn truncation_keeps_earliest_and_reports_omitted() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..10 {
seed_row(
&db,
&format!("t{i}"),
&format!("Alpha step {i}"),
&meta_empty(),
BASE_MICROS + (i as i64) * 1_000_000,
&["Alpha"],
i as f32,
);
}
drain(&db);
let out = db.recall_thread(NS, &["Alpha"], 4).unwrap();
assert_eq!(out.total, 10);
assert_eq!(out.omitted, 6);
assert_eq!(
out.items.iter().map(|i| i.rid.as_str()).collect::<Vec<_>>(),
vec!["t0", "t1", "t2", "t3"],
"the earliest are kept, never a similarity sample"
);
assert_eq!(
out.items.iter().map(|i| i.position).collect::<Vec<_>>(),
vec![1, 2, 3, 4],
"full-thread numbering"
);
}
#[test]
fn multi_entity_rows_appear_once_with_matched_subset() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_row(
&db,
"d1",
"Alpha alone here",
&meta_empty(),
BASE_MICROS,
&["Alpha"],
1.0,
);
seed_row(
&db,
"d2",
"Alpha met Beta today",
&meta_empty(),
BASE_MICROS + 1_000_000,
&["Alpha", "Beta"],
2.0,
);
seed_row(
&db,
"d3",
"Beta alone here",
&meta_empty(),
BASE_MICROS + 2_000_000,
&["Beta"],
3.0,
);
drain(&db);
let out = db.recall_thread(NS, &["Alpha", "Beta"], 10).unwrap();
assert_eq!(out.total, 3);
assert_eq!(out.items.len(), 3, "d2 appears exactly once");
assert_eq!(out.items[0].entities, vec!["Alpha".to_string()]);
assert_eq!(
out.items[1].entities,
vec!["Alpha".to_string(), "Beta".to_string()]
);
assert_eq!(out.items[2].entities, vec!["Beta".to_string()]);
let dup = db.recall_thread(NS, &["Alpha", "alpha"], 10).unwrap();
assert_eq!(dup.items[0].entities, vec!["Alpha".to_string()]);
let none = db.recall_thread(NS, &[], 10).unwrap();
assert_eq!((none.total, none.omitted, none.items.len()), (0, 0, 0));
let unknown = db.recall_thread(NS, &["Nobody"], 10).unwrap();
assert_eq!(unknown.total, 0);
assert_eq!(
db.recall_thread("other_ns", &["Alpha"], 10).unwrap().total,
0
);
}
#[test]
fn tombstoned_and_superseded_rows_are_excluded() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for (i, rid) in ["v1", "v2", "v3"].iter().enumerate() {
seed_row(
&db,
rid,
&format!("Alpha visibility row {i}"),
&meta_empty(),
BASE_MICROS + (i as i64) * 1_000_000,
&["Alpha"],
i as f32,
);
}
drain(&db);
assert_eq!(db.recall_thread(NS, &["Alpha"], 10).unwrap().total, 3);
assert!(db.forget("v2").unwrap());
let out = db.recall_thread(NS, &["Alpha"], 10).unwrap();
assert_eq!(out.total, 2, "forgotten row excluded");
assert_eq!(
out.items.iter().map(|i| i.rid.as_str()).collect::<Vec<_>>(),
vec!["v1", "v3"]
);
assert_eq!(
out.items.iter().map(|i| i.position).collect::<Vec<_>>(),
vec![1, 2],
"positions renumber over the eligible thread"
);
assert!(db.status_read_policy());
db.link(
"v3",
&crate::types::RecordLink {
target_rid: "v1".to_string(),
link_type: crate::types::LinkType::Supersedes,
},
)
.unwrap();
let out = db.recall_thread(NS, &["Alpha"], 10).unwrap();
assert_eq!(out.total, 1, "superseded row excluded");
assert_eq!(out.items[0].rid, "v3");
assert_eq!(out.items[0].position, 1);
}
#[test]
fn replication_parity_on_follower() {
use crate::replication::{apply_ops, extract_ops_since};
let leader = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..6 {
let (name, tag) = if i % 2 == 0 {
("Alpha", "milestone")
} else {
("Beta", "sync")
};
seed_row(
&leader,
&format!("p{i}"),
&format!("{name} {tag} number {i}"),
&serde_json::json!({"source_turn": i}),
BASE_MICROS + (i as i64) * 1_000_000,
&[name],
i as f32,
);
}
leader.relate("Alpha", "Beta", "related_to", 1.0).unwrap();
drain(&leader);
let leader_out = leader.recall_thread(NS, &["Alpha"], 100).unwrap();
assert_eq!(leader_out.total, 3, "leader thread complete");
let follower = YantrikDB::new(":memory:", 8).unwrap();
let ops = extract_ops_since(&leader.conn(), None, None, None, 1000).unwrap();
apply_ops(&follower, &ops).unwrap();
let follower_out = follower.recall_thread(NS, &["Alpha"], 100).unwrap();
assert_eq!(
follower_out, leader_out,
"same thread on both sides — items, positions, turns, totals"
);
assert_eq!(
follower.recall_thread(NS, &["Alpha", "Beta"], 100).unwrap(),
leader.recall_thread(NS, &["Alpha", "Beta"], 100).unwrap()
);
}
#[test]
fn every_writer_stamps_the_normalized_entity_key() {
use crate::engine::thread::normalize_entity_name;
use crate::replication::{apply_ops, extract_ops_since};
let leader = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..4 {
seed_row(
&leader,
&format!("c{i}"),
&format!("Münster planning with Alpha and Beta round {i}"),
&meta_empty(),
BASE_MICROS + (i as i64) * 1_000_000,
&["Münster", "Alpha"],
i as f32,
);
}
leader.relate("Alpha", "Beta", "related_to", 1.0).unwrap();
leader.relate("Münster", "Beta", "located_in", 0.7).unwrap();
drain(&leader);
let follower = YantrikDB::new(":memory:", 8).unwrap();
let ops = extract_ops_since(&leader.conn(), None, None, None, 1000).unwrap();
apply_ops(&follower, &ops).unwrap();
let census = |db: &YantrikDB, side: &str| {
let conn = db.conn();
let rows: Vec<(String, Option<String>)> = conn
.prepare("SELECT entity_name, entity_name_norm FROM memory_entities")
.unwrap()
.query_map([], |r| Ok((r.get(0)?, r.get(1)?)))
.unwrap()
.collect::<std::result::Result<Vec<_>, _>>()
.unwrap();
assert!(
!rows.is_empty(),
"{side}: census precondition — the natural paths must have \
produced memory_entities rows"
);
let bad = rows
.iter()
.filter(|(name, norm)| {
norm.as_deref() != Some(normalize_entity_name(name).as_str())
})
.count();
assert_eq!(
bad, 0,
"{side}: a writer inserted memory_entities without the normalized \
key; every writer must bind normalize_entity_name()"
);
};
census(&leader, "leader");
census(&follower, "follower");
}
#[test]
fn non_ascii_entity_resolves_case_insensitively_via_indexed_path() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_row(
&db,
"m_de",
"Planning the MÜNSTER rollout",
&meta_empty(),
BASE_MICROS,
&["MÜNSTER"],
1.0,
);
drain(&db);
let out = db.recall_thread(NS, &["münster"], 10).unwrap();
assert_eq!(
out.total, 1,
"Unicode fold must match, ASCII fold would miss"
);
assert_eq!(out.items[0].rid, "m_de");
assert_eq!(out.items[0].position, 1);
assert_eq!(out.items[0].entities, vec!["münster".to_string()]);
}
#[test]
fn natural_writes_self_heal_a_null_normalized_key() {
use crate::engine::thread::normalize_entity_name;
let db = YantrikDB::new(":memory:", 8).unwrap();
db.conn()
.execute(
"INSERT INTO memory_entities (memory_rid, entity_name) \
VALUES ('m_pre', 'Münster')",
[],
)
.unwrap();
seed_row(
&db,
"m_pre",
"Münster status update",
&meta_empty(),
BASE_MICROS,
&["Münster"],
1.0,
);
drain(&db);
db.conn()
.execute(
"INSERT INTO memory_entities (memory_rid, entity_name) \
VALUES ('m_pre2', 'Alpha')",
[],
)
.unwrap();
db.link_memory_entity("m_pre2", "Alpha").unwrap();
let norm = |rid: &str, name: &str| -> Option<String> {
db.conn()
.query_row(
"SELECT entity_name_norm FROM memory_entities \
WHERE memory_rid = ?1 AND entity_name = ?2",
[rid, name],
|r| r.get(0),
)
.unwrap()
};
assert_eq!(
norm("m_pre", "Münster").as_deref(),
Some(normalize_entity_name("Münster").as_str()),
"the record materializer must repair a pre-existing NULL norm"
);
assert_eq!(
norm("m_pre2", "Alpha").as_deref(),
Some("alpha"),
"link_memory_entity must repair a pre-existing NULL norm"
);
}
use crate::engine::thread::ThreadQuery;
fn q(entities: &[&str], phrases: &[&str], topics: &[&str]) -> ThreadQuery {
ThreadQuery {
entities: entities.iter().map(|s| s.to_string()).collect(),
phrases: phrases.iter().map(|s| s.to_string()).collect(),
topic_rids: topics.iter().map(|s| s.to_string()).collect(),
}
}
#[test]
fn v2_union_dedup_and_multi_route_provenance() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_row(
&db,
"j1",
"Alpha standup notes",
&meta_empty(),
BASE_MICROS,
&["Alpha"],
1.0,
);
seed_row(
&db,
"j2",
"Alpha ran the quarterly sync today",
&meta_empty(),
BASE_MICROS + 1_000_000,
&["Alpha"],
2.0,
);
seed_row(
&db,
"j3",
"Notes from the quarterly sync recap",
&meta_empty(),
BASE_MICROS + 2_000_000,
&["Beta"],
3.0,
);
drain(&db);
let out = db
.recall_thread_v2(NS, &q(&["Alpha"], &["quarterly sync"], &[]), 10)
.unwrap();
assert_eq!(out.total, 3, "union of both routes, deduped");
assert_eq!(out.returned, 3);
assert_eq!(out.omitted, 0);
let rids: Vec<&str> = out.items.iter().map(|i| i.rid.as_str()).collect();
assert_eq!(rids, vec!["j1", "j2", "j3"], "chronological order");
assert_eq!(out.items[0].routes, vec!["entity"]);
assert_eq!(
out.items[1].routes,
vec!["entity", "phrase"],
"both routes, stable order"
);
assert_eq!(out.items[2].routes, vec!["phrase"]);
assert_eq!(out.items[1].entities, vec!["Alpha".to_string()]);
assert_eq!(out.items[1].phrases, vec!["quarterly sync".to_string()]);
assert!(out.items[0].phrases.is_empty());
assert!(out.items[2].entities.is_empty());
assert!(out.items.iter().all(|i| i.topic_rids.is_empty()));
}
#[test]
fn v2_per_anchor_phrase_provenance_in_request_order() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_row(
&db,
"k1",
"budget review and roadmap planning in one session",
&meta_empty(),
BASE_MICROS,
&["Gamma"],
1.0,
);
seed_row(
&db,
"k2",
"roadmap planning only here",
&meta_empty(),
BASE_MICROS + 1_000_000,
&["Gamma"],
2.0,
);
drain(&db);
let out = db
.recall_thread_v2(
NS,
&q(
&[],
&["budget review", "roadmap planning", "budget review"],
&[],
),
10,
)
.unwrap();
assert_eq!(out.total, 2);
assert_eq!(
out.items[0].phrases,
vec!["budget review".to_string(), "roadmap planning".to_string()],
"both matching phrases, request order, duplicate collapsed"
);
assert_eq!(
out.items[1].phrases,
vec!["roadmap planning".to_string()],
"only the phrase that actually matched"
);
assert_eq!(out.items[0].routes, vec!["phrase"]);
}
#[test]
fn v2_three_route_row_carries_all_provenance() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_row(
&db,
"t1",
"Alpha closed the vendor contract",
&serde_json::json!({"source_turn": 4}),
BASE_MICROS,
&["Alpha"],
1.0,
);
drain(&db);
let topic = crate::consolidate::record_synthesis(
&db,
&["t1".to_string()],
"Procurement thread organizer",
Some(&vec_seed(9.0, 8)),
"topic",
"rollup",
&serde_json::json!({}),
"topic:procurement-v1",
)
.unwrap();
let topic_rid = topic["consolidated_rid"].as_str().unwrap().to_string();
let out = db
.recall_thread_v2(
NS,
&q(&["Alpha"], &["vendor contract"], &[topic_rid.as_str()]),
10,
)
.unwrap();
let item = out
.items
.iter()
.find(|i| i.rid == "t1")
.expect("the evidence row is eligible");
assert_eq!(item.routes, vec!["entity", "phrase", "topic"]);
assert_eq!(item.entities, vec!["Alpha".to_string()]);
assert_eq!(item.phrases, vec!["vendor contract".to_string()]);
assert_eq!(item.topic_rids, vec![topic_rid.clone()]);
assert_eq!(item.source_turn, Some(4), "persisted column served");
}
#[test]
fn v2_fts_phrases_are_literals_never_syntax() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_row(
&db,
"l1",
"they said \"hello there\" twice",
&meta_empty(),
BASE_MICROS,
&["Delta"],
1.0,
);
seed_row(
&db,
"l2",
"alpha or beta together",
&meta_empty(),
BASE_MICROS + 1_000_000,
&["Delta"],
2.0,
);
seed_row(
&db,
"l3",
"alpha alone here",
&meta_empty(),
BASE_MICROS + 2_000_000,
&["Delta"],
3.0,
);
drain(&db);
let quoted = db
.recall_thread_v2(NS, &q(&[], &["said \"hello there\""], &[]), 10)
.unwrap();
assert_eq!(
quoted
.items
.iter()
.map(|i| i.rid.as_str())
.collect::<Vec<_>>(),
vec!["l1"],
"embedded double-quote is escaped, not FTS syntax"
);
let or_phrase = db
.recall_thread_v2(NS, &q(&[], &["alpha OR beta"], &[]), 10)
.unwrap();
assert_eq!(
or_phrase
.items
.iter()
.map(|i| i.rid.as_str())
.collect::<Vec<_>>(),
vec!["l2"],
"OR inside a phrase is a literal token — l3 (alpha alone) must NOT match"
);
}
#[test]
fn v2_topic_errors_are_typed_and_leak_nothing() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_row(
&db,
"m1",
"Alpha ordinary row",
&meta_empty(),
BASE_MICROS,
&["Alpha"],
1.0,
);
drain(&db);
let other_rid = db
.record_with_rid(
"m_other",
"Other-ns evidence",
"episodic",
0.5,
0.0,
604800.0,
&meta_empty(),
&vec_seed(5.0, 8),
"other_ns",
0.8,
"general",
"user",
None,
BASE_MICROS,
&[],
"test-model.v1",
None,
crate::provenance::WriteAdmission::Admitted,
)
.map(|_| "m_other".to_string())
.unwrap();
drain(&db);
let cross_topic = crate::consolidate::record_synthesis(
&db,
&[other_rid],
"Other-ns organizer",
Some(&vec_seed(6.0, 8)),
"topic",
"rollup",
&serde_json::json!({}),
"topic:otherns-v1",
)
.unwrap();
let cross_rid = cross_topic["consolidated_rid"]
.as_str()
.unwrap()
.to_string();
let probe = |topic: &str| -> String {
let err = db
.recall_thread_v2(NS, &q(&["Alpha"], &[], &[topic]), 10)
.unwrap_err();
assert!(
matches!(err, crate::error::YantrikDbError::InvalidThreadTopic { .. }),
"typed InvalidThreadTopic, got: {err:?}"
);
err.to_string().replace(topic, "<RID>")
};
let ordinary = probe("m1"); let nonexistent = probe("no-such-rid");
let cross_ns = probe(&cross_rid);
assert_eq!(
ordinary, nonexistent,
"ordinary-memory and nonexistent rids: identical error shape"
);
assert_eq!(
nonexistent, cross_ns,
"nonexistent and cross-namespace rids: identical error shape (no leakage)"
);
}
#[test]
fn v2_totals_positions_and_limit_zero() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..10 {
seed_row(
&db,
&format!("n{i}"),
&format!("Alpha step {i}"),
&meta_empty(),
BASE_MICROS + (i as i64) * 1_000_000,
&["Alpha"],
i as f32,
);
}
drain(&db);
let page = db
.recall_thread_v2(NS, &q(&["Alpha"], &[], &[]), 4)
.unwrap();
assert_eq!(
(page.total, page.returned, page.omitted),
(10, 4, 6),
"full-union total, page-bounded returned"
);
assert_eq!(
page.items.iter().map(|i| i.position).collect::<Vec<_>>(),
vec![1, 2, 3, 4],
"positions contiguous from 1 over the full-thread numbering"
);
let val = serde_json::to_value(&page).unwrap();
assert_eq!(
val["returned"], 4,
"returned is an explicit serialized field"
);
assert_eq!(val["items"].as_array().unwrap().len(), 4);
let zero = db
.recall_thread_v2(NS, &q(&["Alpha"], &[], &[]), 0)
.unwrap();
assert_eq!(
(zero.items.len(), zero.total, zero.returned, zero.omitted),
(0, 10, 0, 10),
"limit=0 must still report the exact full-union total"
);
let none = db.recall_thread_v2(NS, &q(&[], &[], &[]), 10).unwrap();
assert_eq!((none.total, none.returned, none.omitted), (0, 0, 0));
}
#[test]
fn v2_sql_orders_by_persisted_turn_with_nulls_last() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let t = BASE_MICROS;
seed_row(
&db,
"z_turn2",
"Epsilon event two",
&serde_json::json!({"source_turn": 2}),
t,
&["Epsilon"],
1.0,
);
seed_row(
&db,
"a_turn9",
"Epsilon event nine",
&serde_json::json!({"turn_id": 9}),
t,
&["Epsilon"],
2.0,
);
seed_row(
&db,
"b_none",
"Epsilon event without a turn",
&meta_empty(),
t,
&["Epsilon"],
3.0,
);
seed_row(
&db,
"a_none",
"Epsilon event also without a turn",
&serde_json::json!({"source_turn": "not-a-turn"}),
t,
&["Epsilon"],
4.0,
);
seed_row(
&db,
"later_turn0",
"Epsilon later event",
&serde_json::json!({"source_turn": 0}),
t + 1_000_000,
&["Epsilon"],
5.0,
);
drain(&db);
let out = db
.recall_thread_v2(NS, &q(&["Epsilon"], &[], &[]), 10)
.unwrap();
assert_eq!(
out.items.iter().map(|i| i.rid.as_str()).collect::<Vec<_>>(),
vec!["z_turn2", "a_turn9", "a_none", "b_none", "later_turn0"],
"turn asc (2 before 9 despite rid order), then NULLS LAST by rid, \
then created_at dominates"
);
assert_eq!(
out.items.iter().map(|i| i.source_turn).collect::<Vec<_>>(),
vec![Some(2), Some(9), None, None, Some(0)],
"the persisted column's values, never invented"
);
}
#[test]
fn v2_input_caps_are_typed() {
use crate::error::YantrikDbError;
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_row(
&db,
"p1",
"Alpha row",
&meta_empty(),
BASE_MICROS,
&["Alpha"],
1.0,
);
drain(&db);
assert!(
db.recall_thread_v2(NS, &q(&["Alpha"], &[], &[]), 10_000)
.is_ok(),
"limit at the cap works"
);
let over = db
.recall_thread_v2(NS, &q(&["Alpha"], &[], &[]), 10_001)
.unwrap_err();
assert!(matches!(over, YantrikDbError::InvalidInput(_)), "{over:?}");
let many: Vec<String> = (0..65).map(|i| format!("e{i}")).collect();
let many_refs: Vec<&str> = many.iter().map(String::as_str).collect();
let err = db
.recall_thread_v2(NS, &q(&many_refs, &[], &[]), 10)
.unwrap_err();
assert!(matches!(err, YantrikDbError::InvalidInput(_)), "{err:?}");
let big = "x".repeat(513);
let err = db
.recall_thread_v2(NS, &q(&[], &[big.as_str()], &[]), 10)
.unwrap_err();
assert!(matches!(err, YantrikDbError::InvalidInput(_)), "{err:?}");
}
#[test]
fn v2_supersede_exclusion_matches_v1() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for (i, rid) in ["s1", "s2", "s3"].iter().enumerate() {
seed_row(
&db,
rid,
&format!("Alpha supersede row {i}"),
&meta_empty(),
BASE_MICROS + (i as i64) * 1_000_000,
&["Alpha"],
i as f32,
);
}
drain(&db);
assert!(db.status_read_policy());
db.link(
"s3",
&crate::types::RecordLink {
target_rid: "s1".to_string(),
link_type: crate::types::LinkType::Supersedes,
},
)
.unwrap();
let v1 = db.recall_thread(NS, &["Alpha"], 10).unwrap();
let v2 = db
.recall_thread_v2(NS, &q(&["Alpha"], &[], &[]), 10)
.unwrap();
assert_eq!(v1.total, v2.total, "same eligible count");
assert_eq!(
v1.items.iter().map(|i| i.rid.as_str()).collect::<Vec<_>>(),
v2.items.iter().map(|i| i.rid.as_str()).collect::<Vec<_>>(),
"same rows in the same order — the SQL NOT EXISTS is equivalent \
to v1's post-SQL superseded_rids_among"
);
assert!(
!v2.items.iter().any(|i| i.rid == "s1"),
"superseded row excluded in SQL"
);
}
#[test]
fn v2_replication_parity_on_follower() {
use crate::replication::{apply_ops, extract_ops_since};
let leader = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..6 {
let (name, tag) = if i % 2 == 0 {
("Alpha", "milestone")
} else {
("Beta", "sync")
};
seed_row(
&leader,
&format!("r{i}"),
&format!("{name} {tag} number {i}"),
&serde_json::json!({"source_turn": i}),
BASE_MICROS + (i as i64) * 1_000_000,
&[name],
i as f32,
);
}
leader.relate("Alpha", "Beta", "related_to", 1.0).unwrap();
drain(&leader);
let query = q(&["Alpha"], &["sync"], &[]);
let leader_out = leader.recall_thread_v2(NS, &query, 100).unwrap();
assert_eq!(leader_out.total, 6, "3 entity rows + 3 phrase rows");
assert!(
leader_out.items.iter().any(|i| i.routes == vec!["phrase"]),
"phrase-only rows present"
);
let follower = YantrikDB::new(":memory:", 8).unwrap();
let ops = extract_ops_since(&leader.conn(), None, None, None, 1000).unwrap();
apply_ops(&follower, &ops).unwrap();
let follower_out = follower.recall_thread_v2(NS, &query, 100).unwrap();
assert_eq!(
follower_out, leader_out,
"identical v2 result on both sides — order, turns, routes, \
per-anchor fields, totals, returned, omitted"
);
}
#[test]
fn v2_encrypted_phrase_route_is_typed_error() {
use crate::error::YantrikDbError;
let db = YantrikDB::new_encrypted(":memory:", 8, &[7u8; 32]).unwrap();
seed_row(
&db,
"e1",
"Alpha encrypted row",
&serde_json::json!({"source_turn": 1}),
BASE_MICROS,
&["Alpha"],
1.0,
);
drain(&db);
let err = db
.recall_thread_v2(NS, &q(&["Alpha"], &["anything"], &[]), 10)
.unwrap_err();
match err {
YantrikDbError::CapabilityUnavailable { capability, .. } => {
assert_eq!(capability, "phrase_thread_route");
}
other => panic!("expected CapabilityUnavailable, got {other:?}"),
}
let ok = db
.recall_thread_v2(NS, &q(&["Alpha"], &[], &[]), 10)
.unwrap();
assert_eq!(ok.total, 1);
assert_eq!(ok.items[0].source_turn, Some(1));
}
}