1use crate::{
2 ClaimId, DailyReviewId, EvidenceId, EvidenceRelation, MemoryRef, PersonId, RetrievalItem,
3 RetrievalPack, SourceId, SourceKind, TenantId, Timestamp,
4};
5use rusqlite::{Connection, OptionalExtension, Transaction, params};
6use serde::{Deserialize, Serialize};
7use sha2::{Digest, Sha256};
8use std::{
9 collections::{HashMap, HashSet},
10 path::Path,
11};
12
13#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
14enum RetrievalTarget {
15 Source(SourceId),
16 Evidence(EvidenceId),
17 Claim(ClaimId),
18}
19
20pub trait Embedder {
21 type Error: std::error::Error + Send + Sync + 'static;
22
23 fn embed(&self, input: &str) -> std::result::Result<Embedding, Self::Error>;
24}
25
26#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
27pub struct Embedding {
28 pub vector: Vec<f32>,
29 pub model: String,
30 pub version: String,
31 pub input_hash: String,
32 pub normalization: VectorNormalization,
33 pub distance: VectorDistance,
34}
35
36#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
37#[serde(rename_all = "snake_case")]
38pub enum VectorNormalization {
39 None,
40 L2,
41}
42
43#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
44#[serde(rename_all = "snake_case")]
45pub enum VectorDistance {
46 Cosine,
47 Dot,
48 Euclidean,
49}
50
51#[derive(Clone, Debug, Eq, Hash, PartialEq, Serialize, Deserialize)]
52#[serde(tag = "kind", content = "id", rename_all = "snake_case")]
53pub enum EmbeddingTarget {
54 Source(SourceId),
55 Evidence(EvidenceId),
56 Claim(ClaimId),
57}
58
59#[derive(Debug, Deserialize)]
60pub struct EmbeddingInput {
61 pub tenant_id: TenantId,
62 pub person_id: PersonId,
63 pub target: EmbeddingTarget,
64 pub embedding: Embedding,
65}
66
67#[derive(Debug, Serialize)]
68pub struct StoredEmbedding {
69 pub target: EmbeddingTarget,
70 pub dimension: usize,
71 pub target_revision: i64,
72 pub input_hash: String,
73 pub created_at: Timestamp,
74}
75
76#[derive(Debug, Deserialize)]
77pub struct ProjectionAuditInput {
78 pub tenant_id: TenantId,
79 pub person_id: PersonId,
80 pub model: String,
81 pub version: String,
82 #[serde(default = "default_limit")]
83 pub limit: u32,
84}
85
86#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
87#[serde(rename_all = "snake_case")]
88pub enum ProjectionState {
89 Missing,
90 Stale,
91}
92
93#[derive(Debug, Serialize)]
94pub struct ProjectionInput {
95 pub target: EmbeddingTarget,
96 pub text: String,
97 pub target_revision: i64,
98 pub input_hash: String,
99}
100
101#[derive(Debug, Serialize)]
102pub struct ProjectionIssue {
103 pub state: ProjectionState,
104 pub input: ProjectionInput,
105 pub stored_target_revision: Option<i64>,
106 pub stored_input_hash: Option<String>,
107 pub stored_created_at: Option<Timestamp>,
108}
109
110#[derive(Debug, thiserror::Error)]
111pub enum Error {
112 #[error(transparent)]
113 Sql(#[from] rusqlite::Error),
114 #[error(transparent)]
115 Json(#[from] serde_json::Error),
116 #[error("{0}")]
117 Invalid(String),
118 #[error("record not found")]
119 NotFound,
120}
121
122pub type Result<T> = std::result::Result<T, Error>;
123
124#[derive(Debug, Deserialize)]
125pub struct RememberInput {
126 pub tenant_id: TenantId,
127 pub person_id: PersonId,
128 #[serde(default)]
129 pub ingestion_key: Option<String>,
130 pub kind: SourceKind,
131 pub text: String,
132 pub captured_at: Timestamp,
133 pub claim: Option<ClaimInput>,
134}
135
136#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
137pub struct TranscriptLocator {
138 pub device_id: String,
139 pub provider: String,
140 pub stream_id: String,
141 pub segment_id: String,
142 pub start_ms: u64,
143 pub end_ms: u64,
144}
145
146#[derive(Debug, Deserialize)]
147pub struct EvidenceLocatorInput {
148 pub tenant_id: TenantId,
149 pub person_id: PersonId,
150 pub evidence_id: EvidenceId,
151}
152
153#[derive(Debug, Deserialize)]
154pub struct RememberRequest {
155 #[serde(flatten)]
156 pub memory: RememberInput,
157 #[serde(default)]
158 pub locator: Option<TranscriptLocator>,
159}
160
161#[derive(Debug, Deserialize)]
162pub struct ClaimInput {
163 pub subject: String,
164 pub predicate: String,
165 pub value: String,
166 pub valid_from: Timestamp,
167}
168
169#[derive(Debug, Serialize)]
170pub struct Remembered {
171 pub source_id: SourceId,
172 pub evidence_id: EvidenceId,
173 pub claim_id: Option<ClaimId>,
174}
175
176#[derive(Debug, Deserialize)]
177pub struct SearchInput {
178 pub tenant_id: TenantId,
179 pub person_id: PersonId,
180 pub query: String,
181 #[serde(default = "default_limit")]
182 pub limit: u32,
183 #[serde(default)]
184 pub query_embedding: Option<DenseQuery>,
185}
186
187#[derive(Debug, Deserialize)]
188pub struct GetInput {
189 pub tenant_id: TenantId,
190 pub person_id: PersonId,
191 pub target: EmbeddingTarget,
192}
193
194#[derive(Debug, Deserialize)]
195pub struct DenseQuery {
196 pub vector: Vec<f32>,
197 pub model: String,
198 pub version: String,
199}
200
201#[derive(Debug, Deserialize)]
202pub struct CorrectInput {
203 pub tenant_id: TenantId,
204 pub person_id: PersonId,
205 pub claim_id: ClaimId,
206 pub text: String,
207 pub value: String,
208 pub occurred_at: Timestamp,
209}
210
211#[derive(Debug, Serialize)]
212pub struct Corrected {
213 pub source_id: SourceId,
214 pub evidence_id: EvidenceId,
215 pub claim_id: ClaimId,
216 pub superseded_claim_id: ClaimId,
217}
218
219#[derive(Debug, Deserialize)]
220pub struct DeleteInput {
221 pub tenant_id: TenantId,
222 pub person_id: PersonId,
223 pub source_id: SourceId,
224 pub deleted_at: Timestamp,
225}
226
227#[derive(Debug, Serialize)]
228pub struct Deleted {
229 pub source_id: SourceId,
230 pub evidence_count: u64,
231 pub claim_count: u64,
232}
233
234#[derive(Debug, Deserialize)]
235pub struct ReviewInput {
236 pub tenant_id: TenantId,
237 pub person_id: PersonId,
238 pub day: String,
239 pub summary: String,
240 pub evidence_ids: Vec<EvidenceId>,
241 pub recorded_at: Timestamp,
242}
243
244#[derive(Debug, Deserialize)]
245pub struct ReviewsInput {
246 pub tenant_id: TenantId,
247 pub person_id: PersonId,
248 #[serde(default = "default_limit")]
249 pub limit: u32,
250}
251
252#[derive(Debug, Serialize)]
253pub struct StoredReview {
254 pub id: DailyReviewId,
255}
256
257#[derive(Debug, Serialize)]
258pub struct ReviewRecord {
259 pub id: DailyReviewId,
260 pub day: String,
261 pub summary: String,
262 pub evidence_ids: Vec<EvidenceId>,
263 pub recorded_at: Timestamp,
264}
265
266pub struct MemoryDb {
267 connection: Connection,
268}
269
270impl MemoryDb {
271 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
272 let connection = Connection::open(path)?;
273 let mut database = Self { connection };
274 database.migrate()?;
275 Ok(database)
276 }
277
278 pub fn remember(&mut self, input: RememberInput) -> Result<Remembered> {
279 self.remember_with_locator(input, None)
280 }
281
282 pub fn remember_with_locator(
283 &mut self,
284 input: RememberInput,
285 locator: Option<TranscriptLocator>,
286 ) -> Result<Remembered> {
287 require_scope(&input.tenant_id, &input.person_id)?;
288 require_text("text", &input.text)?;
289 if let Some(locator) = &locator {
290 validate_transcript_locator(locator)?;
291 }
292 if let Some(key) = &input.ingestion_key {
293 require_text("ingestion_key", key)?;
294 }
295 let transaction = self.connection.transaction()?;
296 let source_id = SourceId(new_id(&transaction)?);
297 let evidence_id = EvidenceId(new_id(&transaction)?);
298 let kind = serde_json::to_string(&input.kind)?;
299 let inserted = transaction.execute(
300 "INSERT OR IGNORE INTO sources(id, tenant_id, person_id, ingestion_key, revision, kind, content, captured_at, recorded_at) VALUES(?1, ?2, ?3, ?4, 1, ?5, ?6, ?7, ?7)",
301 params![source_id.0, input.tenant_id.0, input.person_id.0, input.ingestion_key, kind, input.text, input.captured_at],
302 )?;
303 if inserted == 0 {
304 let replay = transaction.query_row(
305 "SELECT s.id, e.id, c.id, s.kind, s.content, s.captured_at, s.deleted_at, e.deleted_at, c.subject, c.predicate, c.value, c.valid_from FROM sources s JOIN evidence e ON e.source_id = s.id AND e.tenant_id = s.tenant_id AND e.person_id = s.person_id LEFT JOIN claim_evidence ce ON ce.evidence_id = e.id AND ce.tenant_id = e.tenant_id AND ce.person_id = e.person_id LEFT JOIN claims c ON c.id = ce.claim_id AND c.tenant_id = ce.tenant_id AND c.person_id = ce.person_id WHERE s.tenant_id = ?1 AND s.person_id = ?2 AND s.ingestion_key = ?3 ORDER BY e.id, c.id LIMIT 1",
306 params![input.tenant_id.0, input.person_id.0, input.ingestion_key],
307 |row| {
308 Ok((
309 Remembered {
310 source_id: SourceId(row.get(0)?),
311 evidence_id: EvidenceId(row.get(1)?),
312 claim_id: row.get::<_, Option<String>>(2)?.map(ClaimId),
313 },
314 row.get::<_, String>(3)?,
315 row.get::<_, String>(4)?,
316 row.get::<_, i64>(5)?,
317 row.get::<_, Option<i64>>(6)?,
318 row.get::<_, Option<i64>>(7)?,
319 row.get::<_, Option<String>>(8)?,
320 row.get::<_, Option<String>>(9)?,
321 row.get::<_, Option<String>>(10)?,
322 row.get::<_, Option<i64>>(11)?,
323 ))
324 },
325 ).optional()?.ok_or(Error::NotFound)?;
326 if replay.4.is_some() || replay.5.is_some() {
327 return Err(Error::NotFound);
328 }
329 let stored_locator = transaction
330 .query_row(
331 "SELECT device_id, provider, stream_id, segment_id, start_ms, end_ms FROM evidence_locators WHERE tenant_id = ?1 AND person_id = ?2 AND evidence_id = ?3",
332 params![input.tenant_id.0, input.person_id.0, replay.0.evidence_id.0],
333 |row| {
334 Ok(TranscriptLocator {
335 device_id: row.get(0)?,
336 provider: row.get(1)?,
337 stream_id: row.get(2)?,
338 segment_id: row.get(3)?,
339 start_ms: row.get(4)?,
340 end_ms: row.get(5)?,
341 })
342 },
343 )
344 .optional()?;
345 let stored_claim = match (&replay.6, &replay.7, &replay.8, replay.9) {
346 (Some(subject), Some(predicate), Some(value), Some(valid_from)) => {
347 Some((subject, predicate, value, valid_from))
348 }
349 (None, None, None, None) => None,
350 _ => return Err(Error::NotFound),
351 };
352 let input_claim = input.claim.as_ref().map(|claim| {
353 (
354 &claim.subject,
355 &claim.predicate,
356 &claim.value,
357 claim.valid_from,
358 )
359 });
360 if replay.1 != kind
361 || replay.2 != input.text
362 || replay.3 != input.captured_at
363 || stored_claim != input_claim
364 || stored_locator.as_ref() != locator.as_ref()
365 {
366 return Err(Error::Invalid(
367 "ingestion_key conflicts with different memory payload".to_owned(),
368 ));
369 }
370 transaction.commit()?;
371 return Ok(replay.0);
372 }
373 transaction.execute(
374 "INSERT INTO source_fts(source_id, tenant_id, person_id, content) VALUES(?1, ?2, ?3, ?4)",
375 params![source_id.0, input.tenant_id.0, input.person_id.0, input.text],
376 )?;
377 transaction.execute(
378 "INSERT INTO evidence(id, tenant_id, person_id, source_id, source_revision, quote, recorded_at) VALUES(?1, ?2, ?3, ?4, 1, ?5, ?6)",
379 params![evidence_id.0, input.tenant_id.0, input.person_id.0, source_id.0, input.text, input.captured_at],
380 )?;
381 if let Some(locator) = locator {
382 transaction.execute(
383 "INSERT INTO evidence_locators(tenant_id, person_id, evidence_id, device_id, provider, stream_id, segment_id, start_ms, end_ms) VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
384 params![input.tenant_id.0, input.person_id.0, evidence_id.0, locator.device_id, locator.provider, locator.stream_id, locator.segment_id, locator.start_ms, locator.end_ms],
385 )?;
386 }
387 let claim_id = input
388 .claim
389 .map(|claim| {
390 insert_claim(
391 &transaction,
392 &input.tenant_id,
393 &input.person_id,
394 &evidence_id,
395 claim,
396 input.captured_at,
397 )
398 })
399 .transpose()?;
400 transaction.commit()?;
401 Ok(Remembered {
402 source_id,
403 evidence_id,
404 claim_id,
405 })
406 }
407
408 pub fn evidence_locator(
409 &self,
410 input: EvidenceLocatorInput,
411 ) -> Result<Option<TranscriptLocator>> {
412 require_scope(&input.tenant_id, &input.person_id)?;
413 let locator = self
414 .connection
415 .query_row(
416 "SELECT l.device_id, l.provider, l.stream_id, l.segment_id, l.start_ms, l.end_ms FROM evidence_locators l JOIN evidence e ON e.id = l.evidence_id AND e.tenant_id = l.tenant_id AND e.person_id = l.person_id WHERE l.tenant_id = ?1 AND l.person_id = ?2 AND l.evidence_id = ?3 AND e.deleted_at IS NULL",
417 params![input.tenant_id.0, input.person_id.0, input.evidence_id.0],
418 |row| {
419 Ok(TranscriptLocator {
420 device_id: row.get(0)?,
421 provider: row.get(1)?,
422 stream_id: row.get(2)?,
423 segment_id: row.get(3)?,
424 start_ms: row.get(4)?,
425 end_ms: row.get(5)?,
426 })
427 },
428 )
429 .optional()?;
430 Ok(locator)
431 }
432
433 pub fn search(&self, input: SearchInput) -> Result<RetrievalPack> {
434 require_scope(&input.tenant_id, &input.person_id)?;
435 require_text("query", &input.query)?;
436 let limit = bounded_limit(input.limit);
437 let candidate_limit = limit * 4;
438 let query = format!("\"{}\"", input.query.replace('"', "\"\""));
439 let mut statement = self.connection.prepare(
440 "SELECT s.id, c.id
441 FROM source_fts
442 JOIN sources s ON s.id = source_fts.source_id AND s.tenant_id = source_fts.tenant_id AND s.person_id = source_fts.person_id
443 JOIN evidence e ON e.source_id = s.id AND e.tenant_id = s.tenant_id AND e.person_id = s.person_id AND e.deleted_at IS NULL
444 LEFT JOIN claim_evidence ce ON ce.evidence_id = e.id AND ce.tenant_id = e.tenant_id AND ce.person_id = e.person_id
445 LEFT JOIN claims c ON c.id = ce.claim_id AND c.tenant_id = ce.tenant_id AND c.person_id = ce.person_id AND c.status = 'accepted'
446 WHERE source_fts MATCH ?1 AND source_fts.tenant_id = ?2 AND source_fts.person_id = ?3 AND s.deleted_at IS NULL
447 AND (c.id IS NOT NULL OR NOT EXISTS (
448 SELECT 1 FROM evidence live_e
449 JOIN claim_evidence live_ce ON live_ce.evidence_id = live_e.id AND live_ce.tenant_id = live_e.tenant_id AND live_ce.person_id = live_e.person_id
450 JOIN claims live_c ON live_c.id = live_ce.claim_id AND live_c.tenant_id = live_ce.tenant_id AND live_c.person_id = live_ce.person_id
451 WHERE live_e.source_id = s.id AND live_e.tenant_id = s.tenant_id AND live_e.person_id = s.person_id AND live_e.deleted_at IS NULL AND live_c.status = 'accepted'
452 ))
453 ORDER BY bm25(source_fts), s.id, c.id LIMIT ?4",
454 )?;
455 let rows = statement.query_map(
456 params![query, input.tenant_id.0, input.person_id.0, candidate_limit],
457 |row| {
458 let source_id = row.get::<_, String>(0)?;
459 Ok(match row.get::<_, Option<String>>(1)? {
460 Some(claim_id) => RetrievalTarget::Claim(ClaimId(claim_id)),
461 None => RetrievalTarget::Source(SourceId(source_id)),
462 })
463 },
464 )?;
465 let lexical = rows.collect::<std::result::Result<Vec<_>, _>>()?;
466 let dense = input
467 .query_embedding
468 .as_ref()
469 .map(|query| self.dense_claims(&input.tenant_id, &input.person_id, query))
470 .transpose()?
471 .unwrap_or_default();
472 let ranked = reciprocal_rank_fusion(&lexical, &dense, limit as usize);
473 let mut items = Vec::with_capacity(ranked.len());
474 for (target, relevance_basis_points) in ranked {
475 items.push(self.retrieval_item(
476 &input.tenant_id,
477 &input.person_id,
478 target,
479 relevance_basis_points,
480 )?);
481 }
482 let gaps = if items.is_empty() {
483 vec!["no cited memory matched".to_owned()]
484 } else {
485 Vec::new()
486 };
487 Ok(RetrievalPack {
488 query: input.query,
489 items,
490 gaps,
491 })
492 }
493
494 pub fn get(&self, input: GetInput) -> Result<RetrievalItem> {
495 require_scope(&input.tenant_id, &input.person_id)?;
496 let target = match input.target {
497 EmbeddingTarget::Source(id) => RetrievalTarget::Source(id),
498 EmbeddingTarget::Evidence(id) => RetrievalTarget::Evidence(id),
499 EmbeddingTarget::Claim(id) => RetrievalTarget::Claim(id),
500 };
501 self.retrieval_item(&input.tenant_id, &input.person_id, target, 10_000)
502 }
503
504 fn dense_claims(
505 &self,
506 tenant_id: &TenantId,
507 person_id: &PersonId,
508 query: &DenseQuery,
509 ) -> Result<Vec<RetrievalTarget>> {
510 validate_dense_query(query)?;
511 let mut statement = self.connection.prepare(
512 "SELECT embeddings.target_kind, embeddings.target_id, embeddings.dimension, embeddings.normalization, embeddings.distance, embeddings.vector, embeddings.target_revision, embeddings.input_hash
513 FROM embeddings
514 WHERE embeddings.tenant_id = ?1 AND embeddings.person_id = ?2 AND embeddings.model = ?3 AND embeddings.version = ?4",
515 )?;
516 let rows = statement.query_map(
517 params![tenant_id.0, person_id.0, query.model, query.version],
518 |row| {
519 Ok((
520 row.get::<_, String>(0)?,
521 row.get::<_, String>(1)?,
522 row.get::<_, usize>(2)?,
523 row.get::<_, String>(3)?,
524 row.get::<_, String>(4)?,
525 row.get::<_, String>(5)?,
526 row.get::<_, i64>(6)?,
527 row.get::<_, String>(7)?,
528 ))
529 },
530 )?;
531 let mut scores = HashMap::<RetrievalTarget, f32>::new();
532 let mut lane = None;
533 for row in rows {
534 let (
535 target_kind,
536 target_id,
537 dimension,
538 normalization,
539 distance,
540 vector,
541 target_revision,
542 input_hash,
543 ) = row?;
544 let target = embedding_target(&target_kind, &target_id)?;
545 let current = match self.projection_input(tenant_id, person_id, target) {
546 Ok(current) => current,
547 Err(Error::NotFound) => continue,
548 Err(error) => return Err(error),
549 };
550 if current.target_revision != target_revision || current.input_hash != input_hash {
551 continue;
552 }
553 if dimension != query.vector.len() {
554 return Err(Error::Invalid(format!(
555 "query embedding dimension {} does not match stored dimension {dimension}",
556 query.vector.len()
557 )));
558 }
559 let normalization: VectorNormalization = serde_json::from_str(&normalization)?;
560 let distance: VectorDistance = serde_json::from_str(&distance)?;
561 let configuration = (normalization, distance, dimension);
562 if lane.as_ref().is_some_and(|lane| lane != &configuration) {
563 return Err(Error::Invalid(
564 "stored embeddings mix incompatible vector configurations".to_owned(),
565 ));
566 }
567 lane = Some(configuration.clone());
568 let vector: Vec<f32> = serde_json::from_str(&vector)?;
569 if vector.len() != dimension || vector.iter().any(|value| !value.is_finite()) {
570 return Err(Error::Invalid("stored embedding is invalid".to_owned()));
571 }
572 let score = vector_score(&query.vector, &vector, &configuration.1)?;
573 for target in self.retrieval_targets_for_embedding(
574 tenant_id,
575 person_id,
576 &target_kind,
577 &target_id,
578 )? {
579 scores
580 .entry(target)
581 .and_modify(|existing| *existing = existing.max(score))
582 .or_insert(score);
583 }
584 }
585 let mut ranked = scores.into_iter().collect::<Vec<_>>();
586 ranked.sort_by(|left, right| {
587 right
588 .1
589 .total_cmp(&left.1)
590 .then_with(|| left.0.cmp(&right.0))
591 });
592 Ok(ranked.into_iter().map(|(target, _)| target).collect())
593 }
594
595 fn retrieval_targets_for_embedding(
596 &self,
597 tenant_id: &TenantId,
598 person_id: &PersonId,
599 target_kind: &str,
600 target_id: &str,
601 ) -> Result<Vec<RetrievalTarget>> {
602 let sql = match target_kind {
603 "claim" => {
604 "SELECT id FROM claims WHERE id = ?1 AND tenant_id = ?2 AND person_id = ?3 AND status = 'accepted'"
605 }
606 "evidence" => {
607 "SELECT c.id FROM evidence e JOIN sources s ON s.id = e.source_id AND s.tenant_id = e.tenant_id AND s.person_id = e.person_id LEFT JOIN claim_evidence ce ON ce.evidence_id = e.id AND ce.tenant_id = e.tenant_id AND ce.person_id = e.person_id LEFT JOIN claims c ON c.id = ce.claim_id AND c.tenant_id = ce.tenant_id AND c.person_id = ce.person_id AND c.status = 'accepted' WHERE e.id = ?1 AND e.tenant_id = ?2 AND e.person_id = ?3 AND e.deleted_at IS NULL AND s.deleted_at IS NULL ORDER BY c.id"
608 }
609 "source" => {
610 "SELECT DISTINCT c.id FROM sources s JOIN evidence e ON e.source_id = s.id AND e.tenant_id = s.tenant_id AND e.person_id = s.person_id LEFT JOIN claim_evidence ce ON ce.evidence_id = e.id AND ce.tenant_id = e.tenant_id AND ce.person_id = e.person_id LEFT JOIN claims c ON c.id = ce.claim_id AND c.tenant_id = ce.tenant_id AND c.person_id = ce.person_id AND c.status = 'accepted' WHERE s.id = ?1 AND s.tenant_id = ?2 AND s.person_id = ?3 AND s.deleted_at IS NULL AND e.deleted_at IS NULL ORDER BY c.id"
611 }
612 _ => {
613 return Err(Error::Invalid(
614 "stored embedding target is invalid".to_owned(),
615 ));
616 }
617 };
618 let mut statement = self.connection.prepare(sql)?;
619 let rows = statement.query_map(params![target_id, tenant_id.0, person_id.0], |row| {
620 row.get::<_, Option<String>>(0)
621 })?;
622 let rows = rows.collect::<std::result::Result<Vec<_>, _>>()?;
623 if rows.is_empty() {
624 return Ok(Vec::new());
625 }
626 let claims = rows
627 .into_iter()
628 .flatten()
629 .map(|id| RetrievalTarget::Claim(ClaimId(id)))
630 .collect::<Vec<_>>();
631 if !claims.is_empty() {
632 return Ok(claims);
633 }
634 Ok(match target_kind {
635 "source" => vec![RetrievalTarget::Source(SourceId(target_id.to_owned()))],
636 "evidence" => vec![RetrievalTarget::Evidence(EvidenceId(target_id.to_owned()))],
637 "claim" => Vec::new(),
638 _ => unreachable!(),
639 })
640 }
641
642 fn retrieval_item(
643 &self,
644 tenant_id: &TenantId,
645 person_id: &PersonId,
646 target: RetrievalTarget,
647 relevance_basis_points: u16,
648 ) -> Result<RetrievalItem> {
649 let (sql, id) = match &target {
650 RetrievalTarget::Claim(id) => (
651 "SELECT c.subject || ' ' || c.predicate || ' ' || c.value, ce.evidence_id
652 FROM claims c
653 JOIN claim_evidence ce ON ce.claim_id = c.id AND ce.tenant_id = c.tenant_id AND ce.person_id = c.person_id
654 JOIN evidence e ON e.id = ce.evidence_id AND e.tenant_id = ce.tenant_id AND e.person_id = ce.person_id
655 JOIN sources s ON s.id = e.source_id AND s.tenant_id = e.tenant_id AND s.person_id = e.person_id
656 WHERE c.id = ?1 AND c.tenant_id = ?2 AND c.person_id = ?3 AND c.status = 'accepted' AND e.deleted_at IS NULL AND s.deleted_at IS NULL
657 ORDER BY ce.evidence_id LIMIT 1",
658 &id.0,
659 ),
660 RetrievalTarget::Source(id) => (
661 "SELECT s.content, e.id FROM sources s JOIN evidence e ON e.source_id = s.id AND e.tenant_id = s.tenant_id AND e.person_id = s.person_id WHERE s.id = ?1 AND s.tenant_id = ?2 AND s.person_id = ?3 AND s.deleted_at IS NULL AND e.deleted_at IS NULL ORDER BY e.id LIMIT 1",
662 &id.0,
663 ),
664 RetrievalTarget::Evidence(id) => (
665 "SELECT e.quote, e.id FROM evidence e JOIN sources s ON s.id = e.source_id AND s.tenant_id = e.tenant_id AND s.person_id = e.person_id WHERE e.id = ?1 AND e.tenant_id = ?2 AND e.person_id = ?3 AND e.deleted_at IS NULL AND s.deleted_at IS NULL",
666 &id.0,
667 ),
668 };
669 let memory = match &target {
670 RetrievalTarget::Claim(id) => MemoryRef::Claim(id.clone()),
671 RetrievalTarget::Source(id) => MemoryRef::Source(id.clone()),
672 RetrievalTarget::Evidence(id) => MemoryRef::Evidence(id.clone()),
673 };
674 self.connection
675 .query_row(sql, params![id, tenant_id.0, person_id.0], |row| {
676 Ok(RetrievalItem {
677 memory,
678 excerpt: row.get(0)?,
679 relevance_basis_points,
680 evidence_ids: vec![EvidenceId(row.get(1)?)],
681 })
682 })
683 .optional()?
684 .ok_or(Error::NotFound)
685 }
686
687 pub fn correct(&mut self, input: CorrectInput) -> Result<Corrected> {
688 require_scope(&input.tenant_id, &input.person_id)?;
689 require_text("correction text", &input.text)?;
690 require_text("value", &input.value)?;
691 let transaction = self.connection.transaction()?;
692 let old = transaction
693 .query_row(
694 "SELECT subject, predicate, valid_from, recorded_from FROM claims WHERE id = ?1 AND tenant_id = ?2 AND person_id = ?3 AND status = 'accepted'",
695 params![input.claim_id.0, input.tenant_id.0, input.person_id.0],
696 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?, row.get::<_, i64>(2)?, row.get::<_, i64>(3)?)),
697 )
698 .optional()?
699 .ok_or(Error::NotFound)?;
700 if input.occurred_at <= old.3 {
701 return Err(Error::Invalid(
702 "occurred_at must be after the original claim was recorded".to_owned(),
703 ));
704 }
705 let source_id = SourceId(new_id(&transaction)?);
706 let evidence_id = EvidenceId(new_id(&transaction)?);
707 transaction.execute(
708 "INSERT INTO sources(id, tenant_id, person_id, revision, kind, content, captured_at, recorded_at) VALUES(?1, ?2, ?3, 1, '\"user_correction\"', ?4, ?5, ?5)",
709 params![source_id.0, input.tenant_id.0, input.person_id.0, input.text, input.occurred_at],
710 )?;
711 transaction.execute(
712 "INSERT INTO source_fts(source_id, tenant_id, person_id, content) VALUES(?1, ?2, ?3, ?4)",
713 params![source_id.0, input.tenant_id.0, input.person_id.0, input.text],
714 )?;
715 transaction.execute(
716 "INSERT INTO evidence(id, tenant_id, person_id, source_id, source_revision, quote, recorded_at) VALUES(?1, ?2, ?3, ?4, 1, ?5, ?6)",
717 params![evidence_id.0, input.tenant_id.0, input.person_id.0, source_id.0, input.text, input.occurred_at],
718 )?;
719 transaction.execute(
720 "UPDATE claims SET status = 'superseded', recorded_until = ?1 WHERE id = ?2 AND tenant_id = ?3 AND person_id = ?4",
721 params![input.occurred_at, input.claim_id.0, input.tenant_id.0, input.person_id.0],
722 )?;
723 let claim_id = insert_claim(
724 &transaction,
725 &input.tenant_id,
726 &input.person_id,
727 &evidence_id,
728 ClaimInput {
729 subject: old.0,
730 predicate: old.1,
731 value: input.value,
732 valid_from: input.occurred_at.max(old.2),
733 },
734 input.occurred_at,
735 )?;
736 transaction.commit()?;
737 Ok(Corrected {
738 source_id,
739 evidence_id,
740 claim_id,
741 superseded_claim_id: input.claim_id,
742 })
743 }
744
745 pub fn delete_source(&mut self, input: DeleteInput) -> Result<Deleted> {
746 require_scope(&input.tenant_id, &input.person_id)?;
747 let transaction = self.connection.transaction()?;
748 let recorded_at = transaction
749 .query_row(
750 "SELECT recorded_at FROM sources WHERE id = ?1 AND tenant_id = ?2 AND person_id = ?3 AND deleted_at IS NULL",
751 params![input.source_id.0, input.tenant_id.0, input.person_id.0],
752 |row| row.get::<_, i64>(0),
753 )
754 .optional()?
755 .ok_or(Error::NotFound)?;
756 if input.deleted_at < recorded_at {
757 return Err(Error::Invalid(
758 "deleted_at cannot predate source recording".to_owned(),
759 ));
760 }
761 let changed = transaction.execute(
762 "UPDATE sources SET deleted_at = ?1, revision = revision + 1 WHERE id = ?2 AND tenant_id = ?3 AND person_id = ?4 AND deleted_at IS NULL",
763 params![input.deleted_at, input.source_id.0, input.tenant_id.0, input.person_id.0],
764 )?;
765 if changed == 0 {
766 return Err(Error::NotFound);
767 }
768 transaction.execute(
769 "DELETE FROM source_fts WHERE source_id = ?1 AND tenant_id = ?2 AND person_id = ?3",
770 params![input.source_id.0, input.tenant_id.0, input.person_id.0],
771 )?;
772 let evidence_count = transaction.execute(
773 "UPDATE evidence SET deleted_at = ?1 WHERE source_id = ?2 AND tenant_id = ?3 AND person_id = ?4 AND deleted_at IS NULL",
774 params![input.deleted_at, input.source_id.0, input.tenant_id.0, input.person_id.0],
775 )? as u64;
776 let claim_count = transaction.execute(
777 "UPDATE claims SET status = 'rejected', recorded_until = ?1
778 WHERE tenant_id = ?2 AND person_id = ?3 AND status = 'accepted'
779 AND id IN (SELECT ce.claim_id FROM claim_evidence ce JOIN evidence e ON e.id = ce.evidence_id AND e.tenant_id = ce.tenant_id AND e.person_id = ce.person_id WHERE e.source_id = ?4)
780 AND NOT EXISTS (SELECT 1 FROM claim_evidence live_ce JOIN evidence live_e ON live_e.id = live_ce.evidence_id AND live_e.tenant_id = live_ce.tenant_id AND live_e.person_id = live_ce.person_id WHERE live_ce.claim_id = claims.id AND live_e.deleted_at IS NULL)",
781 params![input.deleted_at, input.tenant_id.0, input.person_id.0, input.source_id.0],
782 )? as u64;
783 transaction.execute(
784 "DELETE FROM embeddings WHERE tenant_id = ?1 AND person_id = ?2 AND ((target_kind = 'source' AND target_id = ?3) OR (target_kind = 'evidence' AND target_id IN (SELECT id FROM evidence WHERE source_id = ?3 AND tenant_id = ?1 AND person_id = ?2)) OR (target_kind = 'claim' AND target_id IN (SELECT id FROM claims WHERE tenant_id = ?1 AND person_id = ?2 AND status = 'rejected')))",
785 params![input.tenant_id.0, input.person_id.0, input.source_id.0],
786 )?;
787 transaction.execute(
788 "DELETE FROM daily_reviews WHERE tenant_id = ?1 AND person_id = ?2 AND EXISTS (SELECT 1 FROM json_each(evidence_ids) citation JOIN evidence e ON e.id = citation.value WHERE e.source_id = ?3 AND e.tenant_id = ?1 AND e.person_id = ?2)",
789 params![input.tenant_id.0, input.person_id.0, input.source_id.0],
790 )?;
791 transaction.commit()?;
792 Ok(Deleted {
793 source_id: input.source_id,
794 evidence_count,
795 claim_count,
796 })
797 }
798
799 pub fn store_review(&mut self, input: ReviewInput) -> Result<StoredReview> {
800 require_scope(&input.tenant_id, &input.person_id)?;
801 require_text("day", &input.day)?;
802 require_text("summary", &input.summary)?;
803 if input.evidence_ids.is_empty() {
804 return Err(Error::Invalid("review needs evidence_ids".to_owned()));
805 }
806 let transaction = self.connection.transaction()?;
807 for evidence_id in &input.evidence_ids {
808 let found: bool = transaction.query_row(
809 "SELECT EXISTS(SELECT 1 FROM evidence WHERE id = ?1 AND tenant_id = ?2 AND person_id = ?3 AND deleted_at IS NULL)",
810 params![evidence_id.0, input.tenant_id.0, input.person_id.0],
811 |row| row.get(0),
812 )?;
813 if !found {
814 return Err(Error::Invalid(format!(
815 "evidence {} is unavailable",
816 evidence_id.0
817 )));
818 }
819 }
820 let id = DailyReviewId(new_id(&transaction)?);
821 transaction.execute(
822 "INSERT INTO daily_reviews(id, tenant_id, person_id, day, summary, evidence_ids, recorded_at) VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7)",
823 params![id.0, input.tenant_id.0, input.person_id.0, input.day, input.summary, serde_json::to_string(&input.evidence_ids)?, input.recorded_at],
824 )?;
825 transaction.commit()?;
826 Ok(StoredReview { id })
827 }
828
829 pub fn reviews(&self, input: ReviewsInput) -> Result<Vec<ReviewRecord>> {
830 require_scope(&input.tenant_id, &input.person_id)?;
831 let mut statement = self.connection.prepare(
832 "SELECT id, day, summary, evidence_ids, recorded_at FROM daily_reviews WHERE tenant_id = ?1 AND person_id = ?2 ORDER BY day DESC, recorded_at DESC LIMIT ?3",
833 )?;
834 let rows = statement.query_map(
835 params![
836 input.tenant_id.0,
837 input.person_id.0,
838 bounded_limit(input.limit)
839 ],
840 |row| {
841 let json: String = row.get(3)?;
842 let evidence_ids = serde_json::from_str(&json).map_err(|error| {
843 rusqlite::Error::FromSqlConversionFailure(
844 3,
845 rusqlite::types::Type::Text,
846 Box::new(error),
847 )
848 })?;
849 Ok(ReviewRecord {
850 id: DailyReviewId(row.get(0)?),
851 day: row.get(1)?,
852 summary: row.get(2)?,
853 evidence_ids,
854 recorded_at: row.get(4)?,
855 })
856 },
857 )?;
858 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
859 }
860
861 pub fn rebuild_embedding<E: Embedder>(
862 &mut self,
863 tenant_id: TenantId,
864 person_id: PersonId,
865 target: EmbeddingTarget,
866 embedder: &E,
867 ) -> Result<StoredEmbedding> {
868 require_scope(&tenant_id, &person_id)?;
869 let projection = self.projection_input(&tenant_id, &person_id, target.clone())?;
870 let embedding = embedder
871 .embed(&projection.text)
872 .map_err(|error| Error::Invalid(format!("embedder failed: {error}")))?;
873 self.upsert_embedding(EmbeddingInput {
874 tenant_id,
875 person_id,
876 target,
877 embedding,
878 })
879 }
880
881 pub fn upsert_embedding(&mut self, input: EmbeddingInput) -> Result<StoredEmbedding> {
882 require_scope(&input.tenant_id, &input.person_id)?;
883 validate_embedding(&input.embedding)?;
884 let projection =
885 self.projection_input(&input.tenant_id, &input.person_id, input.target.clone())?;
886 if input.embedding.input_hash != projection.input_hash {
887 return Err(Error::Invalid(format!(
888 "embedding input_hash does not match current target input; expected {}",
889 projection.input_hash
890 )));
891 }
892 let (target_kind, target_id) = embedding_target_parts(&input.target);
893 let dimension = input.embedding.vector.len();
894 let created_at = self
895 .connection
896 .query_row("SELECT unixepoch()", [], |row| row.get(0))?;
897 self.connection.execute(
898 "INSERT INTO embeddings(tenant_id, person_id, target_kind, target_id, model, version, dimension, input_hash, target_revision, created_at, normalization, distance, vector)
899 VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)
900 ON CONFLICT(tenant_id, person_id, target_kind, target_id, model, version) DO UPDATE SET dimension = excluded.dimension, input_hash = excluded.input_hash, target_revision = excluded.target_revision, created_at = excluded.created_at, normalization = excluded.normalization, distance = excluded.distance, vector = excluded.vector",
901 params![input.tenant_id.0, input.person_id.0, target_kind, target_id, input.embedding.model, input.embedding.version, dimension, projection.input_hash, projection.target_revision, created_at, serde_json::to_string(&input.embedding.normalization)?, serde_json::to_string(&input.embedding.distance)?, serde_json::to_string(&input.embedding.vector)?],
902 )?;
903 Ok(StoredEmbedding {
904 target: input.target,
905 dimension,
906 target_revision: projection.target_revision,
907 input_hash: projection.input_hash,
908 created_at,
909 })
910 }
911
912 pub fn projection_input(
913 &self,
914 tenant_id: &TenantId,
915 person_id: &PersonId,
916 target: EmbeddingTarget,
917 ) -> Result<ProjectionInput> {
918 require_scope(tenant_id, person_id)?;
919 let (table, id, expression, revision, live) = match &target {
920 EmbeddingTarget::Source(id) => (
921 "sources",
922 &id.0,
923 "content",
924 "revision",
925 "deleted_at IS NULL",
926 ),
927 EmbeddingTarget::Evidence(id) => (
928 "evidence",
929 &id.0,
930 "quote",
931 "source_revision",
932 "deleted_at IS NULL",
933 ),
934 EmbeddingTarget::Claim(id) => (
935 "claims",
936 &id.0,
937 "subject || ' ' || predicate || ' ' || value",
938 "recorded_from",
939 "status = 'accepted'",
940 ),
941 };
942 let (text, target_revision) = self
943 .connection
944 .query_row(
945 &format!("SELECT {expression}, {revision} FROM {table} WHERE id = ?1 AND tenant_id = ?2 AND person_id = ?3 AND {live}"),
946 params![id, tenant_id.0, person_id.0],
947 |row| Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?)),
948 )
949 .optional()?
950 .ok_or(Error::NotFound)?;
951 Ok(ProjectionInput {
952 target,
953 input_hash: input_hash(&text),
954 text,
955 target_revision,
956 })
957 }
958
959 pub fn projection_issues(&self, input: ProjectionAuditInput) -> Result<Vec<ProjectionIssue>> {
960 require_scope(&input.tenant_id, &input.person_id)?;
961 require_text("embedding model", &input.model)?;
962 require_text("embedding version", &input.version)?;
963 let mut statement = self.connection.prepare(
964 "SELECT target_kind, target_id FROM (
965 SELECT 'source' AS target_kind, id AS target_id FROM sources WHERE tenant_id = ?1 AND person_id = ?2 AND deleted_at IS NULL
966 UNION ALL SELECT 'evidence', e.id FROM evidence e JOIN sources s ON s.id = e.source_id AND s.tenant_id = e.tenant_id AND s.person_id = e.person_id WHERE e.tenant_id = ?1 AND e.person_id = ?2 AND e.deleted_at IS NULL AND s.deleted_at IS NULL
967 UNION ALL SELECT 'claim', c.id FROM claims c WHERE c.tenant_id = ?1 AND c.person_id = ?2 AND c.status = 'accepted'
968 ) ORDER BY target_kind, target_id",
969 )?;
970 let targets = statement
971 .query_map(params![input.tenant_id.0, input.person_id.0], |row| {
972 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
973 })?
974 .collect::<std::result::Result<Vec<_>, _>>()?;
975 let mut issues = Vec::new();
976 let limit = bounded_limit(input.limit) as usize;
977 for (kind, id) in targets {
978 let projection = self.projection_input(
979 &input.tenant_id,
980 &input.person_id,
981 embedding_target(&kind, &id)?,
982 )?;
983 let stored = self
984 .connection
985 .query_row(
986 "SELECT target_revision, input_hash, created_at FROM embeddings WHERE tenant_id = ?1 AND person_id = ?2 AND target_kind = ?3 AND target_id = ?4 AND model = ?5 AND version = ?6",
987 params![input.tenant_id.0, input.person_id.0, kind, embedding_target_parts(&projection.target).1, input.model, input.version],
988 |row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?, row.get::<_, i64>(2)?)),
989 )
990 .optional()?;
991 let state = match &stored {
992 None => ProjectionState::Missing,
993 Some((revision, hash, _))
994 if *revision != projection.target_revision
995 || hash != &projection.input_hash =>
996 {
997 ProjectionState::Stale
998 }
999 Some(_) => continue,
1000 };
1001 issues.push(ProjectionIssue {
1002 state,
1003 input: projection,
1004 stored_target_revision: stored.as_ref().map(|value| value.0),
1005 stored_input_hash: stored.as_ref().map(|value| value.1.clone()),
1006 stored_created_at: stored.map(|value| value.2),
1007 });
1008 if issues.len() == limit {
1009 break;
1010 }
1011 }
1012 Ok(issues)
1013 }
1014
1015 fn migrate(&mut self) -> Result<()> {
1016 self.connection.execute_batch(
1017 "PRAGMA foreign_keys = ON;
1018 BEGIN IMMEDIATE;
1019 CREATE TABLE IF NOT EXISTS sources(id TEXT PRIMARY KEY, tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, ingestion_key TEXT, revision INTEGER NOT NULL, kind TEXT NOT NULL, content TEXT NOT NULL, captured_at INTEGER NOT NULL, recorded_at INTEGER NOT NULL, deleted_at INTEGER);
1020 CREATE INDEX IF NOT EXISTS sources_scope ON sources(tenant_id, person_id, id);
1021 CREATE VIRTUAL TABLE IF NOT EXISTS source_fts USING fts5(source_id UNINDEXED, tenant_id UNINDEXED, person_id UNINDEXED, content, tokenize='unicode61');
1022 CREATE TABLE IF NOT EXISTS evidence(id TEXT PRIMARY KEY, tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, source_id TEXT NOT NULL REFERENCES sources(id), source_revision INTEGER NOT NULL, quote TEXT NOT NULL, recorded_at INTEGER NOT NULL, deleted_at INTEGER);
1023 CREATE INDEX IF NOT EXISTS evidence_scope ON evidence(tenant_id, person_id, source_id);
1024 CREATE TABLE IF NOT EXISTS evidence_locators(tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, evidence_id TEXT NOT NULL REFERENCES evidence(id), device_id TEXT NOT NULL, provider TEXT NOT NULL, stream_id TEXT NOT NULL, segment_id TEXT NOT NULL, start_ms INTEGER NOT NULL, end_ms INTEGER NOT NULL, PRIMARY KEY(tenant_id, person_id, evidence_id));
1025 CREATE TABLE IF NOT EXISTS claims(id TEXT PRIMARY KEY, tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, subject TEXT NOT NULL, predicate TEXT NOT NULL, value TEXT NOT NULL, valid_from INTEGER NOT NULL, valid_until INTEGER, recorded_from INTEGER NOT NULL, recorded_until INTEGER, status TEXT NOT NULL);
1026 CREATE INDEX IF NOT EXISTS claims_scope ON claims(tenant_id, person_id, status);
1027 CREATE TABLE IF NOT EXISTS claim_evidence(tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, claim_id TEXT NOT NULL REFERENCES claims(id), evidence_id TEXT NOT NULL REFERENCES evidence(id), relation TEXT NOT NULL, confidence_basis_points INTEGER NOT NULL, PRIMARY KEY(tenant_id, person_id, claim_id, evidence_id));
1028 CREATE TABLE IF NOT EXISTS daily_reviews(id TEXT PRIMARY KEY, tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, day TEXT NOT NULL, summary TEXT NOT NULL, evidence_ids TEXT NOT NULL, recorded_at INTEGER NOT NULL);
1029 CREATE INDEX IF NOT EXISTS reviews_scope ON daily_reviews(tenant_id, person_id, day);
1030 CREATE TABLE IF NOT EXISTS embeddings(tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, target_kind TEXT NOT NULL, target_id TEXT NOT NULL, model TEXT NOT NULL, version TEXT NOT NULL, dimension INTEGER NOT NULL, input_hash TEXT NOT NULL, target_revision INTEGER NOT NULL, created_at INTEGER NOT NULL, normalization TEXT NOT NULL, distance TEXT NOT NULL, vector TEXT NOT NULL, PRIMARY KEY(tenant_id, person_id, target_kind, target_id, model, version));
1031 CREATE INDEX IF NOT EXISTS embeddings_scope ON embeddings(tenant_id, person_id, target_kind, target_id);
1032 PRAGMA user_version = 1;
1033 COMMIT;",
1034 )?;
1035 for (column, definition) in [
1036 ("target_revision", "INTEGER NOT NULL DEFAULT 0"),
1037 ("created_at", "INTEGER NOT NULL DEFAULT 0"),
1038 ] {
1039 let exists: bool = self.connection.query_row(
1040 "SELECT EXISTS(SELECT 1 FROM pragma_table_info('embeddings') WHERE name = ?1)",
1041 [column],
1042 |row| row.get(0),
1043 )?;
1044 if !exists {
1045 self.connection.execute(
1046 &format!("ALTER TABLE embeddings ADD COLUMN {column} {definition}"),
1047 [],
1048 )?;
1049 }
1050 }
1051 let has_ingestion_key: bool = self.connection.query_row(
1052 "SELECT EXISTS(SELECT 1 FROM pragma_table_info('sources') WHERE name = 'ingestion_key')",
1053 [],
1054 |row| row.get(0),
1055 )?;
1056 if !has_ingestion_key {
1057 self.connection
1058 .execute("ALTER TABLE sources ADD COLUMN ingestion_key TEXT", [])?;
1059 }
1060 self.connection.execute(
1061 "CREATE UNIQUE INDEX IF NOT EXISTS sources_ingestion_key ON sources(tenant_id, person_id, ingestion_key) WHERE ingestion_key IS NOT NULL",
1062 [],
1063 )?;
1064 Ok(())
1065 }
1066}
1067
1068fn insert_claim(
1069 transaction: &Transaction<'_>,
1070 tenant_id: &TenantId,
1071 person_id: &PersonId,
1072 evidence_id: &EvidenceId,
1073 claim: ClaimInput,
1074 recorded_at: Timestamp,
1075) -> Result<ClaimId> {
1076 require_text("claim subject", &claim.subject)?;
1077 require_text("claim predicate", &claim.predicate)?;
1078 require_text("claim value", &claim.value)?;
1079 let id = ClaimId(new_id(transaction)?);
1080 transaction.execute(
1081 "INSERT INTO claims(id, tenant_id, person_id, subject, predicate, value, valid_from, recorded_from, status) VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 'accepted')",
1082 params![id.0, tenant_id.0, person_id.0, claim.subject, claim.predicate, claim.value, claim.valid_from, recorded_at],
1083 )?;
1084 let relation = serde_json::to_string(&EvidenceRelation::Supports)?;
1085 transaction.execute(
1086 "INSERT INTO claim_evidence(tenant_id, person_id, claim_id, evidence_id, relation, confidence_basis_points) VALUES(?1, ?2, ?3, ?4, ?5, 10000)",
1087 params![tenant_id.0, person_id.0, id.0, evidence_id.0, relation],
1088 )?;
1089 Ok(id)
1090}
1091
1092fn validate_transcript_locator(locator: &TranscriptLocator) -> Result<()> {
1093 for (field, value) in [
1094 ("locator.device_id", locator.device_id.as_str()),
1095 ("locator.provider", locator.provider.as_str()),
1096 ("locator.stream_id", locator.stream_id.as_str()),
1097 ("locator.segment_id", locator.segment_id.as_str()),
1098 ] {
1099 require_text(field, value)?;
1100 }
1101 if locator.start_ms >= locator.end_ms {
1102 return Err(Error::Invalid(
1103 "locator end_ms must be greater than start_ms".to_owned(),
1104 ));
1105 }
1106 if locator.end_ms > i64::MAX as u64 {
1107 return Err(Error::Invalid(
1108 "locator end_ms exceeds the storage range".to_owned(),
1109 ));
1110 }
1111 Ok(())
1112}
1113
1114fn new_id(transaction: &Transaction<'_>) -> Result<String> {
1115 Ok(transaction.query_row("SELECT lower(hex(randomblob(16)))", [], |row| row.get(0))?)
1116}
1117
1118fn require_scope(tenant_id: &TenantId, person_id: &PersonId) -> Result<()> {
1119 require_text("tenant_id", &tenant_id.0)?;
1120 require_text("person_id", &person_id.0)
1121}
1122
1123fn require_text(field: &str, value: &str) -> Result<()> {
1124 if value.trim().is_empty() {
1125 return Err(Error::Invalid(format!("{field} must not be empty")));
1126 }
1127 Ok(())
1128}
1129
1130fn validate_embedding(embedding: &Embedding) -> Result<()> {
1131 require_text("embedding model", &embedding.model)?;
1132 require_text("embedding version", &embedding.version)?;
1133 require_text("embedding input_hash", &embedding.input_hash)?;
1134 if embedding.vector.is_empty() || embedding.vector.iter().any(|value| !value.is_finite()) {
1135 return Err(Error::Invalid(
1136 "embedding vector must contain finite values".to_owned(),
1137 ));
1138 }
1139 Ok(())
1140}
1141
1142fn input_hash(text: &str) -> String {
1143 format!("sha256:{:x}", Sha256::digest(text.as_bytes()))
1144}
1145
1146fn validate_dense_query(query: &DenseQuery) -> Result<()> {
1147 require_text("query embedding model", &query.model)?;
1148 require_text("query embedding version", &query.version)?;
1149 if query.vector.is_empty() || query.vector.iter().any(|value| !value.is_finite()) {
1150 return Err(Error::Invalid(
1151 "query embedding vector must contain finite values".to_owned(),
1152 ));
1153 }
1154 Ok(())
1155}
1156
1157fn vector_score(query: &[f32], stored: &[f32], distance: &VectorDistance) -> Result<f32> {
1158 let dot = query
1159 .iter()
1160 .zip(stored)
1161 .map(|(left, right)| left * right)
1162 .sum::<f32>();
1163 let score = match distance {
1164 VectorDistance::Dot => dot,
1165 VectorDistance::Euclidean => -query
1166 .iter()
1167 .zip(stored)
1168 .map(|(left, right)| (left - right).powi(2))
1169 .sum::<f32>()
1170 .sqrt(),
1171 VectorDistance::Cosine => {
1172 let query_norm = query.iter().map(|value| value * value).sum::<f32>().sqrt();
1173 let stored_norm = stored.iter().map(|value| value * value).sum::<f32>().sqrt();
1174 if query_norm == 0.0 || stored_norm == 0.0 {
1175 return Err(Error::Invalid(
1176 "cosine embeddings must have non-zero magnitude".to_owned(),
1177 ));
1178 }
1179 dot / (query_norm * stored_norm)
1180 }
1181 };
1182 if !score.is_finite() {
1183 return Err(Error::Invalid(
1184 "embedding similarity is not finite".to_owned(),
1185 ));
1186 }
1187 Ok(score)
1188}
1189
1190fn reciprocal_rank_fusion(
1191 lexical: &[RetrievalTarget],
1192 dense: &[RetrievalTarget],
1193 limit: usize,
1194) -> Vec<(RetrievalTarget, u16)> {
1195 let mut scores = HashMap::<RetrievalTarget, u32>::new();
1196 for ranking in [lexical, dense] {
1197 let mut seen = HashSet::new();
1198 for (offset, id) in ranking.iter().enumerate() {
1199 if seen.insert(id) {
1200 *scores.entry(id.clone()).or_default() += 1_000_000 / (61 + offset as u32);
1201 }
1202 }
1203 }
1204 let maximum = scores.values().copied().max().unwrap_or(1);
1205 let mut ranked = scores
1206 .into_iter()
1207 .map(|(id, score)| (id, ((score as u64 * 10_000) / maximum as u64) as u16))
1208 .collect::<Vec<_>>();
1209 ranked.sort_by(|left, right| right.1.cmp(&left.1).then_with(|| left.0.cmp(&right.0)));
1210 ranked.truncate(limit);
1211 ranked
1212}
1213
1214fn embedding_target_parts(target: &EmbeddingTarget) -> (&'static str, &str) {
1215 match target {
1216 EmbeddingTarget::Source(id) => ("source", &id.0),
1217 EmbeddingTarget::Evidence(id) => ("evidence", &id.0),
1218 EmbeddingTarget::Claim(id) => ("claim", &id.0),
1219 }
1220}
1221
1222fn embedding_target(kind: &str, id: &str) -> Result<EmbeddingTarget> {
1223 Ok(match kind {
1224 "source" => EmbeddingTarget::Source(SourceId(id.to_owned())),
1225 "evidence" => EmbeddingTarget::Evidence(EvidenceId(id.to_owned())),
1226 "claim" => EmbeddingTarget::Claim(ClaimId(id.to_owned())),
1227 _ => {
1228 return Err(Error::Invalid(
1229 "stored embedding target is invalid".to_owned(),
1230 ));
1231 }
1232 })
1233}
1234
1235const fn default_limit() -> u32 {
1236 10
1237}
1238const fn bounded_limit(limit: u32) -> u32 {
1239 if limit == 0 {
1240 10
1241 } else if limit > 100 {
1242 100
1243 } else {
1244 limit
1245 }
1246}
1247
1248#[cfg(test)]
1249mod tests {
1250 use super::*;
1251
1252 fn remember(tenant: &str, person: &str, value: &str) -> RememberInput {
1253 RememberInput {
1254 tenant_id: TenantId(tenant.into()),
1255 person_id: PersonId(person.into()),
1256 ingestion_key: None,
1257 kind: SourceKind::Conversation,
1258 text: format!("Sam works at {value}"),
1259 captured_at: 10,
1260 claim: Some(ClaimInput {
1261 subject: "Sam".into(),
1262 predicate: "employer".into(),
1263 value: value.into(),
1264 valid_from: 10,
1265 }),
1266 }
1267 }
1268
1269 fn remember_raw(tenant: &str, person: &str, text: &str) -> RememberInput {
1270 RememberInput {
1271 tenant_id: TenantId(tenant.into()),
1272 person_id: PersonId(person.into()),
1273 ingestion_key: None,
1274 kind: SourceKind::Conversation,
1275 text: text.into(),
1276 captured_at: 10,
1277 claim: None,
1278 }
1279 }
1280
1281 fn hash_for(db: &MemoryDb, target: EmbeddingTarget) -> String {
1282 db.projection_input(&TenantId("a".into()), &PersonId("sam".into()), target)
1283 .unwrap()
1284 .input_hash
1285 }
1286
1287 #[test]
1288 fn ingestion_keys_are_idempotent_within_scope() {
1289 let mut db = MemoryDb {
1290 connection: Connection::open_in_memory().unwrap(),
1291 };
1292 db.migrate().unwrap();
1293 let mut first = remember_raw("a", "sam", "Remember once");
1294 first.ingestion_key = Some("turn-1".into());
1295 let stored = db.remember(first).unwrap();
1296 let mut replay = remember_raw("a", "sam", "Remember once");
1297 replay.ingestion_key = Some("turn-1".into());
1298 let replayed = db.remember(replay).unwrap();
1299 assert_eq!(stored.source_id, replayed.source_id);
1300 assert_eq!(stored.evidence_id, replayed.evidence_id);
1301 assert_eq!(
1302 db.connection
1303 .query_row("SELECT count(*) FROM sources", [], |row| row
1304 .get::<_, u64>(0))
1305 .unwrap(),
1306 1
1307 );
1308 let mut other_scope = remember_raw("b", "sam", "Remember once");
1309 other_scope.ingestion_key = Some("turn-1".into());
1310 assert_ne!(
1311 stored.source_id,
1312 db.remember(other_scope).unwrap().source_id
1313 );
1314 }
1315
1316 #[test]
1317 fn ingestion_key_rejects_changed_payload() {
1318 let mut db = MemoryDb {
1319 connection: Connection::open_in_memory().unwrap(),
1320 };
1321 db.migrate().unwrap();
1322 let mut first = remember_raw("a", "sam", "Remember once");
1323 first.ingestion_key = Some("turn-1".into());
1324 db.remember(first).unwrap();
1325 let changed_text = remember_raw("a", "sam", "Changed replay content");
1326 let mut changed_kind = remember_raw("a", "sam", "Remember once");
1327 changed_kind.kind = SourceKind::Screen;
1328 let mut changed_time = remember_raw("a", "sam", "Remember once");
1329 changed_time.captured_at = 11;
1330 let mut changed_claim = remember_raw("a", "sam", "Remember once");
1331 changed_claim.claim = Some(ClaimInput {
1332 subject: "Sam".into(),
1333 predicate: "status".into(),
1334 value: "focused".into(),
1335 valid_from: 10,
1336 });
1337 for mut changed in [changed_text, changed_kind, changed_time, changed_claim] {
1338 changed.ingestion_key = Some("turn-1".into());
1339 assert!(matches!(db.remember(changed), Err(Error::Invalid(_))));
1340 }
1341 let with_claim = |value: &str| {
1342 let mut input = remember_raw("a", "sam", "Claim capture");
1343 input.ingestion_key = Some("claim-turn".into());
1344 input.claim = Some(ClaimInput {
1345 subject: "Sam".into(),
1346 predicate: "status".into(),
1347 value: value.into(),
1348 valid_from: 10,
1349 });
1350 input
1351 };
1352 let stored = db.remember(with_claim("focused")).unwrap();
1353 assert_eq!(
1354 stored.claim_id,
1355 db.remember(with_claim("focused")).unwrap().claim_id
1356 );
1357 assert!(matches!(
1358 db.remember(with_claim("distracted")),
1359 Err(Error::Invalid(_))
1360 ));
1361 }
1362
1363 #[test]
1364 fn ingestion_key_rejects_deleted_memory() {
1365 let mut db = MemoryDb {
1366 connection: Connection::open_in_memory().unwrap(),
1367 };
1368 db.migrate().unwrap();
1369 let mut first = remember_raw("a", "sam", "Remember once");
1370 first.ingestion_key = Some("turn-1".into());
1371 let stored = db.remember(first).unwrap();
1372 db.delete_source(DeleteInput {
1373 tenant_id: TenantId("a".into()),
1374 person_id: PersonId("sam".into()),
1375 source_id: stored.source_id,
1376 deleted_at: 20,
1377 })
1378 .unwrap();
1379 let mut replay = remember_raw("a", "sam", "Remember once");
1380 replay.ingestion_key = Some("turn-1".into());
1381 assert!(matches!(db.remember(replay), Err(Error::NotFound)));
1382 }
1383
1384 #[test]
1385 fn raw_sources_are_scoped_cited_and_deleted_from_retrieval() {
1386 let mut db = MemoryDb {
1387 connection: Connection::open_in_memory().unwrap(),
1388 };
1389 db.migrate().unwrap();
1390 let raw = db
1391 .remember(remember_raw("a", "sam", "The launch code is marigold"))
1392 .unwrap();
1393 db.remember(remember_raw("b", "sam", "Marigold belongs elsewhere"))
1394 .unwrap();
1395
1396 let found = db
1397 .search(SearchInput {
1398 tenant_id: TenantId("a".into()),
1399 person_id: PersonId("sam".into()),
1400 query: "marigold".into(),
1401 limit: 5,
1402 query_embedding: None,
1403 })
1404 .unwrap();
1405 assert_eq!(found.items.len(), 1);
1406 assert_eq!(
1407 found.items[0].memory,
1408 MemoryRef::Source(raw.source_id.clone())
1409 );
1410 assert_eq!(found.items[0].evidence_ids, vec![raw.evidence_id]);
1411
1412 db.delete_source(DeleteInput {
1413 tenant_id: TenantId("a".into()),
1414 person_id: PersonId("sam".into()),
1415 source_id: raw.source_id.clone(),
1416 deleted_at: 20,
1417 })
1418 .unwrap();
1419 assert!(matches!(
1420 db.get(GetInput {
1421 tenant_id: TenantId("a".into()),
1422 person_id: PersonId("sam".into()),
1423 target: EmbeddingTarget::Source(raw.source_id),
1424 }),
1425 Err(Error::NotFound)
1426 ));
1427 assert!(
1428 db.search(SearchInput {
1429 tenant_id: TenantId("a".into()),
1430 person_id: PersonId("sam".into()),
1431 query: "marigold".into(),
1432 limit: 5,
1433 query_embedding: None,
1434 })
1435 .unwrap()
1436 .items
1437 .is_empty()
1438 );
1439 }
1440
1441 #[test]
1442 fn accepted_claim_replaces_its_source_in_retrieval() {
1443 let mut db = MemoryDb {
1444 connection: Connection::open_in_memory().unwrap(),
1445 };
1446 db.migrate().unwrap();
1447 let remembered = db.remember(remember("a", "sam", "Acme")).unwrap();
1448 let target = EmbeddingTarget::Source(remembered.source_id);
1449 let input_hash = hash_for(&db, target.clone());
1450 db.upsert_embedding(EmbeddingInput {
1451 tenant_id: TenantId("a".into()),
1452 person_id: PersonId("sam".into()),
1453 target,
1454 embedding: Embedding {
1455 vector: vec![1.0, 0.0],
1456 model: "test/model".into(),
1457 version: "1".into(),
1458 input_hash,
1459 normalization: VectorNormalization::L2,
1460 distance: VectorDistance::Cosine,
1461 },
1462 })
1463 .unwrap();
1464 let found = db
1465 .search(SearchInput {
1466 tenant_id: TenantId("a".into()),
1467 person_id: PersonId("sam".into()),
1468 query: "Acme".into(),
1469 limit: 5,
1470 query_embedding: Some(DenseQuery {
1471 vector: vec![1.0, 0.0],
1472 model: "test/model".into(),
1473 version: "1".into(),
1474 }),
1475 })
1476 .unwrap();
1477 assert_eq!(found.items.len(), 1);
1478 assert_eq!(
1479 found.items[0].memory,
1480 MemoryRef::Claim(remembered.claim_id.unwrap())
1481 );
1482 }
1483
1484 #[test]
1485 fn dense_evidence_without_a_claim_is_retrievable() {
1486 let mut db = MemoryDb {
1487 connection: Connection::open_in_memory().unwrap(),
1488 };
1489 db.migrate().unwrap();
1490 let raw = db
1491 .remember(remember_raw("a", "sam", "Quiet desk near a window"))
1492 .unwrap();
1493 let target = EmbeddingTarget::Evidence(raw.evidence_id.clone());
1494 let input_hash = hash_for(&db, target.clone());
1495 db.upsert_embedding(EmbeddingInput {
1496 tenant_id: TenantId("a".into()),
1497 person_id: PersonId("sam".into()),
1498 target,
1499 embedding: Embedding {
1500 vector: vec![1.0, 0.0],
1501 model: "test/model".into(),
1502 version: "1".into(),
1503 input_hash,
1504 normalization: VectorNormalization::L2,
1505 distance: VectorDistance::Cosine,
1506 },
1507 })
1508 .unwrap();
1509 let found = db
1510 .search(SearchInput {
1511 tenant_id: TenantId("a".into()),
1512 person_id: PersonId("sam".into()),
1513 query: "unmatched lexical phrase".into(),
1514 limit: 5,
1515 query_embedding: Some(DenseQuery {
1516 vector: vec![1.0, 0.0],
1517 model: "test/model".into(),
1518 version: "1".into(),
1519 }),
1520 })
1521 .unwrap();
1522 assert_eq!(found.items.len(), 1);
1523 assert_eq!(found.items[0].memory, MemoryRef::Evidence(raw.evidence_id));
1524
1525 db.delete_source(DeleteInput {
1526 tenant_id: TenantId("a".into()),
1527 person_id: PersonId("sam".into()),
1528 source_id: raw.source_id,
1529 deleted_at: 20,
1530 })
1531 .unwrap();
1532 assert!(
1533 db.search(SearchInput {
1534 tenant_id: TenantId("a".into()),
1535 person_id: PersonId("sam".into()),
1536 query: "unmatched lexical phrase".into(),
1537 limit: 5,
1538 query_embedding: Some(DenseQuery {
1539 vector: vec![1.0, 0.0],
1540 model: "test/model".into(),
1541 version: "1".into(),
1542 }),
1543 })
1544 .unwrap()
1545 .items
1546 .is_empty()
1547 );
1548 }
1549
1550 #[test]
1551 fn lifecycle_is_scoped_cited_and_propagates_deletion() {
1552 let mut db = MemoryDb {
1553 connection: Connection::open_in_memory().unwrap(),
1554 };
1555 db.migrate().unwrap();
1556 let first = db.remember(remember("a", "sam", "Acme")).unwrap();
1557 db.remember(remember("b", "sam", "Other")).unwrap();
1558 let found = db
1559 .search(SearchInput {
1560 tenant_id: TenantId("a".into()),
1561 person_id: PersonId("sam".into()),
1562 query: "Acme".into(),
1563 limit: 5,
1564 query_embedding: None,
1565 })
1566 .unwrap();
1567 assert_eq!(found.items.len(), 1);
1568 assert_eq!(found.items[0].evidence_ids, vec![first.evidence_id.clone()]);
1569 let corrected = db
1570 .correct(CorrectInput {
1571 tenant_id: TenantId("a".into()),
1572 person_id: PersonId("sam".into()),
1573 claim_id: first.claim_id.unwrap(),
1574 text: "I moved to Beta".into(),
1575 value: "Beta".into(),
1576 occurred_at: 20,
1577 })
1578 .unwrap();
1579 let deleted = db
1580 .delete_source(DeleteInput {
1581 tenant_id: TenantId("a".into()),
1582 person_id: PersonId("sam".into()),
1583 source_id: corrected.source_id,
1584 deleted_at: 30,
1585 })
1586 .unwrap();
1587 assert_eq!((deleted.evidence_count, deleted.claim_count), (1, 1));
1588 assert!(
1589 db.search(SearchInput {
1590 tenant_id: TenantId("a".into()),
1591 person_id: PersonId("sam".into()),
1592 query: "Beta".into(),
1593 limit: 5,
1594 query_embedding: None,
1595 })
1596 .unwrap()
1597 .items
1598 .is_empty()
1599 );
1600 }
1601
1602 #[test]
1603 fn reviews_require_live_same_scope_evidence() {
1604 let mut db = MemoryDb {
1605 connection: Connection::open_in_memory().unwrap(),
1606 };
1607 db.migrate().unwrap();
1608 let remembered = db.remember(remember("a", "sam", "Acme")).unwrap();
1609 db.store_review(ReviewInput {
1610 tenant_id: TenantId("a".into()),
1611 person_id: PersonId("sam".into()),
1612 day: "2026-07-21".into(),
1613 summary: "Worked at Acme".into(),
1614 evidence_ids: vec![remembered.evidence_id.clone()],
1615 recorded_at: 20,
1616 })
1617 .unwrap();
1618 assert_eq!(
1619 db.reviews(ReviewsInput {
1620 tenant_id: TenantId("a".into()),
1621 person_id: PersonId("sam".into()),
1622 limit: 1
1623 })
1624 .unwrap()
1625 .len(),
1626 1
1627 );
1628 db.delete_source(DeleteInput {
1629 tenant_id: TenantId("a".into()),
1630 person_id: PersonId("sam".into()),
1631 source_id: remembered.source_id,
1632 deleted_at: 30,
1633 })
1634 .unwrap();
1635 assert!(
1636 db.reviews(ReviewsInput {
1637 tenant_id: TenantId("a".into()),
1638 person_id: PersonId("sam".into()),
1639 limit: 1
1640 })
1641 .unwrap()
1642 .is_empty()
1643 );
1644 }
1645
1646 #[test]
1647 fn embedding_projection_is_validated_and_scoped() {
1648 let mut db = MemoryDb {
1649 connection: Connection::open_in_memory().unwrap(),
1650 };
1651 db.migrate().unwrap();
1652 let remembered = db.remember(remember("a", "sam", "Acme")).unwrap();
1653 let target = EmbeddingTarget::Evidence(remembered.evidence_id.clone());
1654 let embedding = Embedding {
1655 vector: vec![0.1, 0.2],
1656 model: "provider/model".into(),
1657 version: "1".into(),
1658 input_hash: hash_for(&db, target.clone()),
1659 normalization: VectorNormalization::L2,
1660 distance: VectorDistance::Cosine,
1661 };
1662 let stored = db
1663 .upsert_embedding(EmbeddingInput {
1664 tenant_id: TenantId("a".into()),
1665 person_id: PersonId("sam".into()),
1666 target,
1667 embedding: embedding.clone(),
1668 })
1669 .unwrap();
1670 assert_eq!(stored.dimension, 2);
1671 assert!(matches!(
1672 db.upsert_embedding(EmbeddingInput {
1673 tenant_id: TenantId("b".into()),
1674 person_id: PersonId("sam".into()),
1675 target: EmbeddingTarget::Evidence(remembered.evidence_id),
1676 embedding,
1677 }),
1678 Err(Error::NotFound)
1679 ));
1680 }
1681
1682 #[test]
1683 fn search_fuses_lexical_and_real_dense_ranks_deterministically() {
1684 let mut db = MemoryDb {
1685 connection: Connection::open_in_memory().unwrap(),
1686 };
1687 db.migrate().unwrap();
1688 let lexical = db.remember(remember("a", "sam", "Jazz Club")).unwrap();
1689 let dense = db.remember(remember("a", "sam", "Music Venue")).unwrap();
1690 for (claim_id, vector) in [
1691 (lexical.claim_id.clone().unwrap(), vec![0.9, 0.1]),
1692 (dense.claim_id.unwrap(), vec![1.0, 0.0]),
1693 ] {
1694 let target = EmbeddingTarget::Claim(claim_id);
1695 let input_hash = hash_for(&db, target.clone());
1696 db.upsert_embedding(EmbeddingInput {
1697 tenant_id: TenantId("a".into()),
1698 person_id: PersonId("sam".into()),
1699 target,
1700 embedding: Embedding {
1701 vector,
1702 model: "test/model".into(),
1703 version: "1".into(),
1704 input_hash,
1705 normalization: VectorNormalization::L2,
1706 distance: VectorDistance::Cosine,
1707 },
1708 })
1709 .unwrap();
1710 }
1711
1712 let found = db
1713 .search(SearchInput {
1714 tenant_id: TenantId("a".into()),
1715 person_id: PersonId("sam".into()),
1716 query: "Jazz Club".into(),
1717 limit: 2,
1718 query_embedding: Some(DenseQuery {
1719 vector: vec![1.0, 0.0],
1720 model: "test/model".into(),
1721 version: "1".into(),
1722 }),
1723 })
1724 .unwrap();
1725
1726 assert_eq!(
1727 found.items[0].memory,
1728 MemoryRef::Claim(lexical.claim_id.unwrap())
1729 );
1730 assert_eq!(found.items.len(), 2);
1731 assert!(!found.items[0].evidence_ids.is_empty());
1732 }
1733
1734 #[test]
1735 fn stale_projections_are_excluded_and_reported_with_current_inputs() {
1736 let mut db = MemoryDb {
1737 connection: Connection::open_in_memory().unwrap(),
1738 };
1739 db.migrate().unwrap();
1740 let raw = db
1741 .remember(remember_raw("a", "sam", "A quiet desk"))
1742 .unwrap();
1743 let claimed = db.remember(remember("a", "sam", "Acme")).unwrap();
1744 let other = db.remember(remember("b", "sam", "Other")).unwrap();
1745 let targets = [
1746 EmbeddingTarget::Source(raw.source_id.clone()),
1747 EmbeddingTarget::Evidence(raw.evidence_id.clone()),
1748 EmbeddingTarget::Claim(claimed.claim_id.clone().unwrap()),
1749 ];
1750 for target in &targets {
1751 db.upsert_embedding(EmbeddingInput {
1752 tenant_id: TenantId("a".into()),
1753 person_id: PersonId("sam".into()),
1754 target: target.clone(),
1755 embedding: Embedding {
1756 vector: vec![1.0, 0.0],
1757 model: "test/model".into(),
1758 version: "1".into(),
1759 input_hash: hash_for(&db, target.clone()),
1760 normalization: VectorNormalization::L2,
1761 distance: VectorDistance::Cosine,
1762 },
1763 })
1764 .unwrap();
1765 }
1766 db.connection
1767 .execute(
1768 "UPDATE sources SET content = 'A changed desk', revision = revision + 1 WHERE id = ?1",
1769 [&raw.source_id.0],
1770 )
1771 .unwrap();
1772 db.connection
1773 .execute(
1774 "UPDATE evidence SET quote = 'Changed evidence' WHERE id = ?1",
1775 [&raw.evidence_id.0],
1776 )
1777 .unwrap();
1778 db.connection
1779 .execute(
1780 "UPDATE claims SET value = 'Changed employer' WHERE id = ?1",
1781 [&claimed.claim_id.as_ref().unwrap().0],
1782 )
1783 .unwrap();
1784
1785 let found = db
1786 .search(SearchInput {
1787 tenant_id: TenantId("a".into()),
1788 person_id: PersonId("sam".into()),
1789 query: "no lexical match".into(),
1790 limit: 10,
1791 query_embedding: Some(DenseQuery {
1792 vector: vec![1.0, 0.0],
1793 model: "test/model".into(),
1794 version: "1".into(),
1795 }),
1796 })
1797 .unwrap();
1798 assert!(found.items.is_empty());
1799
1800 let issues = db
1801 .projection_issues(ProjectionAuditInput {
1802 tenant_id: TenantId("a".into()),
1803 person_id: PersonId("sam".into()),
1804 model: "test/model".into(),
1805 version: "1".into(),
1806 limit: 100,
1807 })
1808 .unwrap();
1809 let stale = issues
1810 .iter()
1811 .filter(|issue| issue.state == ProjectionState::Stale)
1812 .map(|issue| &issue.input.target)
1813 .collect::<HashSet<_>>();
1814 assert_eq!(stale, targets.iter().collect::<HashSet<_>>());
1815 assert!(issues.iter().all(|issue| {
1816 issue.input.target != EmbeddingTarget::Source(other.source_id.clone())
1817 && issue.input.target != EmbeddingTarget::Evidence(other.evidence_id.clone())
1818 && issue.input.target != EmbeddingTarget::Claim(other.claim_id.clone().unwrap())
1819 }));
1820 assert_eq!(
1821 db.projection_issues(ProjectionAuditInput {
1822 tenant_id: TenantId("a".into()),
1823 person_id: PersonId("sam".into()),
1824 model: "test/model".into(),
1825 version: "1".into(),
1826 limit: 2,
1827 })
1828 .unwrap()
1829 .len(),
1830 2
1831 );
1832 }
1833
1834 #[test]
1835 fn embedding_rejects_input_hash_for_different_text() {
1836 let mut db = MemoryDb {
1837 connection: Connection::open_in_memory().unwrap(),
1838 };
1839 db.migrate().unwrap();
1840 let remembered = db
1841 .remember(remember_raw("a", "sam", "Current text"))
1842 .unwrap();
1843 let result = db.upsert_embedding(EmbeddingInput {
1844 tenant_id: TenantId("a".into()),
1845 person_id: PersonId("sam".into()),
1846 target: EmbeddingTarget::Source(remembered.source_id),
1847 embedding: Embedding {
1848 vector: vec![1.0],
1849 model: "test/model".into(),
1850 version: "1".into(),
1851 input_hash: input_hash("different text"),
1852 normalization: VectorNormalization::L2,
1853 distance: VectorDistance::Cosine,
1854 },
1855 });
1856 assert!(matches!(result, Err(Error::Invalid(_))));
1857 }
1858
1859 #[test]
1860 fn migration_marks_existing_projections_for_lifecycle_revalidation() {
1861 let connection = Connection::open_in_memory().unwrap();
1862 connection
1863 .execute_batch(
1864 "CREATE TABLE embeddings(tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, target_kind TEXT NOT NULL, target_id TEXT NOT NULL, model TEXT NOT NULL, version TEXT NOT NULL, dimension INTEGER NOT NULL, input_hash TEXT NOT NULL, normalization TEXT NOT NULL, distance TEXT NOT NULL, vector TEXT NOT NULL, PRIMARY KEY(tenant_id, person_id, target_kind, target_id, model, version));
1865 INSERT INTO embeddings VALUES('a', 'sam', 'source', 'old', 'model', '1', 1, 'sha256:old', '\"l2\"', '\"cosine\"', '[1.0]');",
1866 )
1867 .unwrap();
1868 let mut db = MemoryDb { connection };
1869 db.migrate().unwrap();
1870 let lifecycle = db
1871 .connection
1872 .query_row(
1873 "SELECT target_revision, created_at FROM embeddings WHERE target_id = 'old'",
1874 [],
1875 |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
1876 )
1877 .unwrap();
1878 assert_eq!(lifecycle, (0, 0));
1879 }
1880}