use super::*;
#[test]
fn idempotent_retry_returns_original_rid_and_writes_nothing() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let write = |db: &YantrikDB| {
db.record_with_idempotency(
"the deploy uses zorbium for caching",
"semantic",
0.7,
0.0,
604800.0,
&serde_json::json!({"k": "v"}),
&vec_seed(1.0, 8),
"idem_ns",
0.8,
"work",
"user",
None,
Some("client-req-001"),
None,
)
};
let rid1 = write(&db).unwrap();
let rows = |db: &YantrikDB| -> i64 {
db.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE namespace = 'idem_ns'",
[],
|r| r.get(0),
)
.unwrap()
};
let ops = |db: &YantrikDB| -> i64 {
db.conn()
.query_row("SELECT COUNT(*) FROM oplog", [], |r| r.get(0))
.unwrap()
};
let stats_count = |db: &YantrikDB| -> i64 {
db.conn()
.query_row(
"SELECT count FROM namespace_importance_stats WHERE namespace = 'idem_ns'",
[],
|r| r.get(0),
)
.unwrap()
};
let (r1, o1, s1, p1) = (
rows(&db),
ops(&db),
stats_count(&db),
db.count_pending_ops().unwrap(),
);
assert_eq!(r1, 1, "first keyed write lands one row");
assert_eq!(s1, 1, "first keyed write advances stats once");
let rid2 = write(&db).unwrap();
assert_eq!(rid2, rid1, "retry returns the ORIGINAL rid");
assert_eq!(rows(&db), r1, "retry wrote no second row");
assert_eq!(ops(&db), o1, "retry wrote no oplog op of any kind");
assert_eq!(stats_count(&db), s1, "retry advanced no calibration stats");
assert_eq!(
db.count_pending_ops().unwrap(),
p1,
"retry enqueued no pending op"
);
let claims: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM idempotency_claims", [], |r| r.get(0))
.unwrap();
assert_eq!(claims, 1, "one claim row, state committed");
for _ in 0..2 {
db.record(
"keyless duplicate content",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"idem_ns",
0.8,
"work",
"user",
None,
)
.unwrap();
}
assert_eq!(rows(&db), r1 + 2, "keyless writes still append");
}
#[test]
fn idempotency_conflict_on_different_payload() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let write = |db: &YantrikDB, importance: f64| {
db.record_with_idempotency(
"conflicting payload probe",
"semantic",
importance,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"conf_ns",
0.8,
"work",
"user",
None,
Some("key-x"),
None,
)
};
let rid1 = write(&db, 0.7).unwrap();
let err = write(&db, 0.9).expect_err("scalar diff under the same key must conflict");
match &err {
crate::error::YantrikDbError::IdempotencyConflict { existing_rid, .. } => {
assert_eq!(existing_rid, &rid1, "conflict names the existing record");
}
other => panic!("expected IdempotencyConflict, got {other:?}"),
}
let rows: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE namespace = 'conf_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(rows, 1, "the conflicting write persisted nothing");
let stats: i64 = db
.conn()
.query_row(
"SELECT count FROM namespace_importance_stats WHERE namespace = 'conf_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(stats, 1, "the conflicting write advanced no stats");
}
#[test]
fn idempotency_holds_on_and_across_the_queued_route() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let write = |db: &YantrikDB, key: &str| {
db.record_with_idempotency(
"queued idempotent probe",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"q_ns",
0.8,
"work",
"user",
None,
Some(key),
None,
)
};
db.write_router.switch_to_queueing();
let rid_q = write(&db, "q-key").unwrap();
let pending_after_first = db.count_pending_ops().unwrap();
let rid_q2 = write(&db, "q-key").unwrap();
assert_eq!(rid_q2, rid_q, "queued retry returns the original rid");
assert_eq!(
db.count_pending_ops().unwrap(),
pending_after_first,
"queued retry enqueued no second op"
);
db.write_router.switch_to_normal();
let rid_s = write(&db, "cross-key").unwrap();
db.write_router.switch_to_queueing();
let pending_before = db.count_pending_ops().unwrap();
let rid_s2 = write(&db, "cross-key").unwrap();
assert_eq!(rid_s2, rid_s, "cross-route retry resolves to the sync rid");
assert_eq!(
db.count_pending_ops().unwrap(),
pending_before,
"cross-route retry enqueued nothing"
);
db.write_router.switch_to_normal();
}
#[test]
fn idempotency_digest_uses_raw_importance_not_calibrated() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..8 {
db.record(
&format!("sat-{i}"),
"semantic",
1.0,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0 + i as f32 * 0.01, 8),
"sat_ns",
0.8,
"work",
"user",
None,
)
.unwrap();
}
let write = |db: &YantrikDB| {
db.record_with_idempotency(
"saturated keyed write",
"semantic",
0.9, 0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"sat_ns",
0.8,
"work",
"user",
None,
Some("sat-key"),
None,
)
};
let rid1 = write(&db).unwrap();
let stored: f64 = db
.conn()
.query_row(
"SELECT importance FROM memories WHERE rid = ?1",
rusqlite::params![rid1],
|r| r.get(0),
)
.unwrap();
assert!(
stored < 0.9,
"precondition: calibration deflated the stored importance, got {stored}"
);
let rid2 = write(&db).expect("identical RAW retry must be a hit, not a conflict");
assert_eq!(rid2, rid1);
}
#[test]
fn invalid_idempotency_keys_are_rejected() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for bad in ["", " ", &"k".repeat(513)] {
let err = db
.record_with_idempotency(
"probe",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"work",
"user",
None,
Some(bad),
None,
)
.expect_err("bad key must be refused");
assert!(
matches!(
err,
crate::error::YantrikDbError::InvalidIdempotencyKey { .. }
),
"expected InvalidIdempotencyKey, got {err:?}"
);
}
let rows: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM memories", [], |r| r.get(0))
.unwrap();
assert_eq!(rows, 0, "refused keys wrote nothing");
}
#[cfg(feature = "bundled-embedder")]
#[test]
fn keyed_duplicate_resolves_even_under_backpressure() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let dim = db.embedding_dim();
let keyed = |db: &YantrikDB| {
db.record_with_idempotency(
"the keyed write that must stay resolvable",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(9.0, 8),
"bp_idem_ns",
0.8,
"work",
"user",
None,
Some("bp-key"),
None,
)
};
let rid = keyed(&db).unwrap();
let mut saturated = false;
for i in 0..400 {
let emb: Vec<f32> = (0..dim).map(|j| ((i + j) as f32) * 0.001).collect();
match db.record(
&format!("bp-filler-{i}"),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
) {
Ok(_) => {}
Err(crate::error::YantrikDbError::Backpressure { .. }) => {
saturated = true;
break;
}
Err(e) => panic!("unexpected: {e:?}"),
}
}
assert!(saturated, "delta never saturated");
let rid2 = keyed(&db).expect(
"a keyed duplicate writes nothing and must resolve to its Hit, \
not die on Backpressure",
);
assert_eq!(rid2, rid, "the hit returns the original rid");
let err = db
.record_with_idempotency(
"a genuinely new keyed write",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(10.0, 8),
"bp_idem_ns",
0.8,
"work",
"user",
None,
Some("bp-key-new"),
None,
)
.expect_err("a NEW keyed write under saturation must still backpressure");
assert!(
matches!(err, crate::error::YantrikDbError::Backpressure { .. }),
"expected Backpressure for the new keyed write, got {err:?}"
);
db.write_router.switch_to_queueing();
let rid3 = keyed(&db).expect("queued dup under saturation must resolve");
assert_eq!(rid3, rid, "queued dup returns the original rid");
db.write_router.switch_to_normal();
}
#[test]
fn keyed_duplicate_resolves_under_pending_queue_saturation() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let write = |db: &YantrikDB, key: &str| {
db.record_with_idempotency(
"pending-saturation probe",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"pq_ns",
0.8,
"work",
"user",
None,
Some(key),
None,
)
};
let rid = write(&db, "pq-key").unwrap();
let real = db
.pending_op_count
.swap(1_000_000, std::sync::atomic::Ordering::SeqCst);
let rid2 = write(&db, "pq-key").expect("keyed dup must resolve under pending-queue saturation");
assert_eq!(rid2, rid, "dup returns the original rid");
let err = write(&db, "pq-key-new")
.expect_err("a NEW keyed write under pending saturation must backpressure");
assert!(
matches!(err, crate::error::YantrikDbError::Backpressure { .. }),
"expected Backpressure from the locked check, got {err:?}"
);
db.write_router.switch_to_queueing();
let rid3 = write(&db, "pq-key").expect("queued dup must resolve under pending saturation");
assert_eq!(rid3, rid);
let err = write(&db, "pq-key-new-q")
.expect_err("a NEW keyed queued write under pending saturation must backpressure");
assert!(
matches!(err, crate::error::YantrikDbError::Backpressure { .. }),
"expected Backpressure on the queued route, got {err:?}"
);
db.write_router.switch_to_normal();
db.pending_op_count
.store(real, std::sync::atomic::Ordering::SeqCst);
}
#[test]
fn record_text_keyed_retry_hits_without_embedding_again() {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
struct CountingEmbedder(Arc<AtomicUsize>);
impl crate::types::Embedder for CountingEmbedder {
fn embed(
&self,
t: &str,
) -> std::result::Result<Vec<f32>, Box<dyn std::error::Error + Send + Sync>> {
self.0.fetch_add(1, Ordering::SeqCst);
let mut v = vec![0.1; 8];
v[0] = (t.len() % 7) as f32 * 0.1;
Ok(v)
}
fn dim(&self) -> usize {
8
}
}
let calls = Arc::new(AtomicUsize::new(0));
let mut db = YantrikDB::new(":memory:", 8).unwrap();
db.set_embedder(Box::new(CountingEmbedder(Arc::clone(&calls))))
.unwrap();
let write = |db: &YantrikDB| {
db.record_text_with_idempotency(
"keyed text write",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
"rt_ns",
0.8,
"work",
"user",
None,
Some("rt-key"),
None,
)
};
let rid = write(&db).unwrap();
let embeds_after_first = calls.load(Ordering::SeqCst);
assert!(embeds_after_first >= 1, "first write embeds");
let rid2 = write(&db).unwrap();
assert_eq!(rid2, rid, "retry returns the original rid");
assert_eq!(
calls.load(Ordering::SeqCst),
embeds_after_first,
"the duplicate retry must resolve at the probe, BEFORE the embed — \
zero additional embedder invocations"
);
let rows: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE namespace = 'rt_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(rows, 1, "one row total");
}
#[test]
fn same_key_across_record_and_record_text_is_a_conflict() {
struct Fake;
impl crate::types::Embedder for Fake {
fn embed(
&self,
t: &str,
) -> std::result::Result<Vec<f32>, Box<dyn std::error::Error + Send + Sync>> {
let mut v = vec![0.1; 8];
v[0] = (t.len() % 7) as f32 * 0.1;
Ok(v)
}
fn dim(&self) -> usize {
8
}
}
let mut db = YantrikDB::new(":memory:", 8).unwrap();
db.set_embedder(Box::new(Fake)).unwrap();
db.record_with_idempotency(
"cross surface text",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"xs_ns",
0.8,
"work",
"user",
None,
Some("xs-key"),
None,
)
.unwrap();
let err = db
.record_text_with_idempotency(
"cross surface text",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
"xs_ns",
0.8,
"work",
"user",
None,
Some("xs-key"),
None,
)
.expect_err("cross-surface same-key must conflict");
assert!(
matches!(
err,
crate::error::YantrikDbError::IdempotencyConflict { .. }
),
"expected IdempotencyConflict, got {err:?}"
);
}
#[test]
fn record_text_normalizes_blank_namespace_like_record() {
struct Fake;
impl crate::types::Embedder for Fake {
fn embed(
&self,
t: &str,
) -> std::result::Result<Vec<f32>, Box<dyn std::error::Error + Send + Sync>> {
let mut v = vec![0.1; 8];
v[0] = (t.len() % 7) as f32 * 0.1;
Ok(v)
}
fn dim(&self) -> usize {
8
}
}
let mut db = YantrikDB::new(":memory:", 8).unwrap();
db.set_embedder(Box::new(Fake)).unwrap();
let rid = db
.record_text_with_idempotency(
"blank namespace probe",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
" ",
0.8,
"work",
"user",
None,
Some("ns-key"),
None,
)
.unwrap();
let ns: String = db
.conn()
.query_row(
"SELECT namespace FROM memories WHERE rid = ?1",
rusqlite::params![rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(ns, "default", "blank namespace must normalize to default");
let claim_ns: String = db
.conn()
.query_row(
"SELECT namespace FROM idempotency_claims WHERE idempotency_key = 'ns-key'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
claim_ns, "default",
"the claim must scope under the NORMALIZED namespace"
);
let rid2 = db
.record_text_with_idempotency(
"blank namespace probe",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
"",
0.8,
"work",
"user",
None,
Some("ns-key"),
None,
)
.unwrap();
assert_eq!(
rid2, rid,
"'' and ' ' resolve to the same normalized scope"
);
}
#[test]
fn batch_append_failure_writes_nothing_at_all() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let dim = db.embedding_dim();
db.record_batch(&[RecordInput {
created_at: None,
idempotency_key: None,
text: "Quarterly sync with Klaxonberg about the roadmap".into(),
memory_type: "episodic".into(),
importance: 0.6,
valence: 0.0,
half_life: 604800.0,
metadata: serde_json::json!({}),
embedding: vec_seed(1.5, 8),
namespace: "ctrl_ns".into(),
certainty: 0.8,
domain: "work".into(),
source: "user".into(),
emotional_state: None,
}])
.unwrap();
let ctrl_entities: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM entities WHERE name = 'Klaxonberg'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
ctrl_entities, 1,
"control precondition: batch entity extraction must fire for this \
text shape, or the absence assertions below prove nothing"
);
let mut saturated = false;
for i in 0..400 {
let emb: Vec<f32> = (0..dim).map(|j| ((i + j) as f32) * 0.001).collect();
match db.record(
&format!("filler-{i}"),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"filler_ns",
0.8,
"general",
"user",
None,
) {
Ok(_) => {}
Err(crate::error::YantrikDbError::Backpressure { .. }) => {
saturated = true;
break;
}
Err(e) => panic!("unexpected: {e:?}"),
}
}
assert!(saturated, "delta never saturated");
let ops_before: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM oplog", [], |r| r.get(0))
.unwrap();
let err = db
.record_batch(&[RecordInput {
created_at: None,
idempotency_key: None,
text: "Quarterly sync with Quorvexia about the roadmap".into(),
memory_type: "episodic".into(),
importance: 0.9,
valence: 0.0,
half_life: 604800.0,
metadata: serde_json::json!({}),
embedding: vec_seed(2.0, 8),
namespace: "batch_ns".into(),
certainty: 0.8,
domain: "work".into(),
source: "user".into(),
emotional_state: None,
}])
.expect_err("batch into the saturated delta must fail");
assert!(
matches!(err, crate::error::YantrikDbError::Backpressure { .. }),
"expected Backpressure, got {err:?}"
);
let rows: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE namespace = 'batch_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(rows, 0, "no memory rows for the rejected batch");
let orphan_entities: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM entities WHERE name = 'Quorvexia'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
orphan_entities, 0,
"a rejected batch left an orphaned `entities` row (#92)"
);
let orphan_links: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memory_entities WHERE entity_name = 'Quorvexia'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
orphan_links, 0,
"a rejected batch left orphaned `memory_entities` links (#92)"
);
assert!(
!db.graph_index
.read()
.all_entity_names()
.iter()
.any(|n| n == "Quorvexia"),
"a rejected batch polluted the in-memory graph_index (#92)"
);
let ops_after: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM oplog", [], |r| r.get(0))
.unwrap();
assert_eq!(
ops_after, ops_before,
"no oplog entries for a rejected batch"
);
assert!(
db.conn().is_autocommit(),
"record_batch failure left a savepoint open on the shared conn"
);
}
#[test]
fn record_batch_normalizes_blank_namespace_like_record() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let mk = |text: &str, ns: &str, emb: f32| RecordInput {
created_at: None,
idempotency_key: None,
text: text.into(),
memory_type: "semantic".into(),
importance: 0.7,
valence: 0.0,
half_life: 604800.0,
metadata: serde_json::json!({}),
embedding: vec_seed(emb, 8),
namespace: ns.into(),
certainty: 0.8,
domain: "work".into(),
source: "user".into(),
emotional_state: None,
};
let rids = db
.record_batch(&[
mk("blank ns batch item one", " ", 1.0),
mk("blank ns batch item two", "", 2.0),
])
.unwrap();
assert_eq!(rids.len(), 2);
let in_default: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE namespace = 'default'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
in_default, 2,
"blank-namespace batch rows must land under 'default', like record()"
);
let in_blank: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE TRIM(namespace) = ''",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
in_blank, 0,
"no row may persist under a raw blank namespace partition (#98)"
);
for rid in &rids {
let payload: String = db
.conn()
.query_row(
"SELECT payload FROM oplog WHERE target_rid = ?1 AND op_type = 'record'",
rusqlite::params![rid],
|r| r.get(0),
)
.unwrap();
let v: serde_json::Value = serde_json::from_str(&payload).unwrap();
assert_eq!(
v["namespace"], "default",
"op payload must replicate the NORMALIZED namespace"
);
}
let count: i64 = db
.conn()
.query_row(
"SELECT count FROM namespace_importance_stats WHERE namespace = 'default'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(count, 2, "both observations under the normalized namespace");
}
fn keyed_input(text: &str, ns: &str, emb: f32, key: Option<&str>) -> RecordInput {
RecordInput {
created_at: None,
idempotency_key: key.map(|k| k.to_string()),
text: text.into(),
memory_type: "semantic".into(),
importance: 0.7,
valence: 0.0,
half_life: 604800.0,
metadata: serde_json::json!({}),
embedding: vec_seed(emb, 8),
namespace: ns.into(),
certainty: 0.8,
domain: "work".into(),
source: "user".into(),
emotional_state: None,
}
}
#[test]
fn keyed_batch_retry_returns_original_rids_and_writes_nothing() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let batch = vec![
keyed_input("kb item one", "kb_ns", 1.0, Some("kb-1")),
keyed_input("kb item two", "kb_ns", 2.0, Some("kb-2")),
];
let rids1 = db.record_batch(&batch).unwrap();
assert_eq!(rids1.len(), 2);
let count = |sql: &str| -> i64 { db.conn().query_row(sql, [], |r| r.get(0)).unwrap() };
let rows_before = count("SELECT COUNT(*) FROM memories WHERE namespace = 'kb_ns'");
let ops_before = count("SELECT COUNT(*) FROM oplog WHERE op_type = 'record'");
let stats_before: i64 = db
.conn()
.query_row(
"SELECT count FROM namespace_importance_stats WHERE namespace = 'kb_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(rows_before, 2);
assert_eq!(stats_before, 2);
let rids2 = db.record_batch(&batch).unwrap();
assert_eq!(rids2, rids1, "retry returns the ORIGINAL rids, in order");
assert_eq!(
count("SELECT COUNT(*) FROM memories WHERE namespace = 'kb_ns'"),
rows_before,
"a keyed retry must write no rows"
);
assert_eq!(
count("SELECT COUNT(*) FROM oplog WHERE op_type = 'record'"),
ops_before,
"a keyed retry must log no ops"
);
let stats_after: i64 = db
.conn()
.query_row(
"SELECT count FROM namespace_importance_stats WHERE namespace = 'kb_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
stats_after, stats_before,
"a keyed retry must not advance stats"
);
}
#[test]
fn batch_in_batch_same_key_same_payload_aliases_to_one_write() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rids = db
.record_batch(&[
keyed_input("alias item", "ib_ns", 1.0, Some("ib-dup")),
keyed_input("independent item", "ib_ns", 2.0, None),
keyed_input("alias item", "ib_ns", 1.0, Some("ib-dup")),
])
.unwrap();
assert_eq!(rids.len(), 3, "positional result for every input");
assert_eq!(rids[0], rids[2], "in-batch dup aliases to the first write");
assert_ne!(rids[0], rids[1]);
let rows: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE namespace = 'ib_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(rows, 2, "the aliased dup must not write a second row");
}
#[test]
fn batch_in_batch_same_key_divergent_payload_is_a_conflict() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let err = db
.record_batch(&[
keyed_input("first payload", "dv_ns", 1.0, Some("dv-key")),
keyed_input("DIFFERENT payload", "dv_ns", 2.0, Some("dv-key")),
])
.expect_err("divergent payloads under one key must fail the batch");
assert!(
matches!(
err,
crate::error::YantrikDbError::IdempotencyConflict { .. }
),
"expected IdempotencyConflict, got {err:?}"
);
let rows: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE namespace = 'dv_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(rows, 0, "a conflicted batch must write nothing");
}
#[test]
fn keyed_batch_duplicate_resolves_even_under_backpressure() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let dim = db.embedding_dim();
let batch = vec![
keyed_input("bp keyed one", "bp_ns", 1.0, Some("bp-1")),
keyed_input("bp keyed two", "bp_ns", 2.0, Some("bp-2")),
];
let rids1 = db.record_batch(&batch).unwrap();
let mut saturated = false;
for i in 0..400 {
let emb: Vec<f32> = (0..dim).map(|j| ((i + j) as f32) * 0.001).collect();
match db.record(
&format!("bp-filler-{i}"),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"bp_filler_ns",
0.8,
"general",
"user",
None,
) {
Ok(_) => {}
Err(crate::error::YantrikDbError::Backpressure { .. }) => {
saturated = true;
break;
}
Err(e) => panic!("unexpected: {e:?}"),
}
}
assert!(saturated, "delta never saturated");
let fresh_err = db
.record_batch(&[keyed_input("bp fresh", "bp_ns", 3.0, None)])
.expect_err("fresh batch into saturated delta must fail");
assert!(matches!(
fresh_err,
crate::error::YantrikDbError::Backpressure { .. }
));
let rids2 = db
.record_batch(&batch)
.expect("a fully-duplicate keyed batch must resolve, not backpressure");
assert_eq!(rids2, rids1);
}
#[test]
fn same_key_across_record_and_record_batch_hits_on_identical_payload() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let input = keyed_input("cross surface batch text", "xb_ns", 1.5, Some("xb-key"));
let rid = db
.record_with_idempotency(
&input.text,
&input.memory_type,
input.importance,
input.valence,
input.half_life,
&input.metadata,
&input.embedding,
&input.namespace,
input.certainty,
&input.domain,
&input.source,
input.emotional_state.as_deref(),
Some("xb-key"),
None,
)
.unwrap();
let rids = db.record_batch(&[input]).unwrap();
assert_eq!(
rids[0], rid,
"identical payload under the same key is the same write on either surface"
);
let rows: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE namespace = 'xb_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(rows, 1, "the batch hit must not write a second row");
}
#[test]
fn keyed_batch_claim_binds_to_a_committed_op() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rids = db
.record_batch(&[keyed_input(
"claim binding probe",
"cb_ns",
1.0,
Some("cb-key"),
)])
.unwrap();
let (claim_rid, op_id, state): (String, String, String) = db
.conn()
.query_row(
"SELECT rid, op_id, state FROM idempotency_claims WHERE idempotency_key = 'cb-key' AND namespace = 'cb_ns'",
[],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.expect("a keyed batch item must write a claim");
assert_eq!(claim_rid, rids[0]);
assert_eq!(state, "committed");
let (op_type, applied, target_rid): (String, i64, String) = db
.conn()
.query_row(
"SELECT op_type, applied, target_rid FROM oplog WHERE op_id = ?1",
rusqlite::params![op_id],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.expect("the claim's op_id must exist in the oplog — the claim binds to it");
assert_eq!(op_type, "record");
assert_eq!(applied, 1);
assert_eq!(target_rid, rids[0]);
}
#[test]
fn keyed_batch_duplicate_resolves_during_reembed_cutover() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let batch = vec![
keyed_input("cutover keyed one", "co_ns", 1.0, Some("co-1")),
keyed_input("cutover keyed two", "co_ns", 2.0, Some("co-2")),
];
let rids1 = db.record_batch(&batch).unwrap();
db.write_router.switch_to_queueing();
let err = db
.record_batch(&[keyed_input("cutover fresh", "co_ns", 3.0, None)])
.expect_err("fresh batch must defer during cutover");
assert!(
matches!(
err,
crate::error::YantrikDbError::BatchDeferredDuringReembed { .. }
),
"expected BatchDeferredDuringReembed, got {err:?}"
);
let rids2 = db
.record_batch(&batch)
.expect("a fully-duplicate keyed batch must resolve during cutover");
assert_eq!(rids2, rids1);
}
#[test]
fn partial_hit_batch_writes_only_the_fresh_items() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let keyed = keyed_input("partial hit anchor", "ph_ns", 1.0, Some("ph-key"));
let first = db.record_batch(std::slice::from_ref(&keyed)).unwrap();
let rids = db
.record_batch(&[
keyed.clone(),
keyed_input("partial fresh item", "ph_ns", 2.0, None),
])
.unwrap();
assert_eq!(rids.len(), 2);
assert_eq!(rids[0], first[0], "hit position carries the ORIGINAL rid");
assert_ne!(rids[1], rids[0]);
let count = |sql: &str| -> i64 { db.conn().query_row(sql, [], |r| r.get(0)).unwrap() };
assert_eq!(
count("SELECT COUNT(*) FROM memories WHERE namespace = 'ph_ns'"),
2,
"one row from each batch — the hit wrote nothing"
);
assert_eq!(
count("SELECT COUNT(*) FROM oplog WHERE op_type = 'record'"),
2,
"one record op per WRITTEN item"
);
let stats_count: i64 = db
.conn()
.query_row(
"SELECT count FROM namespace_importance_stats WHERE namespace = 'ph_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
stats_count, 2,
"two written observations, the hit added none"
);
let ns: String = db
.conn()
.query_row(
"SELECT namespace FROM memories WHERE rid = ?1",
rusqlite::params![rids[1]],
|r| r.get(0),
)
.unwrap();
assert_eq!(ns, "ph_ns");
}
#[test]
fn keyed_batch_record_op_replicates_the_key_to_a_peer() {
use crate::replication::{apply_ops, extract_ops_since};
let leader = YantrikDB::new(":memory:", 8).unwrap();
let follower = YantrikDB::new(":memory:", 8).unwrap();
let rids = leader
.record_batch(&[
keyed_input("replicated keyed item", "rk_ns", 1.0, Some("rk-key")),
keyed_input("replicated unkeyed item", "rk_ns", 2.0, None),
])
.unwrap();
let payload_for = |rid: &str| -> serde_json::Value {
let raw: String = leader
.conn()
.query_row(
"SELECT payload FROM oplog WHERE op_type = 'record' AND target_rid = ?1",
rusqlite::params![rid],
|r| r.get(0),
)
.unwrap();
serde_json::from_str(&raw).unwrap()
};
let keyed_payload = payload_for(&rids[0]);
assert_eq!(
keyed_payload["idempotency_key"], "rk-key",
"keyed batch op payload must carry the key for follower mirroring"
);
assert_eq!(
keyed_payload["origin_actor"],
leader.actor_id(),
"keyed batch op payload must carry the origin actor"
);
let unkeyed_payload = payload_for(&rids[1]);
assert!(
unkeyed_payload["idempotency_key"].is_null(),
"unkeyed items carry null, matching record()'s keyless shape"
);
let ops = extract_ops_since(&leader.conn(), None, None, None, 100).unwrap();
apply_ops(&follower, &ops).unwrap();
let (f_key, f_actor): (Option<String>, Option<String>) = follower
.conn()
.query_row(
"SELECT idempotency_key, origin_actor FROM memories WHERE rid = ?1",
rusqlite::params![rids[0]],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(
f_key.as_deref(),
Some("rk-key"),
"follower keyed row must mirror the leader's idempotency_key"
);
assert_eq!(
f_actor.as_deref(),
Some(leader.actor_id()),
"follower keyed row must mirror the leader's origin_actor"
);
}
#[allow(clippy::too_many_arguments)]
fn rwr(
db: &YantrikDB,
rid: &str,
text: &str,
ns: &str,
entities: &[&str],
seq: Option<u64>,
) -> crate::error::Result<()> {
db.record_with_rid(
rid,
text,
"semantic",
0.6,
0.0,
604800.0,
&serde_json::json!({}),
&vec_seed(4.0, 8),
ns,
0.8,
"general",
"user",
None,
1_750_000_000_000_000,
entities,
"test-embedder",
seq,
crate::provenance::WriteAdmission::Origin,
)
}
#[test]
fn record_with_rid_pending_saturation_rejects_before_writing() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.pending_op_count
.swap(1_000_000, std::sync::atomic::Ordering::SeqCst);
let err = rwr(
&db,
"0198c1c2-0000-7000-8000-0000000000aa",
"saturated pending write",
"sat_ns",
&["SaturatedEntity"],
None,
)
.expect_err("pending saturation must reject the write");
assert!(
matches!(err, crate::error::YantrikDbError::Backpressure { .. }),
"expected Backpressure, got {err:?}"
);
let rows: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE namespace = 'sat_ns'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
rows, 0,
"a write rejected for pending saturation must write NOTHING — pre-port the row and vector were already durable when the check ran"
);
let ops: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM oplog WHERE op_type = 'record_with_rid'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(ops, 0, "no op for a rejected write");
assert!(db.conn().is_autocommit(), "no savepoint left open");
}
#[test]
fn record_with_rid_replay_does_not_reenqueue_materialization() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = "0198c1c2-0000-7000-8000-0000000000bb";
rwr(&db, rid, "replayed write", "rp_ns", &["ReplayEntity"], None).unwrap();
rwr(&db, rid, "replayed write", "rp_ns", &["ReplayEntity"], None).unwrap();
let enqueues: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM oplog WHERE op_type = ?1 AND target_rid = ?2",
rusqlite::params![
crate::engine::op_types::OP_MATERIALIZE_RECORD_WITH_RID_POST,
rid
],
|r| r.get(0),
)
.unwrap();
assert_eq!(
enqueues, 1,
"replay of an existing rid re-enqueued its materialization"
);
let ops: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM oplog WHERE op_type = 'record_with_rid' AND target_rid = ?1",
rusqlite::params![rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(
ops, 1,
"exactly one record_with_rid op for one logical write"
);
}
#[test]
fn record_with_rid_enqueue_counts_exactly_one_pending_op() {
use std::sync::atomic::Ordering;
let db = YantrikDB::new(":memory:", 8).unwrap();
let before = db.pending_op_count.load(Ordering::SeqCst);
rwr(
&db,
"0198c1c2-0000-7000-8000-0000000000cc",
"counted enqueue",
"cnt_ns",
&["CountedEntity"],
None,
)
.unwrap();
assert_eq!(
db.pending_op_count.load(Ordering::SeqCst),
before + 1,
"one enqueued materialization op = exactly one count"
);
rwr(
&db,
"0198c1c2-0000-7000-8000-0000000000cd",
"uncounted write",
"cnt_ns",
&[],
None,
)
.unwrap();
assert_eq!(
db.pending_op_count.load(Ordering::SeqCst),
before + 1,
"an entity-less write enqueues nothing and must not count"
);
rwr(
&db,
"0198c1c2-0000-7000-8000-0000000000cc",
"counted enqueue",
"cnt_ns",
&["CountedEntity"],
None,
)
.unwrap();
assert_eq!(
db.pending_op_count.load(Ordering::SeqCst),
before + 1,
"a replay enqueues nothing and must not count"
);
}