use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use khive_storage::note::Note;
use khive_types::Namespace;
use lattice_embed::{EmbedError, EmbeddingModel, EmbeddingService};
use serde_json::json;
use tokio::sync::Notify;
use crate::atomic_prepare::apply_post_commit_effects_with_report;
use crate::atomic_runner::{run_atomic_unit, AtomicOpPlan, AtomicRunOutcome};
use crate::curation::NotePatch;
use crate::embedder_registry::EmbedderProvider;
use crate::note_write::{NoteFence, NoteWriteOptions};
use crate::{KhiveRuntime, NamespaceToken, RuntimeError, RuntimeResult};
const MODEL: &str = "note-version-test";
#[test]
fn note_fence_explicit_absence_and_alias_serialize_canonically() {
use crate::note_write::NoteFences;
for expected in [None, Some(1), Some(i64::MAX)] {
for field in ["expected_version", "version"] {
let mut value = json!({"kind":"head", "key":"fence/serde"});
value[field] = json!(expected);
let fence: NoteFence = serde_json::from_value(value.clone()).unwrap();
assert_eq!(fence.expected_version, expected);
fence.validate().unwrap();
let canonical = json!({
"kind":"head", "key":"fence/serde", "expected_version":expected
});
assert_eq!(serde_json::to_value(&fence).unwrap(), canonical);
let round_trip: NoteFence =
serde_json::from_value(serde_json::to_value(&fence).unwrap()).unwrap();
assert_eq!(round_trip.expected_version, expected);
for listed in [false, true] {
let input = if listed {
json!([value.clone()])
} else {
value.clone()
};
let fences: NoteFences = serde_json::from_value(input).unwrap();
assert_eq!(fences.entries()[0].expected_version, expected);
let serialized = if listed {
json!([canonical.clone()])
} else {
canonical.clone()
};
assert_eq!(serde_json::to_value(fences).unwrap(), serialized);
}
}
}
}
#[test]
fn note_fence_requires_present_version_and_rejects_invalid_forms() {
use crate::note_write::NoteFences;
let missing = json!({"kind":"head", "key":"fence/serde"});
let direct = serde_json::from_value::<NoteFence>(missing.clone()).unwrap_err();
assert!(direct
.to_string()
.contains("fence requires expected_version (positive integer or null)"));
let mut invalid = vec![missing];
for field in ["expected_version", "version"] {
for value in [
json!(0),
json!(-1),
json!(false),
json!("1"),
json!(1.5),
json!([]),
json!({}),
] {
let mut fence = json!({"kind":"head", "key":"fence/serde"});
fence[field] = value;
invalid.push(fence);
}
}
for fence in invalid {
for input in [fence.clone(), json!([fence])] {
let message = serde_json::from_value::<NoteFences>(input)
.unwrap_err()
.to_string();
assert!(
message.contains("fence requires expected_version (positive integer or null)"),
"{message}"
);
assert!(!message.contains("Shape"), "{message}");
assert!(!message.contains("untagged enum"), "{message}");
}
}
for value in [json!(null), json!(1)] {
let duplicate = json!({
"kind":"head", "key":"fence/serde",
"expected_version":value, "version":value
});
let direct = serde_json::from_value::<NoteFence>(duplicate.clone()).unwrap_err();
assert!(direct.to_string().contains("duplicate field"), "{direct}");
for input in [duplicate.clone(), json!([duplicate])] {
let message = serde_json::from_value::<NoteFences>(input)
.unwrap_err()
.to_string();
assert!(message.contains("duplicate field"), "{message}");
assert!(message.contains("expected_version"), "{message}");
}
}
let unknown = json!({
"kind":"head", "key":"fence/serde", "expected_version":null, "extra":true
});
for input in [unknown.clone(), json!([unknown])] {
let message = serde_json::from_value::<NoteFences>(input)
.unwrap_err()
.to_string();
assert!(message.contains("unknown field"), "{message}");
assert!(message.contains("extra"), "{message}");
}
for input in [
json!(null),
json!([null]),
json!([]),
json!(true),
json!("fence"),
] {
let message = serde_json::from_value::<NoteFences>(input)
.unwrap_err()
.to_string();
assert!(message.contains("fence"), "{message}");
assert!(!message.contains("Shape"), "{message}");
assert!(!message.contains("untagged enum"), "{message}");
}
}
#[test]
fn absence_fence_conflict_keeps_canonical_details_and_version_cas_shape() {
use crate::note_write::NoteWriteConflict;
for index in [None, Some(0)] {
for member in [None, Some(2)] {
let error = NoteWriteConflict::Fence {
key: "fence/serde".into(),
expected: None,
current: Some(3),
index,
}
.into_error_at_member(member);
let details = serde_json::to_value(error.details().unwrap()).unwrap();
let mut expected = json!({
"reason":"fence_conflict", "key":"fence/serde",
"expected_version":"absent", "current_version":"3"
});
if let Some(index) = index {
expected["index"] = json!(index.to_string());
}
if let Some(member) = member {
expected["member"] = json!(member.to_string());
}
assert_eq!(details, expected);
}
}
let error = NoteWriteConflict::Version {
expected: 1,
current: 2,
}
.into_error();
assert_eq!(
serde_json::to_value(error.details().unwrap()).unwrap(),
json!({"reason":"version_conflict", "expected_version":"1", "current_version":"2"})
);
}
#[test]
fn ordered_fences_cap_precedes_entry_validation_for_json_and_typed_inputs() {
use crate::note_write::NoteFences;
let entries: Vec<NoteFence> = (0..100)
.map(|index| NoteFence {
key: format!("fence-cap/{index}"),
kind: "head".into(),
expected_version: Some(1),
live_until: None,
id: None,
})
.collect();
let at_cap = NoteFences::Many(entries.clone());
at_cap.validate().unwrap();
serde_json::from_value::<NoteFences>(serde_json::to_value(&at_cap).unwrap()).unwrap();
let mut over_cap = entries;
over_cap.push(NoteFence {
key: "fence-cap/100".into(),
kind: "head".into(),
expected_version: Some(1),
live_until: None,
id: None,
});
for malformed in [false, true] {
let mut entries = over_cap.clone();
if malformed {
entries[0].expected_version = Some(0);
}
let fences = NoteFences::Many(entries);
let error = fences.validate().unwrap_err();
assert!(matches!(error, RuntimeError::InvalidInput(_)), "{error}");
for message in [
error.to_string(),
serde_json::from_value::<NoteFences>(serde_json::to_value(&fences).unwrap())
.unwrap_err()
.to_string(),
] {
assert!(message.contains("at most 100 entries"), "{message}");
assert!(message.contains("sent 101"), "{message}");
}
}
let malformed = json!(vec![serde_json::Value::Null; 101]);
let message = serde_json::from_value::<NoteFences>(malformed)
.unwrap_err()
.to_string();
assert!(message.contains("at most 100 entries"), "{message}");
assert!(message.contains("sent 101"), "{message}");
}
#[tokio::test]
async fn ordered_fences_observe_prior_write_in_same_transaction() {
let (runtime, token, _) = fixture();
let lease_a = create(&runtime, &token, "lease/a", None).await;
let lease_b = create(&runtime, &token, "lease/b", None).await;
let target = create(&runtime, &token, "target", None).await;
let lease_update = crate::atomic_prepare::prepare_update(
&runtime,
&token,
&json!({"id":lease_b.id,"content":"{\"renewed\":true}","expected_version":1}),
None,
)
.await
.unwrap();
let target_update = crate::atomic_prepare::prepare_update(
&runtime,
&token,
&json!({"id":target.id,"content":"{\"bad\":true}","fence":[
{"kind":"head","key":"lease/a","expected_version":1},
{"kind":"head","key":"lease/b","expected_version":1}]}),
None,
)
.await
.unwrap();
let outcome = run_atomic_unit(runtime.sql().as_ref(), vec![lease_update, target_update])
.await
.unwrap();
let AtomicRunOutcome::RolledBack {
failed_op_index,
failure: crate::atomic_runner::AtomicOpFailure::NoteConflict(conflict),
..
} = outcome
else {
panic!("fence must see the earlier update and roll back: {outcome:?}");
};
assert_eq!(failed_op_index, 1);
let detail = serde_json::to_value(conflict.into_error().details().unwrap()).unwrap();
assert_eq!(detail["index"], "1");
assert_eq!(detail["current_version"], "2");
for note in [lease_a, lease_b, target] {
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(note.id)
.await
.unwrap()
.unwrap(),
note
);
}
}
#[tokio::test]
async fn absence_fence_observes_prior_create_in_same_transaction_and_rolls_back() {
use crate::atomic_message::{AtomicNoteOptions, AtomicNoteSpec};
use crate::note_create::{prepare_note_create, KeyPublication};
use crate::note_write::NoteFences;
for listed in [false, true] {
let (runtime, token, _) = fixture();
let target = create(&runtime, &token, "absence/target", None).await;
let (mut prepared, _) = prepare_note_create(
&runtime,
AtomicNoteSpec {
token: &token,
id: None,
kind: "head",
name: None,
content: "{}",
properties: None,
},
AtomicNoteOptions {
key: Some("absence/holder"),
embed: Some(false),
..Default::default()
},
&[],
KeyPublication::AtInsert,
)
.await
.unwrap();
let holder_id = prepared.notes[0].id;
let fence = NoteFence {
key: "absence/holder".into(),
kind: "head".into(),
expected_version: None,
live_until: None,
id: None,
};
let mut update = patch("{\"changed\":true}", 1, None);
update.write_options.fence = Some(if listed {
NoteFences::Many(vec![fence])
} else {
fence.into()
});
let (_, target_plan) = runtime
.prepare_versioned_note_update(&token, target.clone(), update.clone())
.await
.unwrap();
let target_op_index = prepared.plans.len();
prepared
.plans
.push(AtomicOpPlan::Update(Box::new(target_plan)));
let outcome = run_atomic_unit(runtime.sql().as_ref(), prepared.plans)
.await
.unwrap();
let AtomicRunOutcome::RolledBack {
failed_op_index,
failure: crate::atomic_runner::AtomicOpFailure::NoteConflict(conflict),
..
} = outcome
else {
panic!("absence fence must see the earlier create and roll back: {outcome:?}");
};
assert_eq!(failed_op_index, target_op_index);
let details = serde_json::to_value(conflict.into_error().details().unwrap()).unwrap();
let mut expected = json!({
"reason":"fence_conflict", "key":"absence/holder",
"expected_version":"absent", "current_version":"1"
});
if listed {
expected["index"] = json!("0");
}
assert_eq!(details, expected);
assert!(runtime
.notes(&token)
.unwrap()
.get_note(holder_id)
.await
.unwrap()
.is_none());
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(target.id)
.await
.unwrap()
.unwrap(),
target
);
assert_eq!(
runtime
.update_note(&token, target.id, update)
.await
.unwrap()
.version,
2
);
}
}
#[derive(Default)]
struct Service {
started: Notify,
proceed: Notify,
fail: AtomicBool,
calls: AtomicUsize,
}
#[async_trait]
impl EmbeddingService for Service {
async fn embed(
&self,
texts: &[String],
_: EmbeddingModel,
) -> Result<Vec<Vec<f32>>, EmbedError> {
self.calls.fetch_add(1, Ordering::SeqCst);
if texts.iter().any(|text| text.contains("inflight")) {
self.started.notify_one();
self.proceed.notified().await;
}
if self.fail.load(Ordering::SeqCst) {
return Err(EmbedError::InferenceFailed("late test failure".into()));
}
Ok(texts.iter().map(|_| vec![0.5; 4]).collect())
}
fn supports_model(&self, _: EmbeddingModel) -> bool {
true
}
fn name(&self) -> &'static str {
MODEL
}
}
struct Provider(Arc<Service>, String, usize);
#[async_trait]
impl EmbedderProvider for Provider {
fn name(&self) -> &str {
&self.1
}
fn dimensions(&self) -> usize {
self.2
}
async fn build(&self) -> RuntimeResult<Arc<dyn EmbeddingService>> {
Ok(self.0.clone())
}
}
fn fixture() -> (KhiveRuntime, NamespaceToken, Arc<Service>) {
let runtime = KhiveRuntime::memory().unwrap();
runtime.install_kind_registry(
vec!["concept".into()],
vec!["head".into(), "memory".into(), "observation".into()],
);
let token = runtime
.authorize(Namespace::parse("local").unwrap())
.unwrap();
let service = Arc::new(Service::default());
runtime.register_embedder(Provider(service.clone(), MODEL.into(), 4));
runtime.vectors_for_model(&token, MODEL).unwrap();
(runtime, token, service)
}
async fn create(
runtime: &KhiveRuntime,
token: &NamespaceToken,
key: &str,
embed: Option<bool>,
) -> Note {
runtime
.create_note_with_options(
token,
"head",
None,
"{}",
None,
None,
None,
None,
vec![],
None,
NoteWriteOptions {
key: Some(key.into()),
embed,
..Default::default()
},
)
.await
.unwrap()
.0
}
fn patch(content: &str, version: i64, embed: Option<bool>) -> NotePatch {
NotePatch::new(None, Some(content.into()), None, None, None).with_write_options(
NoteWriteOptions {
expected_version: Some(version),
embed,
..Default::default()
},
)
}
fn details(error: RuntimeError) -> serde_json::Value {
let RuntimeError::Khive(error) = error else {
panic!("expected structured error, got {error:?}");
};
serde_json::to_value(error.details().expect("conflict details")).unwrap()
}
async fn vectors(runtime: &KhiveRuntime, token: &NamespaceToken) -> u64 {
runtime
.vectors_for_model(token, MODEL)
.unwrap()
.count()
.await
.unwrap()
}
async fn ann_deletes(runtime: &KhiveRuntime, id: uuid::Uuid) -> i64 {
let value = runtime
.sql()
.reader()
.await
.unwrap()
.query_scalar(crate::note_write::statement(
"SELECT COUNT(*) FROM ann_write_log WHERE subject_id=?1 AND op='delete'",
vec![khive_storage::SqlValue::Text(id.to_string())],
))
.await
.unwrap()
.unwrap();
let khive_storage::SqlValue::Integer(count) = value else {
panic!("invalid count")
};
count
}
async fn seed_attributed_note_vector(
runtime: &KhiveRuntime,
token: &NamespaceToken,
note: &Note,
) -> khive_storage::ContentRef {
use khive_storage::types::VectorRecord;
let fingerprint = VectorRecord::fingerprint_text("known note vector input");
runtime
.vectors_for_model(token, MODEL)
.unwrap()
.insert_batch(vec![VectorRecord {
subject_id: note.id,
kind: khive_types::SubstrateKind::Note,
namespace: note.namespace.clone(),
field: "note.content".into(),
embedding_model: Some(MODEL.into()),
vectors: vec![vec![0.5; 4]],
text_fingerprint: Some(fingerprint.clone()),
updated_at: chrono::Utc::now(),
}])
.await
.unwrap();
fingerprint
}
async fn note_sidecar_count(runtime: &KhiveRuntime, note: &Note) -> i64 {
let result = runtime
.sql()
.reader()
.await
.unwrap()
.query_scalar(crate::note_write::statement(
"SELECT COUNT(*) FROM vector_provenance \
WHERE model_key=?1 AND namespace=?2 AND subject_id=?3",
vec![
khive_storage::SqlValue::Text(crate::config::sanitize_key(MODEL)),
khive_storage::SqlValue::Text(note.namespace.clone()),
khive_storage::SqlValue::Text(note.id.to_string()),
],
))
.await
.unwrap();
let Some(khive_storage::SqlValue::Integer(count)) = result else {
panic!("expected sidecar count, got {result:?}");
};
count
}
async fn note_sidecar_snapshot(runtime: &KhiveRuntime, note: &Note) -> serde_json::Value {
let rows = runtime
.sql()
.reader()
.await
.unwrap()
.query_all(crate::note_write::statement(
"SELECT model_key, subject_id, namespace, embedding_digest, \
text_fingerprint, updated_at \
FROM vector_provenance WHERE model_key=?1 AND subject_id=?2",
vec![
khive_storage::SqlValue::Text(crate::config::sanitize_key(MODEL)),
khive_storage::SqlValue::Text(note.id.to_string()),
],
))
.await
.unwrap();
serde_json::to_value(rows).unwrap()
}
async fn note_vector_blob_hex(runtime: &KhiveRuntime, note: &Note) -> String {
let result = runtime
.sql()
.reader()
.await
.unwrap()
.query_scalar(crate::note_write::statement(
"SELECT hex(embedding) FROM vec_note_version_test \
WHERE namespace=?1 AND subject_id=?2",
vec![
khive_storage::SqlValue::Text(note.namespace.clone()),
khive_storage::SqlValue::Text(note.id.to_string()),
],
))
.await
.unwrap();
let Some(khive_storage::SqlValue::Text(blob)) = result else {
panic!("expected live vector BLOB, got {result:?}");
};
blob
}
#[tokio::test]
async fn atomic_message_same_blob_note_reindex_clears_provenance() {
let (runtime, token, service) = fixture();
let note = create(&runtime, &token, "version/same-blob-reindex", Some(false)).await;
let fingerprint = seed_attributed_note_vector(&runtime, &token, ¬e).await;
let vectors = runtime.vectors_for_model(&token, MODEL).unwrap();
let before_blob = note_vector_blob_hex(&runtime, ¬e).await;
assert_eq!(note_sidecar_count(&runtime, ¬e).await, 1);
let before = vectors.provenance(note.id).await.unwrap().unwrap();
assert_eq!(before.text_fingerprint, Some(fingerprint));
assert!(before.updated_at.is_some());
let updated = runtime
.update_note(
&token,
note.id,
patch("{\"reindexed\":true}", 1, Some(true)),
)
.await
.unwrap();
assert_eq!(updated.version, 2);
assert_eq!(service.calls.load(Ordering::SeqCst), 1);
assert_eq!(vectors.count().await.unwrap(), 1);
assert_eq!(note_vector_blob_hex(&runtime, &updated).await, before_blob);
assert_eq!(note_sidecar_count(&runtime, &updated).await, 0);
let after = vectors.provenance(note.id).await.unwrap().unwrap();
assert_eq!(after.text_fingerprint, None);
assert_eq!(after.updated_at, None);
}
#[tokio::test]
async fn atomic_message_subject_only_move_clears_old_namespace_provenance() {
let (runtime, token, _service) = fixture();
let note = create(&runtime, &token, "version/raw-note-move", Some(false)).await;
seed_attributed_note_vector(&runtime, &token, ¬e).await;
assert_eq!(note_sidecar_count(&runtime, ¬e).await, 1);
let table = format!("vec_{}", crate::config::sanitize_key(MODEL));
let statements = crate::atomic_message::vector_insert_statements(
&table,
"other",
note.id,
"note.content",
MODEL,
&[0.5; 4],
"test-raw-note-move",
);
runtime
.sql()
.writer()
.await
.unwrap()
.execute_batch(statements.into_iter().map(|step| step.statement).collect())
.await
.unwrap();
assert_eq!(note_sidecar_count(&runtime, ¬e).await, 0);
let moved = runtime
.sql()
.reader()
.await
.unwrap()
.query_scalar(crate::note_write::statement(
format!("SELECT COUNT(*) FROM {table} WHERE namespace=?1 AND subject_id=?2"),
vec![
khive_storage::SqlValue::Text("other".into()),
khive_storage::SqlValue::Text(note.id.to_string()),
],
))
.await
.unwrap();
assert!(matches!(moved, Some(khive_storage::SqlValue::Integer(1))));
}
#[tokio::test]
async fn version_guarded_vector_publication_canonicalizes_builtin_aliases() {
let (runtime, token, service) = fixture();
let model = EmbeddingModel::ParaphraseMultilingualMiniLmL12V2;
runtime.register_embedder(Provider(service, model.to_string(), model.dimensions()));
let note = create(&runtime, &token, "version/model-alias", None).await;
let vector = vec![0.5; model.dimensions()];
assert!(runtime
.publish_note_vector_revision(&token, ¬e, "paraphrase", &vector)
.await
.unwrap());
assert_eq!(
runtime
.vectors_for_model(&token, &model.to_string())
.unwrap()
.count()
.await
.unwrap(),
1
);
let identity = runtime
.sql()
.reader()
.await
.unwrap()
.query_scalar(crate::note_write::statement(
"SELECT embedding_model FROM ann_write_log WHERE subject_id=?1 AND op='upsert'",
vec![khive_storage::SqlValue::Text(note.id.to_string())],
))
.await
.unwrap();
assert!(
matches!(identity, Some(khive_storage::SqlValue::Text(name)) if name == model.to_string())
);
}
#[tokio::test]
async fn version_guarded_vector_publication_checks_declared_dimensions() {
let (runtime, token, service) = fixture();
let note = create(&runtime, &token, "version/model-dimensions", None).await;
runtime.register_embedder(Provider(service, MODEL.into(), 8));
let error = runtime
.publish_note_vector_revision(&token, ¬e, MODEL, &[0.5; 4])
.await
.unwrap_err();
assert!(
error.to_string().contains("expected 8 vector dimensions"),
"{error}"
);
assert_eq!(vectors(&runtime, &token).await, 0);
}
#[tokio::test]
async fn version_creation_compensation_preserves_or_removes_attachments_with_revision() {
use khive_storage::attachment::{Attachment, AttachmentSubstrate};
let (runtime, token, _) = fixture();
for newer in [false, true] {
let note = create(
&runtime,
&token,
&format!("version/attachment-{newer}"),
None,
)
.await;
let attachment = Attachment {
record_uuid: note.id,
substrate: AttachmentSubstrate::Note,
role: "source".into(),
content_ref: khive_storage::ContentRef::from_hex("a".repeat(64)).unwrap(),
media_type: None,
size_bytes: None,
created_at: 1,
};
runtime
.sql()
.writer()
.await
.unwrap()
.execute(
khive_db::stores::attachment::attachment_upsert_statement(&attachment).unwrap(),
)
.await
.unwrap();
if newer {
runtime
.update_note(&token, note.id, patch("{\"newer\":true}", 1, None))
.await
.unwrap();
}
assert_eq!(runtime.compensate_note_creation(¬e).await, !newer);
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(note.id)
.await
.unwrap()
.is_some(),
newer
);
let retained = runtime
.sql()
.reader()
.await
.unwrap()
.query_scalar(crate::note_write::statement(
"SELECT COUNT(*) FROM attachments WHERE record_uuid=?1",
vec![khive_storage::SqlValue::Text(note.id.to_string())],
))
.await
.unwrap();
assert!(
matches!(retained, Some(khive_storage::SqlValue::Integer(count)) if count == i64::from(newer))
);
}
}
#[tokio::test]
async fn note_creation_compensation_clears_provenance() {
let (runtime, token, _) = fixture();
let note = create(
&runtime,
&token,
"version/compensate-provenance",
Some(false),
)
.await;
let fingerprint = seed_attributed_note_vector(&runtime, &token, ¬e).await;
let vectors = runtime.vectors_for_model(&token, MODEL).unwrap();
assert_eq!(note_sidecar_count(&runtime, ¬e).await, 1);
assert_eq!(
vectors
.provenance(note.id)
.await
.unwrap()
.unwrap()
.text_fingerprint,
Some(fingerprint)
);
assert!(runtime.compensate_note_creation(¬e).await);
assert!(runtime
.notes(&token)
.unwrap()
.get_note(note.id)
.await
.unwrap()
.is_none());
assert_eq!(note_sidecar_count(&runtime, ¬e).await, 0);
assert!(vectors.provenance(note.id).await.unwrap().is_none());
let stale = create(
&runtime,
&token,
"version/stale-compensate-provenance",
Some(false),
)
.await;
runtime
.update_note(&token, stale.id, patch("{\"newer\":true}", 1, None))
.await
.unwrap();
let stale_fingerprint = seed_attributed_note_vector(&runtime, &token, &stale).await;
assert!(!runtime.compensate_note_creation(&stale).await);
assert_eq!(note_sidecar_count(&runtime, &stale).await, 1);
assert_eq!(
vectors
.provenance(stale.id)
.await
.unwrap()
.unwrap()
.text_fingerprint,
Some(stale_fingerprint)
);
}
#[tokio::test]
async fn note_vector_purge_scopes_sidecar_to_model_and_namespace() {
let (runtime, token, _) = fixture();
let note = create(&runtime, &token, "version/purge-scope", Some(false)).await;
seed_attributed_note_vector(&runtime, &token, ¬e).await;
let mut writer = runtime.sql().writer().await.unwrap();
writer
.execute(crate::note_write::statement(
"UPDATE vector_provenance SET namespace='foreign' \
WHERE model_key=?1 AND subject_id=?2",
vec![
khive_storage::SqlValue::Text(crate::config::sanitize_key(MODEL)),
khive_storage::SqlValue::Text(note.id.to_string()),
],
))
.await
.unwrap();
writer
.execute(crate::note_write::statement(
"INSERT INTO vector_provenance \
(model_key, subject_id, namespace, embedding_digest, text_fingerprint, updated_at) \
VALUES ('unrelated_model', ?1, 'local', ?2, NULL, NULL)",
vec![
khive_storage::SqlValue::Text(note.id.to_string()),
khive_storage::SqlValue::Text("0".repeat(64)),
],
))
.await
.unwrap();
drop(writer);
runtime
.update_note(&token, note.id, patch("{\"off\":true}", 1, Some(false)))
.await
.unwrap();
assert_eq!(vectors(&runtime, &token).await, 0);
let remaining = runtime
.sql()
.reader()
.await
.unwrap()
.query_all(crate::note_write::statement(
"SELECT model_key, namespace FROM vector_provenance WHERE subject_id=?1 \
ORDER BY model_key",
vec![khive_storage::SqlValue::Text(note.id.to_string())],
))
.await
.unwrap();
assert_eq!(remaining.len(), 2);
assert!(matches!(
remaining[0].get("model_key"),
Some(khive_storage::SqlValue::Text(key)) if key == &crate::config::sanitize_key(MODEL)
));
assert!(matches!(
remaining[0].get("namespace"),
Some(khive_storage::SqlValue::Text(namespace)) if namespace == "foreign"
));
assert!(matches!(
remaining[1].get("model_key"),
Some(khive_storage::SqlValue::Text(key)) if key == "unrelated_model"
));
}
#[tokio::test]
async fn version_keyed_create_late_annotation_refusal_is_not_a_key_conflict() {
use crate::atomic_message::{AtomicNoteOptions, AtomicNoteSpec};
use crate::note_create::{prepare_note_create, KeyPublication};
let (runtime, token, _) = fixture();
let target = create(&runtime, &token, "version/annotation-target", None).await;
let (prepared, _) = prepare_note_create(
&runtime,
AtomicNoteSpec {
token: &token,
id: None,
kind: "head",
name: None,
content: "{}",
properties: None,
},
AtomicNoteOptions {
key: Some("version/annotation-child"),
embed: Some(false),
..Default::default()
},
&[target.id],
KeyPublication::AtInsert,
)
.await
.unwrap();
runtime
.notes(&token)
.unwrap()
.delete_note(target.id, khive_storage::DeleteMode::Hard)
.await
.unwrap();
let id = prepared.notes[0].id;
let outcome = run_atomic_unit(runtime.sql().as_ref(), prepared.plans)
.await
.unwrap();
assert!(
matches!(
outcome,
AtomicRunOutcome::RolledBack {
failure: crate::atomic_runner::AtomicOpFailure::GuardFailed { observed: 0, .. },
..
}
),
"annotation refusal must not disclose the rolled-back candidate as key holder: {outcome:?}"
);
assert!(runtime
.notes(&token)
.unwrap()
.get_note(id)
.await
.unwrap()
.is_none());
}
#[tokio::test]
async fn version_cas_rejects_already_stale_and_lost_ack_retry() {
let (runtime, token, _) = fixture();
let note = create(&runtime, &token, "version/cas", None).await;
assert_eq!(note.version, 1);
let updated = runtime
.update_note(&token, note.id, patch("{\"winner\":true}", 1, None))
.await
.unwrap();
assert_eq!(updated.version, 2);
for body in ["{\"loser\":true}", "{\"winner\":true}"] {
let error = runtime
.update_note(&token, note.id, patch(body, 1, None))
.await
.unwrap_err();
assert_eq!(
details(error),
json!({"reason":"version_conflict", "expected_version":"1", "current_version":"2"})
);
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(note.id)
.await
.unwrap()
.unwrap(),
updated
);
}
assert_eq!(vectors(&runtime, &token).await, 0);
}
#[tokio::test]
async fn version_fence_and_prior_operation_roll_back_together() {
let (runtime, token, _) = fixture();
let fence = create(&runtime, &token, "version/fence", None).await;
let target = create(&runtime, &token, "version/target", None).await;
for (key, expected, current) in [
("version/fence", 2, Some("1")),
("version/missing", 1, None),
] {
let mut update = patch("{\"changed\":true}", 1, None);
update.write_options.fence = Some(
NoteFence {
key: key.into(),
kind: "head".into(),
expected_version: Some(expected),
live_until: None,
id: None,
}
.into(),
);
let error = details(
runtime
.update_note(&token, target.id, update)
.await
.unwrap_err(),
);
assert_eq!(error["reason"], "fence_conflict");
assert_eq!(
error
.get("current_version")
.and_then(|value| value.as_str()),
current
);
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(target.id)
.await
.unwrap()
.unwrap(),
target
);
}
let (_, fence_plan) = runtime
.prepare_versioned_note_update(&token, fence.clone(), patch("{\"renewed\":true}", 1, None))
.await
.unwrap();
let mut update = patch("{\"changed\":true}", 1, None);
update.write_options.fence = Some(
NoteFence {
key: fence.key.clone().unwrap(),
kind: "head".into(),
expected_version: Some(1),
live_until: None,
id: None,
}
.into(),
);
let (_, target_plan) = runtime
.prepare_versioned_note_update(&token, target.clone(), update.clone())
.await
.unwrap();
let outcome = run_atomic_unit(
runtime.sql().as_ref(),
vec![
AtomicOpPlan::Update(Box::new(fence_plan)),
AtomicOpPlan::Update(Box::new(target_plan)),
],
)
.await
.unwrap();
assert!(matches!(
outcome,
AtomicRunOutcome::RolledBack {
failed_op_index: 1,
..
}
));
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(fence.id)
.await
.unwrap()
.unwrap(),
fence
);
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(target.id)
.await
.unwrap()
.unwrap(),
target
);
assert_eq!(
runtime
.update_note(&token, target.id, update)
.await
.unwrap()
.version,
2
);
}
#[tokio::test]
async fn version_embedding_defaults_and_off_on_transitions() {
async fn assert_surfaces(
runtime: &KhiveRuntime,
token: &NamespaceToken,
note: &Note,
embedded: bool,
) {
let candidates = runtime
.vectors_for_model(token, MODEL)
.unwrap()
.search(khive_storage::VectorSearchRequest {
query_vectors: vec![vec![0.5; 4]],
top_k: 10,
namespace: Some(token.namespace().as_str().to_string()),
kind: Some(khive_types::SubstrateKind::Note),
embedding_model: Some(MODEL.into()),
filter: None,
backend_hints: None,
})
.await
.unwrap();
assert_eq!(
candidates.iter().any(|hit| hit.subject_id == note.id),
embedded,
"similarity candidacy must follow the committed embedding state"
);
let lexical = runtime
.text_for_notes(token)
.unwrap()
.search(khive_storage::TextSearchRequest {
query: "phase".into(),
mode: khive_storage::TextQueryMode::Plain,
filter: None,
top_k: 10,
snippet_chars: 100,
})
.await
.unwrap();
assert!(
lexical.iter().any(|hit| hit.subject_id == note.id),
"embedding changes must retain lexical search results"
);
assert_eq!(
runtime
.list_notes(token, Some("head"), 10, 0)
.await
.unwrap(),
vec![note.clone()],
"embedding changes must retain the current note in listing"
);
}
let (runtime, token, _) = fixture();
let note = create(&runtime, &token, "version/embed", None).await;
assert_eq!(vectors(&runtime, &token).await, 0);
let note = runtime
.update_note(&token, note.id, patch("{\"phase\":1}", 1, None))
.await
.unwrap();
assert_eq!(vectors(&runtime, &token).await, 0);
let note = runtime
.update_note(
&token,
note.id,
patch("{\"phase\":2}", note.version, Some(true)),
)
.await
.unwrap();
assert_eq!(vectors(&runtime, &token).await, 1);
assert_surfaces(&runtime, &token, ¬e, true).await;
let note = runtime
.update_note(&token, note.id, patch("{\"phase\":3}", note.version, None))
.await
.unwrap();
assert_eq!(vectors(&runtime, &token).await, 1);
let note = runtime
.update_note(
&token,
note.id,
patch("{\"phase\":4}", note.version, Some(false)),
)
.await
.unwrap();
assert_eq!(vectors(&runtime, &token).await, 0);
assert_surfaces(&runtime, &token, ¬e, false).await;
assert_eq!(
runtime
.get_note_by_key(&token, "version/embed", Some("head"), false)
.await
.unwrap(),
note
);
let note = runtime
.update_note(
&token,
note.id,
patch("{\"phase\":5}", note.version, Some(true)),
)
.await
.unwrap();
assert_eq!(vectors(&runtime, &token).await, 1);
assert_surfaces(&runtime, &token, ¬e, true).await;
}
#[tokio::test]
async fn version_delayed_reindex_cannot_reverse_embed_off() {
let (runtime, token, _) = fixture();
let note = create(&runtime, &token, "version/delayed", Some(true)).await;
let (_, plan) = runtime
.prepare_versioned_note_update(&token, note.clone(), patch("{\"old\":true}", 1, Some(true)))
.await
.unwrap();
let AtomicRunOutcome::Committed { post_commit } = run_atomic_unit(
runtime.sql().as_ref(),
vec![AtomicOpPlan::Update(Box::new(plan))],
)
.await
.unwrap() else {
panic!("commit");
};
runtime
.update_note(&token, note.id, patch("{\"off\":true}", 2, Some(false)))
.await
.unwrap();
assert!(
apply_post_commit_effects_with_report(&runtime, &token, post_commit)
.await
.unwrap()
.is_empty()
);
assert_eq!(vectors(&runtime, &token).await, 0);
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(note.id)
.await
.unwrap()
.unwrap()
.content,
"{\"off\":true}"
);
}
#[tokio::test]
async fn version_embedding_inheritance_includes_retired_model_rows() {
assert_writer_time_embedding_inheritance(true).await;
}
#[tokio::test]
async fn version_embedding_inheritance_sees_publication_after_prepare() {
assert_writer_time_embedding_inheritance(false).await;
}
#[tokio::test]
async fn note_embedding_inheritance_preserves_provenance_without_delete() {
let (runtime, token, _) = fixture();
let note = create(&runtime, &token, "version/inherit-provenance", Some(false)).await;
let fingerprint = seed_attributed_note_vector(&runtime, &token, ¬e).await;
let vectors = runtime.vectors_for_model(&token, MODEL).unwrap();
let before_blob = note_vector_blob_hex(&runtime, ¬e).await;
let before_sidecar = note_sidecar_snapshot(&runtime, ¬e).await;
let (_, plan) = runtime
.prepare_versioned_note_update(&token, note.clone(), patch("{\"new\":true}", 1, None))
.await
.unwrap();
let AtomicRunOutcome::Committed { post_commit } = run_atomic_unit(
runtime.sql().as_ref(),
vec![AtomicOpPlan::Update(Box::new(plan))],
)
.await
.unwrap() else {
panic!("commit");
};
assert_eq!(
post_commit.as_slice(),
&[crate::PostCommitEffect::ReindexNote {
note_id: note.id,
version: 2
}]
);
assert_eq!(note_sidecar_count(&runtime, ¬e).await, 1);
assert_eq!(note_vector_blob_hex(&runtime, ¬e).await, before_blob);
assert_eq!(note_sidecar_snapshot(&runtime, ¬e).await, before_sidecar);
assert_eq!(
vectors
.provenance(note.id)
.await
.unwrap()
.unwrap()
.text_fingerprint,
Some(fingerprint)
);
let applied = apply_post_commit_effects_with_report(&runtime, &token, post_commit)
.await
.unwrap();
assert_eq!(applied.len(), 1);
assert_eq!(note_sidecar_count(&runtime, ¬e).await, 0);
assert_eq!(note_vector_blob_hex(&runtime, ¬e).await, before_blob);
let replaced = vectors.provenance(note.id).await.unwrap().unwrap();
assert_eq!(replaced.text_fingerprint, None);
assert_eq!(replaced.updated_at, None);
}
async fn assert_writer_time_embedding_inheritance(retired: bool) {
let (runtime, token, _) = fixture();
let note = create(&runtime, &token, "version/inherited", None).await;
let (_, plan) = runtime
.prepare_versioned_note_update(&token, note.clone(), patch("{\"new\":true}", 1, None))
.await
.unwrap();
if retired {
runtime
.backend()
.vectors_for_namespace("retired_chunks", "retired", 4, "local")
.unwrap()
.insert(
note.id,
khive_types::SubstrateKind::Note,
"local",
"note.content",
vec![vec![0.1; 4]],
)
.await
.unwrap();
} else {
assert!(runtime
.publish_note_vector_revision(&token, ¬e, MODEL, &[0.1; 4])
.await
.unwrap());
}
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(note.id)
.await
.unwrap()
.unwrap()
.version,
1
);
let AtomicRunOutcome::Committed { post_commit } = run_atomic_unit(
runtime.sql().as_ref(),
vec![AtomicOpPlan::Update(Box::new(plan))],
)
.await
.unwrap() else {
panic!("commit")
};
assert_eq!(
post_commit.as_slice(),
&[crate::PostCommitEffect::ReindexNote {
note_id: note.id,
version: 2
}]
);
apply_post_commit_effects_with_report(&runtime, &token, post_commit)
.await
.unwrap();
let stored = runtime.sql().reader().await.unwrap().query_scalar(crate::note_write::statement(
"SELECT vec_to_json(embedding) FROM vec_note_version_test WHERE namespace=?1 AND subject_id=?2",
vec![khive_storage::SqlValue::Text("local".into()), khive_storage::SqlValue::Text(note.id.to_string())],
)).await.unwrap().unwrap();
let khive_storage::SqlValue::Text(stored) = stored else {
panic!("vector JSON")
};
assert_eq!(
serde_json::from_str::<Vec<f32>>(&stored).unwrap(),
vec![0.5; 4]
);
}
#[tokio::test]
async fn version_embedding_inheritance_rollback_has_no_effect_token() {
let (runtime, token, _) = fixture();
let note = create(&runtime, &token, "version/inherit-rollback", None).await;
let (_, plan) = runtime
.prepare_versioned_note_update(&token, note.clone(), patch("{\"new\":true}", 1, None))
.await
.unwrap();
assert!(runtime
.publish_note_vector_revision(&token, ¬e, MODEL, &[0.1; 4])
.await
.unwrap());
let outcome = run_atomic_unit(
runtime.sql().as_ref(),
vec![
AtomicOpPlan::Update(Box::new(plan.clone())),
AtomicOpPlan::Update(Box::new(plan)),
],
)
.await
.unwrap();
assert!(matches!(
outcome,
AtomicRunOutcome::RolledBack {
failed_op_index: 1,
..
}
));
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(note.id)
.await
.unwrap()
.unwrap(),
note
);
assert_eq!(vectors(&runtime, &token).await, 1);
}
#[tokio::test]
async fn version_inflight_reindex_rechecks_inside_its_writer_transaction() {
let (runtime, token, service) = fixture();
let note = create(&runtime, &token, "version/inflight", Some(true)).await;
let rt = runtime.clone();
let tok = token.clone();
let id = note.id;
let update = tokio::spawn(async move {
rt.update_note(&tok, id, patch("{\"inflight\":true}", 1, Some(true)))
.await
});
tokio::time::timeout(Duration::from_secs(10), service.started.notified())
.await
.unwrap();
let off = runtime
.update_note(&token, id, patch("{\"off\":true}", 2, Some(false)))
.await
.unwrap();
service.proceed.notify_one();
tokio::time::timeout(Duration::from_secs(10), update)
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(vectors(&runtime, &token).await, 0);
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(id)
.await
.unwrap()
.unwrap(),
off
);
}
#[tokio::test]
async fn version_embedding_purge_rolls_back_with_a_later_conflict() {
let (runtime, token, _) = fixture();
let note = create(&runtime, &token, "version/rollback", Some(true)).await;
let retired = runtime
.backend()
.vectors_for_namespace(
"old_model_chunks",
"old-model-chunks",
4,
token.namespace().as_str(),
)
.unwrap();
retired
.insert(
note.id,
khive_types::SubstrateKind::Note,
token.namespace().as_str(),
"note.content",
vec![vec![0.5; 4]],
)
.await
.unwrap();
assert_eq!(ann_deletes(&runtime, note.id).await, 0);
let (_, off) = runtime
.prepare_versioned_note_update(
&token,
note.clone(),
patch("{\"off\":true}", 1, Some(false)),
)
.await
.unwrap();
let (_, stale) = runtime
.prepare_versioned_note_update(&token, note.clone(), patch("{\"late\":true}", 1, None))
.await
.unwrap();
assert!(matches!(
run_atomic_unit(
runtime.sql().as_ref(),
vec![
AtomicOpPlan::Update(Box::new(off)),
AtomicOpPlan::Update(Box::new(stale))
]
)
.await
.unwrap(),
AtomicRunOutcome::RolledBack {
failed_op_index: 1,
..
}
));
assert_eq!(vectors(&runtime, &token).await, 1);
assert_eq!(retired.count().await.unwrap(), 1);
assert_eq!(
ann_deletes(&runtime, note.id).await,
0,
"delete log must roll back with vectors"
);
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(note.id)
.await
.unwrap()
.unwrap(),
note
);
}
#[tokio::test]
async fn version_guard_observes_a_write_after_prepare_without_timestamp_change() {
let (runtime, token, _) = fixture();
let note = create(&runtime, &token, "version/prepare-race", None).await;
let (_, plan) = runtime
.prepare_versioned_note_update(&token, note.clone(), patch("{\"loser\":true}", 1, None))
.await
.unwrap();
runtime
.sql()
.writer()
.await
.unwrap()
.execute(crate::note_write::statement(
"UPDATE notes SET properties='{}' WHERE id=?1",
vec![khive_storage::SqlValue::Text(note.id.to_string())],
))
.await
.unwrap();
let current = runtime
.notes(&token)
.unwrap()
.get_note(note.id)
.await
.unwrap()
.unwrap();
assert_eq!(current.updated_at, note.updated_at);
assert_eq!(current.version, 2);
let outcome = run_atomic_unit(
runtime.sql().as_ref(),
vec![AtomicOpPlan::Update(Box::new(plan))],
)
.await
.unwrap();
assert!(matches!(
outcome,
AtomicRunOutcome::RolledBack {
failure: crate::atomic_runner::AtomicOpFailure::NoteConflict(
crate::note_write::NoteWriteConflict::Version {
expected: 1,
current: 2
}
),
..
}
));
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(note.id)
.await
.unwrap()
.unwrap(),
current
);
}
#[tokio::test]
async fn version_off_removes_unregistered_vectors_created_after_prepare() {
let (runtime, token, _) = fixture();
let note = create(&runtime, &token, "version/retired-model", Some(true)).await;
let control = create(&runtime, &token, "version/retired-control", None).await;
let (_, plan) = runtime
.prepare_versioned_note_update(
&token,
note.clone(),
patch("{\"off\":true}", 1, Some(false)),
)
.await
.unwrap();
let retired = runtime
.backend()
.vectors_for_namespace(
"retired_note_model_chunks",
"retired-note-model",
4,
token.namespace().as_str(),
)
.unwrap();
for id in [note.id, control.id] {
retired
.insert(
id,
khive_types::SubstrateKind::Note,
token.namespace().as_str(),
"note.content",
vec![vec![0.5; 4]],
)
.await
.unwrap();
}
assert_eq!(retired.count().await.unwrap(), 2);
let foreign = runtime
.backend()
.vectors_for_namespace(
"retired_note_model_chunks",
"retired-note-model",
4,
"foreign",
)
.unwrap();
let foreign_id = uuid::Uuid::new_v4();
foreign
.insert(
foreign_id,
khive_types::SubstrateKind::Note,
"foreign",
"note.content",
vec![vec![0.5; 4]],
)
.await
.unwrap();
assert!(!runtime
.registered_embedding_model_names()
.contains(&"retired-note-model".into()));
let outcome = run_atomic_unit(
runtime.sql().as_ref(),
vec![AtomicOpPlan::Update(Box::new(plan))],
)
.await
.unwrap();
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
assert_eq!(vectors(&runtime, &token).await, 0);
assert_eq!(
retired.count().await.unwrap(),
1,
"explicit off must remove the target from persisted unregistered model tables"
);
let retained = retired
.batch_exists(&[control.id, note.id], token.namespace().as_str())
.await
.unwrap();
assert!(retained.contains(&control.id));
assert!(!retained.contains(¬e.id));
assert_eq!(foreign.count().await.unwrap(), 1);
assert_eq!(ann_deletes(&runtime, foreign_id).await, 0);
assert_eq!(ann_deletes(&runtime, control.id).await, 0);
assert_eq!(
ann_deletes(&runtime, note.id).await,
2,
"one delete delta per purged model row"
);
}
#[tokio::test]
async fn version_legacy_create_cannot_publish_vectors_after_explicit_off() {
assert_legacy_creation_revision_guard(false, false).await;
}
#[tokio::test]
async fn version_legacy_multimodel_create_cannot_publish_after_explicit_off() {
assert_legacy_creation_revision_guard(true, false).await;
}
#[tokio::test]
async fn version_legacy_embedding_failure_preserves_newer_note() {
assert_legacy_creation_revision_guard(false, true).await;
}
#[tokio::test]
async fn version_legacy_multimodel_embedding_failure_preserves_newer_note() {
assert_legacy_creation_revision_guard(true, true).await;
}
async fn assert_legacy_creation_revision_guard(multimodel: bool, fail: bool) {
let (runtime, token, service) = fixture();
let second = Arc::new(Service::default());
if multimodel {
runtime.register_embedder(Provider(second.clone(), "second-note-model".into(), 4));
runtime
.vectors_for_model(&token, "second-note-model")
.unwrap();
}
let rt = runtime.clone();
let tok = token.clone();
let create = tokio::spawn(async move {
rt.create_note(
&tok,
"observation",
None,
"inflight legacy creation",
None,
None,
vec![],
)
.await
});
tokio::time::timeout(Duration::from_secs(10), service.started.notified())
.await
.unwrap();
if multimodel {
tokio::time::timeout(Duration::from_secs(10), second.started.notified())
.await
.unwrap();
}
let notes = runtime
.list_notes(&token, Some("observation"), 2, 0)
.await
.unwrap();
assert_eq!(notes.len(), 1);
let id = notes[0].id;
let off = runtime
.update_note(
&token,
id,
patch("off while creation embeds", 1, Some(false)),
)
.await
.unwrap();
service.fail.store(fail, Ordering::SeqCst);
service.proceed.notify_one();
second.proceed.notify_one();
let result = tokio::time::timeout(Duration::from_secs(10), create)
.await
.unwrap()
.unwrap();
assert_eq!(result.is_err(), fail, "creation outcome: {result:?}");
assert_eq!(
runtime
.notes(&token)
.unwrap()
.get_note(id)
.await
.unwrap()
.unwrap(),
off
);
assert_eq!(
vectors(&runtime, &token).await,
0,
"creation's stale embedding must not reverse the later off-update"
);
if multimodel {
assert_eq!(
runtime
.vectors_for_model(&token, "second-note-model")
.unwrap()
.count()
.await
.unwrap(),
0
);
}
let hits = runtime
.text_for_notes(&token)
.unwrap()
.search(khive_storage::TextSearchRequest {
query: "off while creation".into(),
mode: khive_storage::TextQueryMode::Plain,
filter: None,
top_k: 10,
snippet_chars: 100,
})
.await
.unwrap();
assert!(
hits.iter().any(|hit| hit.subject_id == id),
"newer lexical document must survive"
);
}
#[path = "note_fence_race_tests.rs"]
mod fence_races;
#[path = "fence_live_until_tests.rs"]
mod fence_live_until;
#[path = "fence_identity_tests.rs"]
mod fence_identity;
#[path = "keyed_create_race_tests.rs"]
mod keyed_create_races;