use super::*;
const MARKER: &str = "source_turn_backfill_complete";
fn marker_value(db: &YantrikDB) -> Option<String> {
db.conn()
.query_row(
"SELECT value FROM meta WHERE key = ?1",
params![MARKER],
|r| r.get(0),
)
.ok()
}
fn column_turn(db: &YantrikDB, rid: &str) -> Option<i64> {
db.conn()
.query_row(
"SELECT source_turn FROM memories WHERE rid = ?1",
params![rid],
|r| r.get(0),
)
.unwrap()
}
fn census_violations(conn: &rusqlite::Connection) -> i64 {
conn.query_row(
"SELECT COUNT(*) FROM memories \
WHERE json_valid(metadata) \
AND source_turn IS NULL \
AND (CASE \
WHEN json_type(metadata, '$.source_turn') = 'integer' \
AND json_extract(metadata, '$.source_turn') >= 0 \
THEN json_extract(metadata, '$.source_turn') \
WHEN json_type(metadata, '$.turn_id') = 'integer' \
AND json_extract(metadata, '$.turn_id') >= 0 \
THEN json_extract(metadata, '$.turn_id') \
ELSE NULL END) IS NOT NULL",
[],
|r| r.get(0),
)
.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");
}
#[test]
fn census_every_memories_writer_stamps_the_source_turn_column() {
use crate::replication::{apply_ops, extract_ops_since};
let leader = YantrikDB::new(":memory:", 8).unwrap();
let mut record = |text: &str, meta: serde_json::Value| -> String {
leader
.record(
text,
"episodic",
0.5,
0.0,
604800.0,
&meta,
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap()
};
let rid_turn = record("plain turn", serde_json::json!({"source_turn": 5}));
let rid_fallback = record("fallback turn", serde_json::json!({"turn_id": 2}));
let rid_fallthrough = record(
"invalid preferred falls through",
serde_json::json!({"source_turn": "bogus", "turn_id": 7}),
);
let rid_invalid = record("invalid only", serde_json::json!({"source_turn": -3}));
drain(&leader);
assert_eq!(column_turn(&leader, &rid_turn), Some(5));
assert_eq!(column_turn(&leader, &rid_fallback), Some(2));
assert_eq!(
column_turn(&leader, &rid_fallthrough),
Some(7),
"or_else fallback: invalid source_turn falls through to turn_id"
);
assert_eq!(
column_turn(&leader, &rid_invalid),
None,
"never invented from an invalid value"
);
leader
.correct(
&rid_turn,
None,
Some(&serde_json::json!({"source_turn": 9})),
None,
None,
"turn was wrong",
)
.unwrap();
assert_eq!(column_turn(&leader, &rid_turn), Some(9));
let generation = leader.search_generation();
leader
.correct_with_embedding(
&rid_fallback,
Some("fallback turn, reworded"),
&vec_seed(3.0, 8),
generation,
Some(&serde_json::json!({"turn_id": 4})),
None,
None,
"reworded",
)
.unwrap();
assert_eq!(column_turn(&leader, &rid_fallback), Some(4));
let follower = YantrikDB::new(":memory:", 8).unwrap();
let ops: Vec<_> = extract_ops_since(&leader.conn(), None, None, None, 100)
.unwrap()
.into_iter()
.filter(|o| o.op_type != "correct" || o.embedding.is_none())
.collect();
apply_ops(&follower, &ops).unwrap();
for (name, db) in [("leader", &leader), ("follower", &follower)] {
let conn = db.conn();
assert_eq!(
census_violations(&conn),
0,
"{name}: a writer persisted turn-bearing metadata without stamping \
the v50 column; every memories writer must bind extract_source_turn()"
);
}
for rid in [&rid_turn, &rid_fallthrough, &rid_invalid] {
assert_eq!(
column_turn(&follower, rid),
column_turn(&leader, rid),
"follower column equals leader column for {rid}"
);
}
}
#[test]
fn replication_uses_canonical_scalar_with_legacy_fallback() {
use crate::replication::{apply_ops, extract_ops_since};
let leader = YantrikDB::new(":memory:", 8).unwrap();
let rid = leader
.record(
"canonical scalar row",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"source_turn": 5}),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let ops = extract_ops_since(&leader.conn(), None, None, None, 100).unwrap();
let record_op = ops
.iter()
.find(|o| o.op_type == "record" && o.target_rid.as_deref() == Some(rid.as_str()))
.expect("the record op exists")
.clone();
assert_eq!(
record_op.payload["source_turn"], 5,
"the leader-derived canonical scalar rides the op payload"
);
let mut canonical = record_op.clone();
canonical.payload["source_turn"] = serde_json::json!(9);
let f1 = YantrikDB::new(":memory:", 8).unwrap();
apply_ops(&f1, &[canonical]).unwrap();
assert_eq!(
column_turn(&f1, &rid),
Some(9),
"the canonical scalar is used directly, never re-derived"
);
let mut legacy = record_op.clone();
legacy
.payload
.as_object_mut()
.unwrap()
.remove("source_turn");
let f2 = YantrikDB::new(":memory:", 8).unwrap();
apply_ops(&f2, &[legacy]).unwrap();
assert_eq!(
column_turn(&f2, &rid),
Some(5),
"legacy fallback parses metadata through the shared extractor"
);
}
#[test]
fn replication_replay_repair_fills_null_never_overwrites() {
use crate::replication::{apply_ops, extract_ops_since};
let leader = YantrikDB::new(":memory:", 8).unwrap();
let rid = leader
.record(
"repair target",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"source_turn": 5}),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let ops = extract_ops_since(&leader.conn(), None, None, None, 100).unwrap();
let record_op = ops
.iter()
.find(|o| o.op_type == "record" && o.target_rid.as_deref() == Some(rid.as_str()))
.unwrap()
.clone();
let follower = YantrikDB::new(":memory:", 8).unwrap();
apply_ops(&follower, &[record_op.clone()]).unwrap();
assert_eq!(column_turn(&follower, &rid), Some(5));
follower
.conn()
.execute(
"UPDATE memories SET source_turn = NULL WHERE rid = ?1",
params![rid],
)
.unwrap();
let mut replay = record_op.clone();
replay.op_id = crate::id::new_id();
apply_ops(&follower, &[replay]).unwrap();
assert_eq!(
column_turn(&follower, &rid),
Some(5),
"replay repair fills a NULL column"
);
follower
.conn()
.execute(
"UPDATE memories SET source_turn = 8 WHERE rid = ?1",
params![rid],
)
.unwrap();
let mut replay2 = record_op.clone();
replay2.op_id = crate::id::new_id();
apply_ops(&follower, &[replay2]).unwrap();
assert_eq!(
column_turn(&follower, &rid),
Some(8),
"replay repair must not overwrite a non-NULL newer value"
);
}
#[test]
fn stamped_write_after_completion_preserves_marker_true() {
let db = YantrikDB::new_encrypted(":memory:", 8, &[7u8; 32]).unwrap();
assert_eq!(
marker_value(&db).as_deref(),
Some("1"),
"fresh store: complete immediately"
);
let rid = db
.record(
"Alpha stamped write",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"source_turn": 3}),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.link_memory_entity(&rid, "Alpha").unwrap();
assert_eq!(
marker_value(&db).as_deref(),
Some("1"),
"an engine-supported stamped write preserves marker=true \
(its own trigger fire is not staleness)"
);
let query = crate::engine::thread::ThreadQuery {
entities: vec!["Alpha".to_string()],
phrases: Vec::new(),
topic_rids: Vec::new(),
};
let out = db
.recall_thread_v2("default", &query, 10)
.expect("strict query must succeed: no error after a stamped write");
assert_eq!(out.total, 1);
assert_eq!(out.items[0].source_turn, Some(3));
}
#[test]
fn stamped_write_while_incomplete_keeps_marker_false_and_strict_query_errors() {
use crate::error::YantrikDbError;
let db = YantrikDB::new_encrypted(":memory:", 8, &[7u8; 32]).unwrap();
let rid1 = db
.record(
"Alpha first encrypted row",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"source_turn": 3}),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.link_memory_entity(&rid1, "Alpha").unwrap();
db.conn()
.execute(
"UPDATE memories SET source_turn = NULL WHERE rid = ?1",
params![rid1],
)
.unwrap();
assert_eq!(marker_value(&db).as_deref(), Some("0"), "trigger staled it");
let rid2 = db
.record(
"Alpha second encrypted row",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"source_turn": 4}),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.link_memory_entity(&rid2, "Alpha").unwrap();
assert_eq!(
marker_value(&db).as_deref(),
Some("0"),
"a stamped write PRESERVES the pre-write state: false stays false"
);
let query = crate::engine::thread::ThreadQuery {
entities: vec!["Alpha".to_string()],
phrases: Vec::new(),
topic_rids: Vec::new(),
};
let err = db.recall_thread_v2("default", &query, 10).unwrap_err();
assert!(
matches!(err, YantrikDbError::MaintenanceRequired { ref operation, .. }
if operation == "maintain_source_turn_backfill"),
"strict query refuses while incomplete: {err:?}"
);
let mut rounds = 0;
loop {
let progress = db.maintain_source_turn_backfill(1_000).unwrap();
rounds += 1;
if progress.complete {
break;
}
assert!(rounds < 100, "maintenance must terminate");
}
assert_eq!(marker_value(&db).as_deref(), Some("1"));
let out = db.recall_thread_v2("default", &query, 10).unwrap();
assert_eq!(out.total, 2);
assert_eq!(
out.items.iter().map(|i| i.source_turn).collect::<Vec<_>>(),
vec![Some(3), Some(4)],
"the decrypt-and-stamp pass restored the raw-NULLed column"
);
}
#[test]
fn raw_metadata_change_is_recomputed_never_null_filled() {
use crate::error::YantrikDbError;
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"Alpha recompute target",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"source_turn": 5}),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.link_memory_entity(&rid, "Alpha").unwrap();
assert_eq!(column_turn(&db, &rid), Some(5));
assert_eq!(marker_value(&db).as_deref(), Some("1"));
let query = crate::engine::thread::ThreadQuery {
entities: vec!["Alpha".to_string()],
phrases: Vec::new(),
topic_rids: Vec::new(),
};
db.conn()
.execute(
"UPDATE memories SET metadata = json_set(metadata, '$.source_turn', 7) \
WHERE rid = ?1",
params![rid],
)
.unwrap();
assert_eq!(marker_value(&db).as_deref(), Some("0"), "trigger staled it");
let err = db.recall_thread_v2("default", &query, 10).unwrap_err();
assert!(
matches!(err, YantrikDbError::MaintenanceRequired { .. }),
"the PLAINTEXT store's strict query also refuses before repair: {err:?}"
);
loop {
if db.maintain_source_turn_backfill(1_000).unwrap().complete {
break;
}
}
assert_eq!(marker_value(&db).as_deref(), Some("1"));
assert_eq!(
column_turn(&db, &rid),
Some(7),
"maintenance RECOMPUTED the stale non-NULL scalar to the current \
metadata's value — a NULL-only fill would have left 5"
);
let out = db.recall_thread_v2("default", &query, 10).unwrap();
assert_eq!(
out.items[0].source_turn,
Some(7),
"never a wrong value under marker=true"
);
db.conn()
.execute(
"UPDATE memories SET metadata = json_remove(metadata, '$.source_turn') \
WHERE rid = ?1",
params![rid],
)
.unwrap();
assert_eq!(marker_value(&db).as_deref(), Some("0"));
loop {
if db.maintain_source_turn_backfill(1_000).unwrap().complete {
break;
}
}
assert_eq!(marker_value(&db).as_deref(), Some("1"));
assert_eq!(
column_turn(&db, &rid),
None,
"metadata no longer carries a valid turn: the recompute must CLEAR \
the column, not keep the stale scalar"
);
let out = db.recall_thread_v2("default", &query, 10).unwrap();
assert_eq!(out.items[0].source_turn, None);
}
#[test]
fn v1_stays_correct_on_stale_store_while_v2_requires_maintenance() {
use crate::error::YantrikDbError;
let db = YantrikDB::new(":memory:", 8).unwrap();
let t = 1_700_000_000.0_f64;
let mut seed = |text: &str, turn: i64, seedv: f32| -> String {
let rid = db
.record_with_idempotency(
text,
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"source_turn": turn}),
&vec_seed(seedv, 8),
"default",
0.8,
"general",
"user",
None,
None,
Some(t), )
.unwrap();
db.link_memory_entity(&rid, "Alpha").unwrap();
rid
};
let rid5 = seed("Alpha event five", 5, 1.0);
let rid3 = seed("Alpha event three", 3, 2.0);
db.conn()
.execute(
"UPDATE memories SET metadata = json_set(metadata, '$.source_turn', 1) \
WHERE rid = ?1",
params![rid5],
)
.unwrap();
assert_eq!(marker_value(&db).as_deref(), Some("0"));
let v1 = db.recall_thread("default", &["Alpha"], 10).unwrap();
assert_eq!(
v1.items.iter().map(|i| i.rid.as_str()).collect::<Vec<_>>(),
vec![rid5.as_str(), rid3.as_str()],
"v1 orders by the rewritten turn (1 before 3) — its own decrypt path"
);
assert_eq!(
v1.items.iter().map(|i| i.source_turn).collect::<Vec<_>>(),
vec![Some(1), Some(3)]
);
let query = crate::engine::thread::ThreadQuery {
entities: vec!["Alpha".to_string()],
phrases: Vec::new(),
topic_rids: Vec::new(),
};
let err = db.recall_thread_v2("default", &query, 10).unwrap_err();
assert!(matches!(err, YantrikDbError::MaintenanceRequired { .. }));
loop {
if db.maintain_source_turn_backfill(1_000).unwrap().complete {
break;
}
}
let v2 = db.recall_thread_v2("default", &query, 10).unwrap();
assert_eq!(
v2.items.iter().map(|i| i.rid.as_str()).collect::<Vec<_>>(),
v1.items.iter().map(|i| i.rid.as_str()).collect::<Vec<_>>(),
"healed v2 equals v1"
);
}
#[test]
fn maintenance_batch_caps_are_typed() {
use crate::error::YantrikDbError;
let db = YantrikDB::new(":memory:", 8).unwrap();
let err = db.maintain_source_turn_backfill(0).unwrap_err();
assert!(matches!(err, YantrikDbError::InvalidInput(_)), "{err:?}");
let err = db.maintain_source_turn_backfill(10_001).unwrap_err();
assert!(matches!(err, YantrikDbError::InvalidInput(_)), "{err:?}");
let progress = db.maintain_source_turn_backfill(10_000).unwrap();
assert!(
progress.complete,
"empty store completes immediately at MAX"
);
}
#[test]
fn migration_from_v49_installs_column_backfill_and_triggers() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap().to_string();
let rid = {
let db = YantrikDB::new(&path, 8).unwrap();
let rid = db
.record(
"migration row",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"source_turn": 5}),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
drop(db);
rid
};
{
let conn = rusqlite::Connection::open(&path).unwrap();
conn.execute_batch(
"DROP TRIGGER memories_source_turn_marker_insert; \
DROP TRIGGER memories_source_turn_marker_update; \
DROP INDEX idx_memories_source_turn; \
ALTER TABLE memories DROP COLUMN source_turn;",
)
.unwrap();
conn.execute(
"DELETE FROM meta WHERE key IN \
('source_turn_backfill_complete', 'source_turn_invalidation_epoch', \
'source_turn_repair_cursor')",
[],
)
.unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '49')",
[],
)
.unwrap();
}
let db = YantrikDB::new(&path, 8).expect("v49 -> v50 upgrade must succeed");
assert_eq!(
column_turn(&db, &rid),
Some(5),
"the open()-time backfill stamped the pre-existing row"
);
assert_eq!(
marker_value(&db).as_deref(),
Some("1"),
"unencrypted store: complete after the backfill drains"
);
db.conn()
.execute(
"UPDATE memories SET metadata = json_set(metadata, '$.source_turn', 7) \
WHERE rid = ?1",
params![rid],
)
.unwrap();
assert_eq!(
marker_value(&db).as_deref(),
Some("0"),
"the migrated store's UPDATE trigger fires on a raw metadata write"
);
loop {
if db.maintain_source_turn_backfill(1_000).unwrap().complete {
break;
}
}
assert_eq!(column_turn(&db, &rid), Some(7), "recomputed to 7");
assert_eq!(marker_value(&db).as_deref(), Some("1"));
db.conn()
.execute(
"UPDATE memories SET metadata = json_remove(metadata, '$.source_turn') \
WHERE rid = ?1",
params![rid],
)
.unwrap();
assert_eq!(
marker_value(&db).as_deref(),
Some("0"),
"5→absent stales too"
);
loop {
if db.maintain_source_turn_backfill(1_000).unwrap().complete {
break;
}
}
assert_eq!(column_turn(&db, &rid), None, "recomputed to NULL");
db.conn()
.execute(
"INSERT INTO memories (rid, type, text, created_at, updated_at, importance, \
half_life, last_access, valence, metadata, namespace) \
VALUES ('raw-insert', 'episodic', 'x', 1.0, 1.0, 0.5, 604800.0, 1.0, 0.0, \
'{\"source_turn\": 2}', 'default')",
[],
)
.unwrap();
assert_eq!(
marker_value(&db).as_deref(),
Some("0"),
"the migrated store's INSERT trigger fires on a raw insert"
);
}
#[test]
fn queued_record_op_carries_canonical_scalar_with_follower_parity() {
use crate::replication::{apply_ops, extract_ops_since};
let leader = YantrikDB::new(":memory:", 8).unwrap();
let rid_turn = leader
.record_queued(
"queued row with a turn",
"episodic",
0.5,
0.5,
0.0,
604800.0,
&serde_json::json!({"source_turn": 7}),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
crate::provenance::GateVerdict::Clean,
None,
None,
)
.unwrap();
let rid_none = leader
.record_queued(
"queued row without a turn",
"episodic",
0.5,
0.5,
0.0,
604800.0,
&serde_json::json!({"note": "turnless"}),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
crate::provenance::GateVerdict::Clean,
None,
None,
)
.unwrap();
let ops = extract_ops_since(&leader.conn(), None, None, None, 100).unwrap();
let op_turn = ops
.iter()
.find(|o| o.op_type == "record" && o.target_rid.as_deref() == Some(rid_turn.as_str()))
.expect("the queued record op exists")
.clone();
assert_eq!(
op_turn.payload["source_turn"], 7,
"the queued payload carries the canonical scalar like the sync route"
);
let op_none = ops
.iter()
.find(|o| o.op_type == "record" && o.target_rid.as_deref() == Some(rid_none.as_str()))
.expect("the turnless queued record op exists")
.clone();
assert_eq!(
op_none.payload.get("source_turn"),
Some(&serde_json::Value::Null),
"no valid turn is PRESENT-NULL (authoritative None), never an absent key"
);
let mut divergent = op_turn.clone();
divergent.payload["source_turn"] = serde_json::json!(11);
let follower = YantrikDB::new(":memory:", 8).unwrap();
apply_ops(&follower, &[divergent, op_none.clone()]).unwrap();
assert_eq!(
column_turn(&follower, &rid_turn),
Some(11),
"the follower uses the queued op's canonical scalar, not a metadata re-parse"
);
assert_eq!(
column_turn(&follower, &rid_none),
None,
"present-null materializes as an (authoritative) NULL column"
);
}
#[test]
fn malformed_canonical_scalar_is_rejected_not_coerced() {
use crate::replication::{apply_ops, extract_ops_since};
let leader = YantrikDB::new(":memory:", 8).unwrap();
let rid = leader
.record(
"malformed scalar probe",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"source_turn": 5}),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let record_op = extract_ops_since(&leader.conn(), None, None, None, 100)
.unwrap()
.into_iter()
.find(|o| o.op_type == "record" && o.target_rid.as_deref() == Some(rid.as_str()))
.expect("the record op exists");
for bad in [
serde_json::json!(-3),
serde_json::json!(5.5),
serde_json::json!("5"),
serde_json::json!({"turn": 5}),
] {
let mut op = record_op.clone();
op.payload["source_turn"] = bad.clone();
let follower = YantrikDB::new(":memory:", 8).unwrap();
let err = apply_ops(&follower, &[op]).expect_err(&format!(
"a present malformed canonical scalar ({bad}) must be rejected"
));
assert!(
matches!(
err,
crate::error::YantrikDbError::InvalidInput(ref msg)
if msg.contains("malformed canonical source_turn")
),
"typed InvalidInput naming the fault, got: {err:?}"
);
let rows: i64 = follower
.conn()
.query_row("SELECT COUNT(*) FROM memories", [], |r| r.get(0))
.unwrap();
assert_eq!(rows, 0, "the rejected op must not have materialized a row");
}
let follower = YantrikDB::new(":memory:", 8).unwrap();
apply_ops(&follower, &[record_op.clone()]).unwrap();
assert_eq!(column_turn(&follower, &rid), Some(5));
leader
.correct(
&rid,
None,
Some(&serde_json::json!({"source_turn": 6})),
None,
None,
"reviewer blocker 2 probe",
)
.unwrap();
let correct_op = extract_ops_since(&leader.conn(), None, None, None, 100)
.unwrap()
.into_iter()
.find(|o| o.op_type.starts_with("correct") && o.target_rid.as_deref() == Some(rid.as_str()))
.expect("the correction op exists");
let mut bad_correct = correct_op;
bad_correct.payload["source_turn"] = serde_json::json!(-1);
let err = apply_ops(&follower, &[bad_correct])
.expect_err("a malformed canonical scalar on a correction op must be rejected");
assert!(
matches!(
err,
crate::error::YantrikDbError::InvalidInput(ref msg)
if msg.contains("malformed canonical source_turn")
),
"typed InvalidInput at the correction apply site, got: {err:?}"
);
assert_eq!(
column_turn(&follower, &rid),
Some(5),
"the rejected correction must not have cleared or changed the column"
);
}