use super::{
ContractOpError, ContractSchemaState, RecallBounds, RecallDiagnostics, RecallMode,
RecallRequest, RememberOutcome, RetrievalPolicy, SchemaAvailability, confirm, conflicts_for,
forget, recall, reject, remember, resolve_prefix, show, use_candidate_once,
};
use saya_store::{ForgetReason, KnowledgeItemStore, SchemaStore, SqliteStateStore};
use saya_types::{
ClaimId, ClaimOrigin, ClaimPayload, ClaimStatus, Column, Database, DatabaseObjectKind,
DatabaseObjectRef, KnowledgeSlot, KnowledgeState, ProfileIdentity, Schema, SchemaTree, Table,
};
use std::{
fs,
path::{Path, PathBuf},
time::{SystemTime, UNIX_EPOCH},
};
use crate::contracts::{ContractConflict, RecallOutcome};
fn temp_root(label: &str) -> PathBuf {
let stamp = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos();
let root = std::env::temp_dir().join(format!(
"saya-contract-ops-{label}-{}-{stamp}",
std::process::id()
));
fs::create_dir_all(&root).unwrap();
root
}
fn profile_a() -> ProfileIdentity {
ProfileIdentity::parse(&format!("p-{}", "a".repeat(64))).unwrap()
}
fn profile_b() -> ProfileIdentity {
ProfileIdentity::parse(&format!("p-{}", "b".repeat(64))).unwrap()
}
const FRESH_NOW: i64 = 1_000_000;
fn avail(tree: SchemaTree) -> SchemaAvailability {
SchemaAvailability::available(tree, 0)
}
fn schemas_for(
profile: &ProfileIdentity,
trees: &[SchemaTree],
) -> Vec<(ProfileIdentity, SchemaAvailability)> {
trees
.iter()
.map(|tree| (profile.clone(), avail(tree.clone())))
.collect()
}
fn object_ref(profile: &ProfileIdentity, name: &str) -> DatabaseObjectRef {
DatabaseObjectRef::new(
profile.clone(),
"catalog",
"public",
name,
DatabaseObjectKind::Table,
)
.unwrap()
}
fn table(cols: &[(&str, &str, bool)]) -> Table {
Table {
name: "orders".into(),
columns: cols
.iter()
.map(|(name, ty, nullable)| Column {
name: (*name).into(),
data_type: (*ty).into(),
nullable: *nullable,
})
.collect(),
}
}
fn schema_tree_for(tables: &[(&str, Table)]) -> SchemaTree {
SchemaTree {
databases: vec![Database {
name: "catalog".into(),
schemas: vec![Schema {
name: "public".into(),
tables: tables
.iter()
.map(|(name, t)| Table {
name: (*name).into(),
columns: t.columns.clone(),
})
.collect(),
}],
}],
}
}
fn schema_tree_for_owned(tables: &[(String, Table)]) -> SchemaTree {
SchemaTree {
databases: vec![Database {
name: "catalog".into(),
schemas: vec![Schema {
name: "public".into(),
tables: tables
.iter()
.map(|(name, t)| Table {
name: name.clone(),
columns: t.columns.clone(),
})
.collect(),
}],
}],
}
}
async fn store_at(db: &Path) -> SqliteStateStore {
let store = SqliteStateStore::new(db);
store
.upsert_schema(profile_a().as_str(), &SchemaTree::default())
.await
.unwrap();
store
}
async fn put_item(
store: &SqliteStateStore,
object: &DatabaseObjectRef,
payload: ClaimPayload,
state: KnowledgeState,
) {
use saya_store::{KnowledgeItemRequest, KnowledgeItemStore};
use saya_types::SchemaBinding;
let slot = slot_for(&payload);
let binding = SchemaBinding::derive(&slot, &payload).expect("slot/payload agree");
let request = KnowledgeItemRequest {
object: object.clone(),
slot,
value: payload,
source: if state == KnowledgeState::Active {
ClaimOrigin::UserExplicit
} else {
ClaimOrigin::AssistantInferred
},
state,
schema_binding_json: serde_json::to_string(&binding).unwrap(),
fingerprint: crate::commands::unobserved_fingerprint(),
};
store.put_knowledge_item(request).await.unwrap();
}
fn slot_for(payload: &ClaimPayload) -> KnowledgeSlot {
match payload {
ClaimPayload::TableDescription { .. } => KnowledgeSlot::TableDescription,
ClaimPayload::TableAlias { .. } => KnowledgeSlot::TableAlias,
ClaimPayload::TableGrain { .. } => KnowledgeSlot::TableGrain,
ClaimPayload::DefaultTimeColumn { .. } => KnowledgeSlot::TableDefaultTime,
ClaimPayload::ColumnDescription { column, .. } => KnowledgeSlot::ColumnDescription {
column: column.clone(),
},
ClaimPayload::ColumnRole { column, .. } => KnowledgeSlot::ColumnRole {
column: column.clone(),
},
_ => panic!("no slot for payload {:?}", payload),
}
}
async fn item_id_for(
store: &SqliteStateStore,
object: &DatabaseObjectRef,
slot: &KnowledgeSlot,
) -> ClaimId {
let item = store
.knowledge_for_object(object)
.await
.expect("knowledge items listed")
.into_iter()
.find(|i| &i.slot == slot)
.unwrap_or_else(|| panic!("no item for slot {slot:?} on {}", object.object()));
ClaimId::parse(&item.id).expect("ki id parses")
}
async fn item_state(store: &SqliteStateStore, id: &ClaimId) -> KnowledgeState {
store
.get_knowledge_item(id.as_str())
.await
.expect("store read")
.expect("item present")
.state
}
async fn seed_active_time_item(
store: &SqliteStateStore,
object: &DatabaseObjectRef,
column: &str,
) -> ClaimId {
put_item(
store,
object,
ClaimPayload::default_time_column(column, None).unwrap(),
KnowledgeState::Active,
)
.await;
item_id_for(store, object, &KnowledgeSlot::TableDefaultTime).await
}
fn recall_request<'a>(
profiles: &'a [ProfileIdentity],
schemas: &'a [(ProfileIdentity, SchemaAvailability)],
terms: &'a [String],
allow_database_context: bool,
bounds: RecallBounds,
) -> RecallRequest<'a> {
RecallRequest {
profiles,
explicit_refs: &[],
terms,
allow_database_context,
schemas,
now_unix_ms: FRESH_NOW,
bounds,
recall_mode: RecallMode::Confirmed,
admit_candidate: None,
policy: RetrievalPolicy::ForModel,
}
}
#[tokio::test]
async fn cross_profile_alias_resolves_only_in_its_own_profile() {
let root = temp_root("cross_profile");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let a = profile_a();
let b = profile_b();
let obj_a = object_ref(&a, "orders");
let obj_b = object_ref(&b, "orders");
put_item(
&store,
&obj_a,
ClaimPayload::table_alias("orders").unwrap(),
KnowledgeState::Active,
)
.await;
put_item(
&store,
&obj_b,
ClaimPayload::table_alias("orders").unwrap(),
KnowledgeState::Active,
)
.await;
let schema_a = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&a),
&[(a.clone(), avail(schema_a))],
&["orders".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(outcome.contracts.len(), 1);
assert_eq!(outcome.contracts[0].object, obj_a);
let schema_b = schema_tree_for(&[("orders", table(&[("id", "int", false)]))]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&b),
&[(b.clone(), avail(schema_b))],
&["orders".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(outcome.contracts.len(), 1);
assert_eq!(outcome.contracts[0].object, obj_b);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn candidate_claims_never_appear_in_recall() {
let root = temp_root("no_candidates");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::table_alias("orders").unwrap(),
KnowledgeState::Active,
)
.await;
put_item(
&store,
&obj,
ClaimPayload::table_alias("secret_alias").unwrap(),
KnowledgeState::Pending,
)
.await;
let schema = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&["orders".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(outcome.contracts.len(), 1);
let claims = &outcome.contracts[0].claims;
assert!(claims.iter().all(|c| c.status.is_recallable()));
assert!(
!claims.iter().any(|c| {
matches!(
&c.value,
ClaimPayload::TableAlias { alias, .. } if alias == "secret_alias"
)
}),
"candidate alias leaked into recall"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn bounds_hold_on_objects_and_claims() {
let root = temp_root("bounds");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let objs: Vec<DatabaseObjectRef> = (0..3)
.map(|i| object_ref(&p, &format!("orders{i}")))
.collect();
for obj in &objs {
put_item(
&store,
obj,
ClaimPayload::table_alias(obj.object()).unwrap(),
KnowledgeState::Active,
)
.await;
}
let heavy = object_ref(&p, "heavy");
for i in 0..4 {
put_item(
&store,
&heavy,
ClaimPayload::table_alias(format!("a{i}")).unwrap(),
KnowledgeState::Active,
)
.await;
}
let terms: Vec<String> = objs.iter().map(|o| o.object().to_string()).collect();
let bounds = RecallBounds {
max_objects: 2,
max_claims_per_object: 12,
max_bytes: 16384,
};
let schema = schema_tree_for(&[
("orders0", table(&[("id", "bigint", false)])),
("orders1", table(&[("id", "bigint", false)])),
("orders2", table(&[("id", "bigint", false)])),
]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&terms,
true,
bounds,
),
)
.await;
assert_eq!(outcome.contracts.len(), 2, "max_objects not honored");
assert!(
outcome.contracts.iter().any(|c| c.truncated),
"object truncation not flagged"
);
let bounds2 = RecallBounds {
max_objects: 5,
max_claims_per_object: 3,
max_bytes: 16384,
};
let schema_h = schema_tree_for(&[("heavy", table(&[("id", "bigint", false)]))]);
let outcome2 = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema_h.clone()))],
&["heavy".to_string()],
true,
bounds2,
),
)
.await;
assert_eq!(outcome2.contracts.len(), 1);
assert_eq!(outcome2.contracts[0].claims.len(), 3);
assert!(outcome2.contracts[0].truncated);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn ambiguous_alias_returns_every_match() {
let root = temp_root("ambiguous");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj1 = object_ref(&p, "table_one");
let obj2 = object_ref(&p, "table_two");
put_item(
&store,
&obj1,
ClaimPayload::table_alias("shared").unwrap(),
KnowledgeState::Active,
)
.await;
put_item(
&store,
&obj2,
ClaimPayload::table_alias("shared").unwrap(),
KnowledgeState::Active,
)
.await;
let schema = schema_tree_for(&[
("table_one", table(&[("id", "bigint", false)])),
("table_two", table(&[("id", "bigint", false)])),
]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&["shared".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(
outcome.contracts.len(),
2,
"ambiguous alias should return both"
);
let names: Vec<&str> = outcome
.contracts
.iter()
.map(|c| c.object.object())
.collect();
assert!(names.contains(&"table_one"));
assert!(names.contains(&"table_two"));
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn privacy_gate_returns_zero_and_counts_excluded() {
let root = temp_root("privacy");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::table_alias("orders").unwrap(),
KnowledgeState::Active,
)
.await;
let schema = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&["orders".to_string()],
false,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(
outcome.contracts.len(),
0,
"privacy gate returned contracts"
);
assert!(
outcome.diagnostics.excluded_by_privacy > 0,
"excluded_by_privacy should be non-zero when matches were suppressed"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn unopenable_store_reports_unavailable_without_erroring() {
let root = temp_root("unavailable");
fs::write(root.join("blocker"), b"x").unwrap();
let bad = root.join("blocker/state.sqlite3");
let store = SqliteStateStore::new(&bad);
let p = profile_a();
let schema = SchemaTree::default();
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&["orders".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert!(outcome.diagnostics.store_unavailable);
assert_eq!(outcome.contracts.len(), 0);
assert_eq!(outcome.diagnostics.considered, 0);
assert_eq!(outcome.diagnostics.selected, 0);
let _ = fs::remove_dir_all(root);
}
#[test]
fn two_table_grain_claims_conflict_but_both_returned() {
use crate::contracts::ContractClaim;
let p = profile_a();
let obj = object_ref(&p, "orders");
let id1 = ClaimId::parse("c-grain0001").unwrap();
let id2 = ClaimId::parse("c-grain0002").unwrap();
let claims = vec![
ContractClaim {
id: id1.clone(),
object: obj.clone(),
value: ClaimPayload::table_grain("one row per order", None).unwrap(),
source: ClaimOrigin::UserExplicit,
status: ClaimStatus::Confirmed,
},
ContractClaim {
id: id2.clone(),
object: obj.clone(),
value: ClaimPayload::table_grain("one row per order line", None).unwrap(),
source: ClaimOrigin::UserExplicit,
status: ClaimStatus::Confirmed,
},
];
let conflicts = conflicts_for(&claims);
let grain_conflicts: Vec<&ContractConflict> = conflicts
.iter()
.filter(|c| c.kind == "table_grain")
.collect();
assert_eq!(
grain_conflicts.len(),
1,
"expected one table_grain conflict"
);
let ids: Vec<ClaimId> = grain_conflicts[0].claim_ids.clone();
assert!(
ids.contains(&id1) && ids.contains(&id2),
"conflict must name both IDs"
);
assert_eq!(claims.len(), 2);
}
#[tokio::test]
async fn two_table_description_claims_do_not_conflict() {
let root = temp_root("desc_no_conflict");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::table_description("sales fact table").unwrap(),
KnowledgeState::Active,
)
.await;
put_item(
&store,
&obj,
ClaimPayload::table_description("updated nightly").unwrap(),
KnowledgeState::Active,
)
.await;
let schema = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&["orders".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
let contract = &outcome.contracts[0];
assert!(
contract
.conflicts
.iter()
.all(|c| c.kind != "table_description"),
"TableDescription must not be flagged as a conflict (SPEC REVIEW)"
);
assert_eq!(contract.claims.len(), 2);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn recall_diagnostics_carry_no_claim_text() {
let root = temp_root("no_text_in_diag");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
const SENTINEL: &str = "SENTINELDIAGTEXT";
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::table_description(format!("orders {SENTINEL} facts")).unwrap(),
KnowledgeState::Active,
)
.await;
let schema = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
let outcome: RecallOutcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&["orders".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
let debug = format!("{:?}", outcome.diagnostics);
assert!(
!debug.contains(SENTINEL),
"diagnostics leaked claim text: {debug}"
);
let diags: RecallDiagnostics = outcome.diagnostics.clone();
assert!(!format!("{diags:?}").contains(SENTINEL));
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn forgotten_claim_disappears_from_recall() {
use saya_store::KnowledgeItemStore;
let root = temp_root("forget_disappears");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::table_alias("keep").unwrap(),
KnowledgeState::Active,
)
.await;
put_item(
&store,
&obj,
ClaimPayload::table_alias("gone").unwrap(),
KnowledgeState::Active,
)
.await;
let items = store
.knowledge_for_object(&obj)
.await
.expect("knowledge items listed");
let find_id = |alias: &str| {
items
.iter()
.find(|i| {
matches!(
&i.value,
ClaimPayload::TableAlias { alias: a, .. } if a == alias
)
})
.map(|i| i.id.clone())
.expect("alias item stored")
};
let keep = ClaimId::parse(&find_id("keep")).unwrap();
let gone_id = find_id("gone");
store
.update_knowledge_item_state(&gone_id, KnowledgeState::Dismissed)
.await
.expect("item dismissed");
let schema = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&["orders".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
let contract = &outcome.contracts[0];
let ids: Vec<ClaimId> = contract.claims.iter().map(|c| c.id.clone()).collect();
assert!(ids.contains(&keep), "kept claim should remain");
assert!(
!ids.iter().any(|id| id.as_str() == gone_id),
"forgotten claim should not appear"
);
let _ = fs::remove_dir_all(root);
}
use super::review_queue;
#[tokio::test]
async fn queue_lists_candidates_not_confirmed_or_forgotten() {
let root = temp_root("queue_membership");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let live = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
put_item(
&store,
&obj,
ClaimPayload::table_alias("confirmed_alias").unwrap(),
KnowledgeState::Active,
)
.await;
put_item(
&store,
&obj,
ClaimPayload::table_alias("cand_alias").unwrap(),
KnowledgeState::Pending,
)
.await;
let cand_id = item_id_for(&store, &obj, &KnowledgeSlot::TableAlias).await;
put_item(
&store,
&obj,
ClaimPayload::table_description("forgotten").unwrap(),
KnowledgeState::Dismissed,
)
.await;
let queued = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live))],
200,
)
.await
.unwrap();
let ids: Vec<ClaimId> = queued.iter().map(|q| q.claim.id.clone()).collect();
assert!(ids.contains(&cand_id), "candidate missing from queue");
for id in &ids {
assert_eq!(
item_state(&store, id).await,
KnowledgeState::Pending,
"non-pending item appeared in the queue"
);
}
assert_eq!(queued.len(), 1);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn queue_orders_oldest_then_slot_then_id() {
let root = temp_root("queue_order");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let live = schema_tree_for(&[
("old", table(&[("id", "bigint", false)])),
("new", table(&[("id", "bigint", false)])),
("both", table(&[("id", "bigint", false)])),
]);
put_item(
&store,
&object_ref(&p, "old"),
ClaimPayload::table_alias("old").unwrap(),
KnowledgeState::Pending,
)
.await;
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
put_item(
&store,
&object_ref(&p, "new"),
ClaimPayload::table_alias("new").unwrap(),
KnowledgeState::Pending,
)
.await;
let queued = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live.clone()))],
200,
)
.await
.unwrap();
let ordered: Vec<&str> = queued.iter().map(|q| q.claim.object.object()).collect();
let old_pos = ordered
.iter()
.position(|o| *o == "old")
.expect("old candidate missing");
let new_pos = ordered
.iter()
.position(|o| *o == "new")
.expect("new candidate missing");
assert!(old_pos < new_pos, "oldest-first order broken: {ordered:?}");
let both = object_ref(&p, "both");
put_item(
&store,
&both,
ClaimPayload::table_alias("both_alias").unwrap(),
KnowledgeState::Pending,
)
.await;
put_item(
&store,
&both,
ClaimPayload::table_description("both description").unwrap(),
KnowledgeState::Pending,
)
.await;
let queued = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live.clone()))],
200,
)
.await
.unwrap();
let both_rows: Vec<&str> = queued
.iter()
.filter(|q| q.claim.object.object() == "both")
.map(|q| q.claim.id.as_str())
.collect();
let alias_id = item_id_for(&store, &both, &KnowledgeSlot::TableAlias)
.await
.as_str()
.to_string();
let desc_id = item_id_for(&store, &both, &KnowledgeSlot::TableDescription)
.await
.as_str()
.to_string();
let alias_pos = both_rows
.iter()
.position(|id| *id == alias_id)
.expect("alias in queue");
let desc_pos = both_rows
.iter()
.position(|id| *id == desc_id)
.expect("description in queue");
assert!(
alias_pos < desc_pos,
"slot tie-break broken (alias should precede description): {both_rows:?}"
);
let run_a = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live.clone()))],
200,
)
.await
.unwrap();
let run_b = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live))],
200,
)
.await
.unwrap();
let ids_a: Vec<ClaimId> = run_a.iter().map(|q| q.claim.id.clone()).collect();
let ids_b: Vec<ClaimId> = run_b.iter().map(|q| q.claim.id.clone()).collect();
assert_eq!(ids_a, ids_b, "queue order was not stable across runs");
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn queue_limit_is_respected_and_clamped_at_200() {
let root = temp_root("queue_limit");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let live = schema_tree_for(&[]);
for i in 0..5 {
put_item(
&store,
&object_ref(&p, &format!("t{i}")),
ClaimPayload::table_alias(format!("t{i}")).unwrap(),
KnowledgeState::Pending,
)
.await;
}
let queued = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live.clone()))],
3,
)
.await
.unwrap();
assert_eq!(queued.len(), 3, "limit not respected");
let queued = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live.clone()))],
0,
)
.await
.unwrap();
assert!(queued.is_empty(), "limit 0 returned candidates");
let queued = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live))],
10_000,
)
.await
.unwrap();
assert_eq!(queued.len(), 5, "over-large limit misbehaved");
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn queue_reports_schema_state_for_a_changed_object() {
let root = temp_root("queue_schema_state");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::column_description("amount", "how much").unwrap(),
KnowledgeState::Pending,
)
.await;
let live_dropped = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
let queued = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live_dropped))],
200,
)
.await
.unwrap();
assert_eq!(queued.len(), 1);
assert_eq!(
queued[0].schema_state,
ContractSchemaState::Stale,
"a candidate whose referenced column is gone must read stale"
);
let live_current = schema_tree_for(&[(
"orders",
table(&[("id", "bigint", false), ("amount", "numeric", false)]),
)]);
let queued = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live_current))],
200,
)
.await
.unwrap();
assert_eq!(queued[0].schema_state, ContractSchemaState::Current);
let queued = review_queue(&store, std::slice::from_ref(&p), &[], 200)
.await
.unwrap();
assert_eq!(
queued[0].schema_state,
ContractSchemaState::LiveSchemaUnavailable
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn queue_evidence_count_is_zero_without_an_evidence_table() {
let root = temp_root("queue_evidence_count");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let live = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
put_item(
&store,
&obj,
ClaimPayload::table_alias("orders").unwrap(),
KnowledgeState::Pending,
)
.await;
let queued = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(live))],
200,
)
.await
.unwrap();
assert_eq!(queued.len(), 1);
assert_eq!(
queued[0].evidence_count, 0,
"no evidence is attached to a knowledge item"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn queue_unopenable_store_errors_unavailable() {
let root = temp_root("queue_unavailable");
fs::write(root.join("blocker"), b"x").unwrap();
let bad = root.join("blocker/state.sqlite3");
let store = SqliteStateStore::new(&bad);
let p = profile_a();
let err = review_queue(
&store,
std::slice::from_ref(&p),
&[(p.clone(), avail(SchemaTree::default()))],
200,
)
.await
.unwrap_err();
assert_eq!(err, ContractOpError::Unavailable);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn resolve_prefix_matches_a_unique_ki_id_prefix() {
let root = temp_root("resolve_prefix_unique");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::table_alias("orders").unwrap(),
KnowledgeState::Pending,
)
.await;
let id = item_id_for(&store, &obj, &KnowledgeSlot::TableAlias).await;
assert_eq!(resolve_prefix(&store, &p, id.as_str()).await.unwrap(), id);
let short: String = id.as_str().chars().take(6).collect();
assert_eq!(
resolve_prefix(&store, &p, &short).await.unwrap(),
id,
"a short unique `ki-` prefix resolves to the item"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn resolve_prefix_refuses_ambiguous_and_missing() {
let root = temp_root("resolve_prefix_refuse");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
put_item(
&store,
&object_ref(&p, "orders"),
ClaimPayload::table_alias("orders").unwrap(),
KnowledgeState::Pending,
)
.await;
put_item(
&store,
&object_ref(&p, "returns"),
ClaimPayload::table_alias("returns").unwrap(),
KnowledgeState::Pending,
)
.await;
let err = resolve_prefix(&store, &p, "ki-").await.unwrap_err();
assert_eq!(
err,
ContractOpError::Conflict,
"a prefix matching more than one item is ambiguous, not a guess"
);
let err = resolve_prefix(&store, &p, "ki-deadbeef").await.unwrap_err();
assert_eq!(err, ContractOpError::NotFound);
let err = resolve_prefix(&store, &p, "c-").await.unwrap_err();
assert_eq!(
err,
ContractOpError::NotFound,
"a legacy `c-` prefix matches no knowledge item"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn review_wrappers_pass_through_and_map_errors() {
let root = temp_root("review_wrappers");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::table_alias("orders").unwrap(),
KnowledgeState::Pending,
)
.await;
let id = item_id_for(&store, &obj, &KnowledgeSlot::TableAlias).await;
let mut chars: Vec<char> = id.as_str().chars().collect();
if let Some(last) = chars.last_mut() {
*last = match *last {
'0'..='8' => char::from_u32(*last as u32 + 1).unwrap(),
'9' | 'a'..='e' => char::from_u32(*last as u32 + 1).unwrap(),
'f' => '0',
_ => '0',
};
}
let fake = ClaimId::parse(&chars.into_iter().collect::<String>()).unwrap();
let err = confirm(&store, &fake).await.unwrap_err();
assert_eq!(err, ContractOpError::NotFound);
let rejected = reject(&store, &id).await.unwrap();
assert_eq!(rejected.status, ClaimStatus::Rejected);
assert_eq!(item_state(&store, &id).await, KnowledgeState::Dismissed);
put_item(
&store,
&obj,
ClaimPayload::table_description("the orders table").unwrap(),
KnowledgeState::Pending,
)
.await;
let second = item_id_for(&store, &obj, &KnowledgeSlot::TableDescription).await;
forget(&store, &second, ForgetReason::UserRequest)
.await
.unwrap();
assert_eq!(item_state(&store, &second).await, KnowledgeState::Dismissed);
let shown = show(
&store,
&obj,
&SchemaAvailability::Missing,
RetrievalPolicy::ForHumanReview,
FRESH_NOW,
)
.await
.unwrap();
assert!(shown.is_none(), "a dismissed-only contract shows None");
let _ = fs::remove_dir_all(root);
}
async fn seed_active_time_item_drifted(
store: &SqliteStateStore,
object: &DatabaseObjectRef,
) -> ClaimId {
let id = seed_active_time_item(store, object, "created_at").await;
let drifted = schema_tree_for(&[(object.object(), table(&[("id", "bigint", false)]))]);
store
.upsert_schema(object.profile().as_str(), &drifted)
.await
.unwrap();
id
}
#[tokio::test]
async fn show_keeps_a_stale_contract_with_state_and_claims_for_a_human() {
let root = temp_root("show_stale_human");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let id = seed_active_time_item_drifted(&store, &obj).await;
let shown = show(
&store,
&obj,
&avail(schema_tree_for(&[(
"orders",
table(&[("id", "bigint", false)]),
)])),
RetrievalPolicy::ForHumanReview,
FRESH_NOW,
)
.await
.unwrap()
.expect("a stale contract is shown to a human, not hidden");
assert_eq!(
shown.schema_state,
ContractSchemaState::Stale,
"the state is reported as stale"
);
let ids: Vec<ClaimId> = shown.claims.iter().map(|c| c.id.clone()).collect();
assert!(
ids.contains(&id),
"the stale item is kept for a human reviewer, got {ids:?}"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn show_for_model_drops_a_stale_contracts_claims_but_names_the_object() {
let root = temp_root("show_stale_model");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
seed_active_time_item_drifted(&store, &obj).await;
let shown = show(
&store,
&obj,
&avail(schema_tree_for(&[(
"orders",
table(&[("id", "bigint", false)]),
)])),
RetrievalPolicy::ForModel,
FRESH_NOW,
)
.await
.unwrap()
.expect("the stale object is reported, not hidden");
assert_eq!(shown.schema_state, ContractSchemaState::Stale);
assert!(
shown.claims.is_empty(),
"no claims to act on for a stale object, got {}",
shown.claims.len()
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn recall_for_model_counts_a_stale_exclusion() {
let root = temp_root("recall_stale_count");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::default_time_column("created_at", None).unwrap(),
KnowledgeState::Active,
)
.await;
let drifted = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(drifted))],
&["orders".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(
outcome.contracts.len(),
0,
"the stale contract is dropped for the model"
);
assert!(
outcome.diagnostics.excluded_by_schema >= 1,
"the stale exclusion is counted: {:?}",
outcome.diagnostics
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn recall_over_many_objects_returns_every_match() {
let root = temp_root("recall_many_objects");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let names: Vec<String> = (0..20).map(|i| format!("orders{i}")).collect();
for name in &names {
let obj = object_ref(&p, name);
put_item(
&store,
&obj,
ClaimPayload::table_alias(name.as_str()).unwrap(),
KnowledgeState::Active,
)
.await;
}
let tables: Vec<(&str, Table)> = names
.iter()
.map(|n| (n.as_str(), table(&[("id", "bigint", false)])))
.collect();
let schema = schema_tree_for(&tables);
let terms: Vec<String> = vec!["orders".to_string()];
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&terms,
true,
RecallBounds {
max_objects: 20,
max_claims_per_object: 12,
max_bytes: 16384,
},
),
)
.await;
assert_eq!(
outcome.contracts.len(),
20,
"bulk recall must return every matching object, got {}",
outcome.contracts.len()
);
let selected: Vec<String> = outcome
.contracts
.iter()
.map(|c| c.object.object().to_string())
.collect();
for name in &names {
assert!(selected.contains(name), "missing {name} in {selected:?}");
}
let _ = fs::remove_dir_all(root);
}
async fn seed_current_time_column(
store: &SqliteStateStore,
obj: &DatabaseObjectRef,
_table: &Table,
column: &str,
) {
put_item(
store,
obj,
ClaimPayload::default_time_column(column, None).unwrap(),
KnowledgeState::Active,
)
.await;
}
#[tokio::test]
async fn model_path_treats_a_stale_by_age_cache_as_live_schema_unavailable() {
let root = temp_root("freshness_model_stale_by_age");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let table = table(&[("id", "bigint", false), ("created_at", "timestamp", false)]);
seed_current_time_column(&store, &obj, &table, "created_at").await;
let stale_by_age = SchemaAvailability::available(schema_tree_for(&[("orders", table)]), 0);
let now = 25 * 60 * 60 * 1000;
let outcome = recall(
&store,
recall_request_with_freshness(
std::slice::from_ref(&p),
&[(p.clone(), stale_by_age)],
&["orders".to_string()],
true,
RecallBounds::defaults(),
now,
RetrievalPolicy::ForModel,
),
)
.await;
assert_eq!(
outcome.contracts.len(),
1,
"a stale-by-age contract is kept for the model, labelled — not dropped"
);
assert_eq!(
outcome.contracts[0].schema_state,
ContractSchemaState::LiveSchemaUnavailable,
"a stale-by-age cache cannot vouch for currency on the model path"
);
assert_eq!(
outcome.diagnostics.excluded_by_schema, 0,
"a stale-by-age contract is not excluded as stale: {:?}",
outcome.diagnostics
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn human_path_uses_a_stale_by_age_cache_to_classify_current() {
let root = temp_root("freshness_human_stale_by_age");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let table = table(&[("id", "bigint", false), ("created_at", "timestamp", false)]);
seed_current_time_column(&store, &obj, &table, "created_at").await;
let stale_by_age = SchemaAvailability::available(schema_tree_for(&[("orders", table)]), 0);
let now = 25 * 60 * 60 * 1000;
let outcome = recall(
&store,
recall_request_with_freshness(
std::slice::from_ref(&p),
&[(p.clone(), stale_by_age)],
&["orders".to_string()],
true,
RecallBounds::defaults(),
now,
RetrievalPolicy::ForHumanReview,
),
)
.await;
assert_eq!(
outcome.contracts.len(),
1,
"the human path shows the contract despite the stale-by-age cache"
);
assert_eq!(
outcome.contracts[0].schema_state,
ContractSchemaState::Current,
"the human path classifies against the stale-by-age cache as current"
);
let _ = fs::remove_dir_all(root);
}
fn recall_request_with_freshness<'a>(
profiles: &'a [ProfileIdentity],
schemas: &'a [(ProfileIdentity, SchemaAvailability)],
terms: &'a [String],
allow_database_context: bool,
bounds: RecallBounds,
now_unix_ms: i64,
policy: RetrievalPolicy,
) -> RecallRequest<'a> {
RecallRequest {
profiles,
explicit_refs: &[],
terms,
allow_database_context,
schemas,
now_unix_ms,
bounds,
recall_mode: RecallMode::Confirmed,
admit_candidate: None,
policy,
}
}
#[tokio::test]
async fn confirm_revalidates_a_stale_claim_against_the_cached_schema() {
let root = temp_root("confirm_revalidates");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let id = seed_active_time_item(&store, &obj, "created_at").await;
let base = schema_tree_for(&[(
"orders",
table(&[("id", "bigint", false), ("created_at", "timestamp", false)]),
)]);
store.upsert_schema(p.as_str(), &base).await.unwrap();
let confirmed = confirm(&store, &id).await.unwrap();
assert_eq!(confirmed.status, ClaimStatus::Confirmed);
assert_eq!(item_state(&store, &id).await, KnowledgeState::Active);
let shown = show(
&store,
&obj,
&avail(base),
RetrievalPolicy::ForHumanReview,
FRESH_NOW,
)
.await
.unwrap()
.expect("a confirmed-against-schema contract is shown");
assert_eq!(
shown.schema_state,
ContractSchemaState::Current,
"a revalidated item reads Current, not Stale"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn confirm_refuses_a_stale_claim_with_no_cached_schema() {
let root = temp_root("confirm_no_schema");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let id = seed_active_time_item(&store, &obj, "created_at").await;
store.invalidate_schema(p.as_str()).await.unwrap();
let err = confirm(&store, &id).await.unwrap_err();
assert_eq!(
err,
ContractOpError::SchemaUnavailable,
"an existing fact cannot be reconfirmed without a schema"
);
assert_eq!(item_state(&store, &id).await, KnowledgeState::Active);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn confirm_refuses_a_stale_claim_whose_referenced_column_is_gone() {
let root = temp_root("confirm_column_gone");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let id = seed_active_time_item(&store, &obj, "created_at").await;
let drifted = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
store.upsert_schema(p.as_str(), &drifted).await.unwrap();
let err = confirm(&store, &id).await.unwrap_err();
assert_eq!(
err,
ContractOpError::ColumnGone,
"do not revive a fact whose referenced column is gone"
);
assert_eq!(item_state(&store, &id).await, KnowledgeState::Active);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn confirm_a_candidate_is_status_only_and_needs_no_schema() {
let root = temp_root("confirm_candidate");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
store.invalidate_schema(p.as_str()).await.unwrap();
put_item(
&store,
&obj,
ClaimPayload::table_description("the orders table").unwrap(),
KnowledgeState::Pending,
)
.await;
let id = item_id_for(&store, &obj, &KnowledgeSlot::TableDescription).await;
let confirmed = confirm(&store, &id).await.unwrap();
assert_eq!(confirmed.status, ClaimStatus::Confirmed);
assert_eq!(item_state(&store, &id).await, KnowledgeState::Active);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn confirm_candidate_against_an_empty_cached_schema_succeeds_not_object_gone() {
let root = temp_root("confirm_empty_schema_candidate");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::table_description("the orders table").unwrap(),
KnowledgeState::Pending,
)
.await;
let id = item_id_for(&store, &obj, &KnowledgeSlot::TableDescription).await;
let confirmed = confirm(&store, &id).await.unwrap();
assert_eq!(confirmed.status, ClaimStatus::Confirmed);
assert_eq!(item_state(&store, &id).await, KnowledgeState::Active);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn confirm_active_against_an_empty_cached_schema_is_schema_unavailable_not_object_gone() {
let root = temp_root("confirm_empty_schema_active");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let id = seed_active_time_item(&store, &obj, "created_at").await;
store
.upsert_schema(p.as_str(), &SchemaTree::default())
.await
.unwrap();
let err = confirm(&store, &id).await.unwrap_err();
assert_eq!(
err,
ContractOpError::SchemaUnavailable,
"an empty cache is no schema to verify against, not proof the table is gone"
);
assert_eq!(item_state(&store, &id).await, KnowledgeState::Active);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn confirm_against_a_populated_schema_lacking_the_table_is_object_gone() {
let root = temp_root("confirm_populated_no_table");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let id = seed_active_time_item(&store, &obj, "created_at").await;
let populated = schema_tree_for(&[("shipments", table(&[("id", "bigint", false)]))]);
store.upsert_schema(p.as_str(), &populated).await.unwrap();
let err = confirm(&store, &id).await.unwrap_err();
assert_eq!(
err,
ContractOpError::ObjectGone,
"a populated schema that lacks the table is evidence the table is gone"
);
assert_eq!(item_state(&store, &id).await, KnowledgeState::Active);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn confirming_a_stale_claim_names_the_missing_column_and_the_repair() {
let root = temp_root("confirm-colgone");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
let id = seed_active_time_item(&store, &obj, "created_at").await;
let drifted = schema_tree_for(&[("orders", table(&[("id", "bigint", false)]))]);
store.upsert_schema(p.as_str(), &drifted).await.unwrap();
let error = confirm(&store, &id).await.unwrap_err();
assert_eq!(error, ContractOpError::ColumnGone);
let rendered = error.to_string();
assert!(
rendered.contains("column"),
"names the obstacle: {rendered}"
);
assert!(
rendered.contains("edit") && rendered.contains("forget"),
"names both repairs: {rendered}"
);
assert!(
!rendered.contains("conflict"),
"must not claim a conflict with another claim: {rendered}"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn plural_prompt_term_selects_the_singular_table() {
let root = temp_root("plural_selects_singular");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "rental");
put_item(
&store,
&obj,
ClaimPayload::table_description("one row per rental").unwrap(),
KnowledgeState::Active,
)
.await;
let schema = schema_tree_for(&[("rental", table(&[("id", "bigint", false)]))]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&["rentals".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(
outcome.contracts.len(),
1,
"plural term `rentals` must select the singular `rental` table, got {}",
outcome.contracts.len()
);
assert_eq!(outcome.contracts[0].object, obj);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn singular_term_still_selects() {
let root = temp_root("singular_still_selects");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "rental");
put_item(
&store,
&obj,
ClaimPayload::table_description("one row per rental").unwrap(),
KnowledgeState::Active,
)
.await;
let schema = schema_tree_for(&[("rental", table(&[("id", "bigint", false)]))]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&["rental".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(outcome.contracts.len(), 1);
assert_eq!(outcome.contracts[0].object, obj);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn irregular_s_words_are_not_mangled() {
let root = temp_root("irregular_s_words");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let status = object_ref(&p, "status");
let address = object_ref(&p, "address");
let staff = object_ref(&p, "staff");
for obj in [&status, &address, &staff] {
put_item(
&store,
obj,
ClaimPayload::table_description("a table").unwrap(),
KnowledgeState::Active,
)
.await;
}
let schema = schema_tree_for(&[
("status", table(&[("id", "bigint", false)])),
("address", table(&[("id", "bigint", false)])),
("staff", table(&[("id", "bigint", false)])),
]);
for (term, want) in [
("status", &status),
("address", &address),
("staff", &staff),
] {
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema.clone()))],
&[term.to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(
outcome.contracts.len(),
1,
"term `{term}` should select exactly one object, got {}",
outcome.contracts.len()
);
assert_eq!(
outcome.contracts[0].object, *want,
"term `{term}` selected the wrong object"
);
}
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn term_naming_no_part_of_an_object_does_not_select_it() {
let root = temp_root("no_overmatch");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let rental = object_ref(&p, "rental");
let customer = object_ref(&p, "customer");
for obj in [&rental, &customer] {
put_item(
&store,
obj,
ClaimPayload::table_description("a table").unwrap(),
KnowledgeState::Active,
)
.await;
}
let schema = schema_tree_for(&[
("rental", table(&[("id", "bigint", false)])),
("customer", table(&[("id", "bigint", false)])),
]);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema.clone()))],
&["public".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(
outcome.contracts.len(),
0,
"schema-name term `public` must not select objects in the public schema, got {}",
outcome.contracts.len()
);
let outcome = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&["rentals".to_string()],
true,
RecallBounds::defaults(),
),
)
.await;
assert_eq!(outcome.contracts.len(), 1);
assert_eq!(outcome.contracts[0].object, rental);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn same_prompt_selects_the_same_objects_in_the_same_order() {
let root = temp_root("selection_determinism");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let names = ["orders_alpha", "orders_beta", "orders_gamma"];
for name in &names {
let obj = object_ref(&p, name);
put_item(
&store,
&obj,
ClaimPayload::table_description("a table").unwrap(),
KnowledgeState::Active,
)
.await;
}
let tables: Vec<(&str, Table)> = names
.iter()
.map(|n| (*n, table(&[("id", "bigint", false)])))
.collect();
let schema = schema_tree_for(&tables);
let terms: Vec<String> = vec!["orders".to_string()];
let bounds = RecallBounds {
max_objects: 10,
max_claims_per_object: 12,
max_bytes: 16384,
};
let a = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema.clone()))],
&terms,
true,
bounds,
),
)
.await;
let b = recall(
&store,
recall_request(
std::slice::from_ref(&p),
&[(p.clone(), avail(schema))],
&terms,
true,
bounds,
),
)
.await;
let names_a: Vec<String> = a
.contracts
.iter()
.map(|c| c.object.object().to_string())
.collect();
let names_b: Vec<String> = b
.contracts
.iter()
.map(|c| c.object.object().to_string())
.collect();
assert_eq!(
names_a, names_b,
"same prompt must select the same objects in the same order"
);
assert_eq!(names_a.len(), 3);
let _ = fs::remove_dir_all(root);
}
fn recall_admitting_one<'a>(
profiles: &'a [ProfileIdentity],
schemas: &'a [(ProfileIdentity, SchemaAvailability)],
terms: &'a [String],
admit: Option<ClaimId>,
) -> RecallRequest<'a> {
RecallRequest {
profiles,
explicit_refs: &[],
terms,
allow_database_context: true,
schemas,
now_unix_ms: FRESH_NOW,
bounds: RecallBounds::defaults(),
recall_mode: RecallMode::Confirmed,
policy: RetrievalPolicy::ForModel,
admit_candidate: admit,
}
}
async fn seed_one_pending_item(
store: &SqliteStateStore,
) -> (ProfileIdentity, DatabaseObjectRef, ClaimId) {
use saya_store::KnowledgeItemStore;
let p = profile_a();
let obj = object_ref(&p, "orders");
let tree = schema_tree_for(&[(
"orders",
table(&[("id", "bigint", false), ("amount", "numeric", false)]),
)]);
store.upsert_schema(p.as_str(), &tree).await.unwrap();
put_item(
store,
&obj,
ClaimPayload::column_description("amount", "how much").unwrap(),
KnowledgeState::Pending,
)
.await;
let id = ClaimId::parse(
&store
.knowledge_for_object(&obj)
.await
.expect("knowledge items listed")
.into_iter()
.find(|i| {
i.slot
== KnowledgeSlot::ColumnDescription {
column: "amount".into(),
}
})
.expect("pending item stored")
.id,
)
.expect("ki id");
(p, obj, id)
}
#[tokio::test]
async fn use_candidate_once_leaves_status_and_origin_unchanged() {
let root = temp_root("use_once_unchanged");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::column_description("amount", "how much").unwrap(),
KnowledgeState::Pending,
)
.await;
let id = item_id_for(
&store,
&obj,
&KnowledgeSlot::ColumnDescription {
column: "amount".into(),
},
)
.await;
let before = store
.get_knowledge_item(id.as_str())
.await
.unwrap()
.expect("item present");
assert_eq!(before.state, KnowledgeState::Pending);
assert_eq!(before.source, ClaimOrigin::AssistantInferred);
use_candidate_once(&store, &id).await.unwrap();
let after = store
.get_knowledge_item(id.as_str())
.await
.unwrap()
.expect("item still present");
assert_eq!(
after.state,
KnowledgeState::Pending,
"using a candidate must not confirm it"
);
assert_eq!(
after.source,
ClaimOrigin::AssistantInferred,
"using a candidate must not change its origin"
);
assert_eq!(after.updated_unix_ms, before.updated_unix_ms);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn use_candidate_once_admits_it_within_the_scope() {
let root = temp_root("use_once_admits");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let (p, obj, id) = seed_one_pending_item(&store).await;
let tree = schema_tree_for(&[(
"orders",
table(&[("id", "bigint", false), ("amount", "numeric", false)]),
)]);
let terms: Vec<String> = vec!["orders".to_string()];
let without = recall(
&store,
recall_admitting_one(
std::slice::from_ref(&p),
&[(p.clone(), avail(tree.clone()))],
&terms,
None,
),
)
.await;
assert!(
without.contracts.is_empty(),
"a candidate is not recallable under confirmed without admission"
);
let with = recall(
&store,
recall_admitting_one(
std::slice::from_ref(&p),
&[(p.clone(), avail(tree))],
&terms,
Some(id.clone()),
),
)
.await;
assert_eq!(
with.contracts.len(),
1,
"the admitted candidate is selected"
);
assert_eq!(with.contracts[0].object, obj);
let admitted: Vec<ClaimId> = with.contracts[0]
.claims
.iter()
.map(|c| c.id.clone())
.collect();
assert!(
admitted.contains(&id),
"the named candidate is in the contract"
);
let stored_admitted = with.contracts[0]
.claims
.iter()
.find(|c| c.id == id)
.expect("the admitted claim is in the contract");
assert_eq!(
stored_admitted.status,
ClaimStatus::Candidate,
"the admitted claim stays Candidate — it still renders unconfirmed"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn use_candidate_once_does_not_persist_across_recalls() {
let root = temp_root("use_once_not_persisted");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let (p, _obj, id) = seed_one_pending_item(&store).await;
let tree = schema_tree_for(&[(
"orders",
table(&[("id", "bigint", false), ("amount", "numeric", false)]),
)]);
let terms: Vec<String> = vec!["orders".to_string()];
let schemas = &[(p.clone(), avail(tree))][..];
let profiles = std::slice::from_ref(&p);
let first = recall(
&store,
recall_admitting_one(profiles, schemas, &terms, Some(id)),
)
.await;
assert_eq!(first.contracts.len(), 1);
let second = recall(
&store,
recall_admitting_one(profiles, schemas, &terms, None),
)
.await;
assert!(
second.contracts.is_empty(),
"the admission is request-scoped: a later recall excludes the candidate again"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn use_candidate_once_refuses_non_candidate_claims() {
let root = temp_root("use_once_refuses");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::table_alias("rejected").unwrap(),
KnowledgeState::Dismissed,
)
.await;
let rejected_id = item_id_for(&store, &obj, &KnowledgeSlot::TableAlias).await;
let err = use_candidate_once(&store, &rejected_id).await.unwrap_err();
assert_eq!(err, ContractOpError::NotACandidate, "dismissed is refused");
assert_eq!(
item_state(&store, &rejected_id).await,
KnowledgeState::Dismissed,
"the refusal changed nothing"
);
put_item(
&store,
&obj,
ClaimPayload::table_description("active fact").unwrap(),
KnowledgeState::Active,
)
.await;
let active_id = item_id_for(&store, &obj, &KnowledgeSlot::TableDescription).await;
let err = use_candidate_once(&store, &active_id).await.unwrap_err();
assert_eq!(err, ContractOpError::NotACandidate, "active is refused");
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn use_candidate_once_refuses_a_confirmed_claim() {
let root = temp_root("use_once_confirmed");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::table_alias("orders").unwrap(),
KnowledgeState::Active,
)
.await;
let confirmed_id = item_id_for(&store, &obj, &KnowledgeSlot::TableAlias).await;
let err = use_candidate_once(&store, &confirmed_id).await.unwrap_err();
assert_eq!(
err,
ContractOpError::NotACandidate,
"a confirmed item is refused, not a silent no-op"
);
assert_eq!(
item_state(&store, &confirmed_id).await,
KnowledgeState::Active
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn use_candidate_once_admits_only_the_named_claim() {
use saya_store::KnowledgeItemStore;
let root = temp_root("use_once_only_named");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let (p, obj, admitted_id) = seed_one_pending_item(&store).await;
put_item(
&store,
&obj,
ClaimPayload::table_alias("sibling").unwrap(),
KnowledgeState::Pending,
)
.await;
let sibling_id = ClaimId::parse(
&store
.knowledge_for_object(&obj)
.await
.expect("knowledge items listed")
.into_iter()
.find(|i| i.slot == KnowledgeSlot::TableAlias)
.expect("sibling item stored")
.id,
)
.expect("ki id");
assert_ne!(sibling_id, admitted_id);
let tree = schema_tree_for(&[(
"orders",
table(&[("id", "bigint", false), ("amount", "numeric", false)]),
)]);
let terms: Vec<String> = vec!["orders".to_string()];
let outcome = recall(
&store,
recall_admitting_one(
std::slice::from_ref(&p),
&[(p.clone(), avail(tree))],
&terms,
Some(admitted_id.clone()),
),
)
.await;
assert_eq!(outcome.contracts.len(), 1, "the object is selected");
let ids: Vec<ClaimId> = outcome.contracts[0]
.claims
.iter()
.map(|c| c.id.clone())
.collect();
assert!(
ids.contains(&admitted_id),
"the named candidate is admitted"
);
assert!(
!ids.contains(&sibling_id),
"the sibling candidate is not admitted — admission is per-claim"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn use_candidate_once_creates_no_evidence() {
let root = temp_root("use_once_no_evidence");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "orders");
put_item(
&store,
&obj,
ClaimPayload::column_description("amount", "how much").unwrap(),
KnowledgeState::Pending,
)
.await;
let id = item_id_for(
&store,
&obj,
&KnowledgeSlot::ColumnDescription {
column: "amount".into(),
},
)
.await;
let before = store
.get_knowledge_item(id.as_str())
.await
.unwrap()
.expect("item present");
use_candidate_once(&store, &id).await.unwrap();
let after = store
.get_knowledge_item(id.as_str())
.await
.unwrap()
.expect("item still present");
assert_eq!(after.state, KnowledgeState::Pending);
assert_eq!(
after.updated_unix_ms, before.updated_unix_ms,
"using a candidate writes nothing; the timestamp is unchanged"
);
let _ = fs::remove_dir_all(root);
}
#[tokio::test]
async fn remember_single_slot_different_value_replaces_and_names_previous() {
let root = temp_root("remember_replace_single_slot");
let db = root.join("state.sqlite3");
let store = store_at(&db).await;
let p = profile_a();
let obj = object_ref(&p, "rental");
let first_grain = ClaimPayload::table_grain("one row per rental", None).unwrap();
let second_grain = ClaimPayload::table_grain("one row per rental per day", None).unwrap();
let unobserved_fp = crate::commands::unobserved_fingerprint();
let outcome1 = remember(&store, &obj, &first_grain, unobserved_fp.clone())
.await
.unwrap();
let id1 = match outcome1 {
RememberOutcome::Stored { id } => id,
other => panic!("first remember should be Stored, got {other:?}"),
};
let items1 = store.knowledge_for_object(&obj).await.unwrap();
assert_eq!(items1.len(), 1);
assert_eq!(items1[0].value, first_grain);
let outcome2 = remember(&store, &obj, &second_grain, unobserved_fp.clone())
.await
.unwrap();
match outcome2 {
RememberOutcome::Replaced { id, previous } => {
assert_eq!(id, id1);
assert_eq!(previous, "one row per rental");
}
other => panic!("expected Replaced, got {other:?}"),
}
let items2 = store.knowledge_for_object(&obj).await.unwrap();
assert_eq!(items2.len(), 1);
assert_eq!(
items2[0].value, second_grain,
"store must contain the second (replacement) value"
);
forget(&store, &id1, ForgetReason::Incorrect).await.unwrap();
let third_grain = ClaimPayload::table_grain("one row per customer rental", None).unwrap();
let outcome3 = remember(&store, &obj, &third_grain, unobserved_fp)
.await
.unwrap();
match outcome3 {
RememberOutcome::Duplicate { id, state } => {
assert_eq!(id, id1);
assert_eq!(state, KnowledgeState::Dismissed);
}
other => panic!("expected Duplicate with Dismissed state, got {other:?}"),
}
let _ = fs::remove_dir_all(root);
}