1use std::sync::Arc;
28
29use async_trait::async_trait;
30use ijima_core::{
31 AcceptedExtraction, AuthorityScope, DiaryEntry, Embedding, Entity, EntityId, EntityRecord,
32 IjimaError, InstanceId, KgStats, KnowledgeGraph, Memory, MemoryId, NamespaceCount, NamespaceId,
33 NamespaceMembership, PalaceGraph, ProjectTaxon, QueuedExtraction, RepoDirectory, Result, Room,
34 SearchHit, Session, SessionId, SessionTurn, Store, StoreStats, TokenRevocation, Triple, Tunnel,
35 TunnelTraversal, embeddings::Embedder, harness::Harness, memory::MemorySource,
36};
37use serde::{Deserialize, Serialize};
38use surrealdb::Surreal;
39use surrealdb::engine::local::{Db, Mem, SurrealKv};
40
41const SURREAL_NS: &str = "ijima";
43const SURREAL_DB: &str = "core";
44
45const MEMORIES_TABLE: &str = "memories";
46const TURNS_TABLE: &str = "session_turns";
47const SESSIONS_TABLE: &str = "sessions";
48const DIARY_TABLE: &str = "diaries";
49const QUEUE_TABLE: &str = "mining_queue";
50const REVOCATIONS_TABLE: &str = "token_revocations";
51const NAMESPACE_MEMBERS_TABLE: &str = "namespace_members";
52const ENTITIES_TABLE: &str = "entities";
54const TRIPLES_TABLE: &str = "triples";
56const REPO_TABLE: &str = "repo_directory";
58
59pub struct SurrealStore {
61 db: Surreal<Db>,
62 embedder: Option<Arc<dyn Embedder>>,
66}
67
68impl SurrealStore {
69 pub async fn open_embedded() -> Result<Self> {
76 Self::open_with_db(new_mem().await?, None).await
77 }
78
79 pub async fn open_embedded_with(embedder: Arc<dyn Embedder>) -> Result<Self> {
85 Self::open_with_db(new_mem().await?, Some(embedder)).await
86 }
87
88 pub async fn open_persistent(path: impl AsRef<std::path::Path>) -> Result<Self> {
106 Self::open_with_db(new_surrealkv(&path).await?, None).await
107 }
108
109 pub async fn open_persistent_with(
116 path: impl AsRef<std::path::Path>,
117 embedder: Arc<dyn Embedder>,
118 ) -> Result<Self> {
119 Self::open_with_db(new_surrealkv(&path).await?, Some(embedder)).await
120 }
121
122 async fn open_with_db(db: Surreal<Db>, embedder: Option<Arc<dyn Embedder>>) -> Result<Self> {
123 db.use_ns(SURREAL_NS)
124 .use_db(SURREAL_DB)
125 .await
126 .map_err(|e| IjimaError::Store {
127 detail: format!("surrealdb use_ns/use_db: {e}"),
128 })?;
129 db.query(Self::INDEX_DDL)
135 .await
136 .map_err(|e| IjimaError::Store {
137 detail: format!("surrealdb define indexes: {e}"),
138 })?;
139 Ok(Self { db, embedder })
140 }
141
142 const INDEX_DDL: &str = r#"
153 DEFINE TABLE IF NOT EXISTS memories SCHEMALESS;
154 DEFINE TABLE IF NOT EXISTS session_turns SCHEMALESS;
155 DEFINE TABLE IF NOT EXISTS sessions SCHEMALESS;
156 DEFINE TABLE IF NOT EXISTS diaries SCHEMALESS;
157 DEFINE TABLE IF NOT EXISTS mining_queue SCHEMALESS;
158 DEFINE TABLE IF NOT EXISTS entities SCHEMALESS;
159 DEFINE TABLE IF NOT EXISTS triples SCHEMALESS;
160 DEFINE TABLE IF NOT EXISTS token_revocations SCHEMALESS;
161 DEFINE TABLE IF NOT EXISTS namespace_members SCHEMALESS;
162 DEFINE INDEX IF NOT EXISTS nmem_ns ON TABLE namespace_members FIELDS namespace;
163 DEFINE INDEX IF NOT EXISTS nmem_ns_princ ON TABLE namespace_members FIELDS namespace, principal;
164 DEFINE INDEX IF NOT EXISTS mem_ns ON TABLE memories FIELDS namespace;
165 DEFINE INDEX IF NOT EXISTS mem_ns_hash ON TABLE memories FIELDS namespace, content_hash;
166 DEFINE INDEX IF NOT EXISTS turns_ns_sess ON TABLE session_turns FIELDS namespace, session_id;
167 DEFINE INDEX IF NOT EXISTS sess_ns ON TABLE sessions FIELDS namespace;
168 DEFINE INDEX IF NOT EXISTS diary_ns_ag ON TABLE diaries FIELDS namespace, agent;
169 DEFINE INDEX IF NOT EXISTS queue_ns ON TABLE mining_queue FIELDS namespace;
170 DEFINE INDEX IF NOT EXISTS ent_ns ON TABLE entities FIELDS namespace;
171 DEFINE INDEX IF NOT EXISTS trip_ns ON TABLE triples FIELDS namespace;
172 DEFINE INDEX IF NOT EXISTS trip_subj ON TABLE triples FIELDS subject;
173 DEFINE INDEX IF NOT EXISTS trip_obj ON TABLE triples FIELDS object;
174 "#;
175
176 async fn project_topic_counts(
179 &self,
180 ns: &NamespaceId,
181 ) -> Result<std::collections::BTreeMap<(String, String), usize>> {
182 #[derive(serde::Serialize, Deserialize)]
183 struct Row {
184 project: String,
185 topic: String,
186 }
187 let mut result = self
188 .db
189 .query(format!(
190 "SELECT project, topic FROM {MEMORIES_TABLE} WHERE namespace = $ns"
191 ))
192 .bind(("ns", ns.as_str().to_string()))
193 .await
194 .map_err(store_err)?;
195 let rows = take_vec::<Row>(&mut result)?;
196 let mut counts = std::collections::BTreeMap::new();
197 for r in rows {
198 *counts.entry((r.project, r.topic)).or_insert(0) += 1;
199 }
200 Ok(counts)
201 }
202
203 async fn project_topic_memories(
206 &self,
207 ns: &NamespaceId,
208 project: &str,
209 topic: &str,
210 limit: usize,
211 ) -> Result<Vec<Memory>> {
212 let mut result = self
213 .db
214 .query(format!(
215 "SELECT memory_id, content, project, topic, source, harness, session_id, namespace,
216 origin, authority, importance, created_at
217 FROM {MEMORIES_TABLE}
218 WHERE namespace = $ns AND project = $proj AND topic = $topic
219 ORDER BY importance DESC, created_at DESC LIMIT $lim"
220 ))
221 .bind(("ns", ns.as_str().to_string()))
222 .bind(("proj", project.to_string()))
223 .bind(("topic", topic.to_string()))
224 .bind(("lim", limit as i64))
225 .await
226 .map_err(store_err)?;
227 let records = take_vec::<MemoryRecord>(&mut result)?;
228 Ok(records.into_iter().map(|r| r.into_memory()).collect())
229 }
230
231 pub async fn export_to(&self, path: impl AsRef<std::path::Path>) -> Result<()> {
235 self.db
236 .export(path.as_ref())
237 .await
238 .map_err(|e| IjimaError::Store {
239 detail: format!("export: {e}"),
240 })?;
241 Ok(())
242 }
243}
244
245async fn new_mem() -> Result<Surreal<Db>> {
246 Surreal::new::<Mem>(())
247 .await
248 .map_err(|e| IjimaError::Store {
249 detail: format!("surrealdb mem init: {e}"),
250 })
251}
252
253async fn new_surrealkv(path: impl AsRef<std::path::Path>) -> Result<Surreal<Db>> {
254 let path = path.as_ref();
255 if let Some(parent) = path.parent() {
256 std::fs::create_dir_all(parent).map_err(|e| IjimaError::Store {
257 detail: format!("mkdir {}: {e}", parent.display()),
258 })?;
259 }
260 Surreal::new::<SurrealKv>(path.to_string_lossy().to_string())
261 .await
262 .map_err(|e| IjimaError::Store {
263 detail: format!("surrealdb surrealkv init at {}: {e}", path.display()),
264 })
265}
266
267fn take_vec<T>(result: &mut surrealdb::IndexedResults) -> Result<Vec<T>>
277where
278 T: serde::Serialize + serde::de::DeserializeOwned + 'static,
279{
280 let wrapped: Vec<surrealdb::types::SerdeWrapper<T>> = result.take(0).map_err(store_err)?;
281 Ok(wrapped.into_iter().map(|w| w.0).collect())
282}
283
284#[derive(Debug, Clone, Serialize, Deserialize)]
289struct MemoryRecord {
290 memory_id: String,
293 content: String,
294 #[serde(default)]
298 content_hash: String,
299 project: String,
300 topic: String,
301 source: MemorySource,
302 harness: Harness,
303 session_id: Option<String>,
304 namespace: String,
305 #[serde(default = "default_record_importance")]
306 importance: f32,
307 #[serde(default)]
308 created_at: String,
309 #[serde(default)]
312 origin: InstanceId,
313 #[serde(default)]
316 authority: AuthorityScope,
317 #[serde(default, skip_serializing_if = "Option::is_none")]
321 embed_model: Option<String>,
322 #[serde(default, skip_serializing_if = "Option::is_none")]
326 embedding: Option<Vec<f32>>,
327}
328
329fn default_record_importance() -> f32 {
330 0.5
331}
332
333impl MemoryRecord {
334 fn from_memory(
335 memory: &Memory,
336 ns: &NamespaceId,
337 embedding: Option<Vec<f32>>,
338 embed_model: Option<String>,
339 ) -> Self {
340 use sha2::{Digest, Sha256};
341 let content_hash = hex(&Sha256::digest(memory.content.as_bytes()));
342 Self {
343 memory_id: memory.id.0.clone(),
344 content: memory.content.clone(),
345 content_hash,
346 project: memory.project.clone(),
347 topic: memory.topic.clone(),
348 source: memory.source,
349 harness: memory.harness,
350 session_id: memory.session_id.clone(),
351 namespace: ns.as_str().to_string(),
352 importance: memory.importance,
353 created_at: memory.created_at.clone(),
354 origin: memory.origin.clone(),
355 authority: memory.authority.clone(),
356 embed_model,
357 embedding,
358 }
359 }
360
361 fn into_memory(self) -> Memory {
362 Memory {
363 id: MemoryId(self.memory_id),
364 content: self.content,
365 project: self.project,
366 topic: self.topic,
367 source: self.source,
368 harness: self.harness,
369 session_id: self.session_id,
370 origin: self.origin,
371 authority: self.authority,
372 importance: self.importance,
373 created_at: self.created_at,
374 }
375 }
376}
377
378#[derive(Debug, Clone, Serialize, Deserialize)]
380struct SessionTurnRecord {
381 session_id: String,
382 turn_index: u64,
383 role: ijima_core::TurnRole,
384 content: String,
385 timestamp: String,
386 namespace: String,
387}
388
389#[derive(Debug, Clone, Serialize, Deserialize)]
393struct DiaryRecord {
394 agent: String,
395 content: String,
396 #[serde(default, skip_serializing_if = "Option::is_none")]
397 topic: Option<String>,
398 timestamp: String,
399 ts: i64,
400 namespace: String,
401}
402
403impl DiaryRecord {
404 fn from_entry(entry: &ijima_core::DiaryEntry, ns: &NamespaceId, now_ms: i64) -> Self {
405 Self {
406 agent: entry.agent.clone(),
407 content: entry.content.clone(),
408 topic: entry.topic.clone(),
409 timestamp: entry.timestamp.clone(),
410 ts: now_ms,
411 namespace: ns.as_str().to_string(),
412 }
413 }
414
415 fn into_entry(self) -> ijima_core::DiaryEntry {
416 ijima_core::DiaryEntry {
417 agent: self.agent,
418 content: self.content,
419 topic: self.topic,
420 timestamp: self.timestamp,
421 }
422 }
423}
424
425#[derive(Debug, Clone, Serialize, Deserialize)]
427struct SessionRecord {
428 session_id: String,
429 harness: String,
430 #[serde(default, skip_serializing_if = "Option::is_none")]
431 channel: Option<String>,
432 started_at: String,
433 #[serde(default, skip_serializing_if = "Option::is_none")]
434 ended_at: Option<String>,
435 namespace: String,
436}
437
438impl SessionRecord {
439 fn from_session(session: &Session, ns: &NamespaceId) -> Self {
440 Self {
441 session_id: session.id.0.clone(),
442 harness: session.harness.as_wire_str().to_string(),
443 channel: session.channel.clone(),
444 started_at: session.started_at.clone(),
445 ended_at: session.ended_at.clone(),
446 namespace: ns.as_str().to_string(),
447 }
448 }
449
450 fn into_session(self) -> Session {
451 Session {
452 id: SessionId(self.session_id),
453 harness: Harness::from_wire_str(&self.harness),
454 channel: self.channel,
455 started_at: self.started_at,
456 ended_at: self.ended_at,
457 }
458 }
459}
460
461#[derive(Debug, Clone, Serialize, Deserialize)]
464struct QueueRecord {
465 queue_id: String,
466 memory_id: String,
468 content: String,
469 project: String,
470 topic: String,
471 #[serde(default)]
472 importance: f32,
473 session_id: Option<String>,
474 harness: String,
475 confidence: f32,
476 source_session_id: String,
477 ts: i64,
478 namespace: String,
479}
480
481impl QueueRecord {
482 fn from_extraction(
483 queue_id: &str,
484 memory: &Memory,
485 confidence: f32,
486 ns: &NamespaceId,
487 now_ms: i64,
488 ) -> Self {
489 Self {
490 queue_id: queue_id.to_string(),
491 memory_id: memory.id.0.clone(),
492 content: memory.content.clone(),
493 project: memory.project.clone(),
494 topic: memory.topic.clone(),
495 importance: memory.importance,
496 session_id: memory.session_id.clone(),
497 harness: memory.harness.as_wire_str().to_string(),
498 confidence,
499 source_session_id: memory.session_id.clone().unwrap_or_default(),
500 ts: now_ms,
501 namespace: ns.as_str().to_string(),
502 }
503 }
504
505 fn into_queued(self) -> QueuedExtraction {
506 let harness = Harness::from_wire_str(&self.harness);
507 let memory = Memory {
508 id: MemoryId(self.memory_id),
509 content: self.content,
510 project: self.project,
511 topic: self.topic,
512 source: MemorySource::Mined,
513 harness,
514 session_id: self.session_id,
515 origin: InstanceId::local(),
516 authority: AuthorityScope::local(),
517 importance: self.importance,
518 created_at: String::new(),
519 };
520 QueuedExtraction {
521 id: self.queue_id,
522 memory,
523 confidence: self.confidence,
524 source_session_id: self.source_session_id,
525 queued_at: self.ts.to_string(),
526 }
527 }
528
529 fn into_memory(self) -> Memory {
530 Memory {
531 id: MemoryId(self.memory_id),
532 content: self.content,
533 project: self.project,
534 topic: self.topic,
535 source: MemorySource::Mined,
536 harness: Harness::from_wire_str(&self.harness),
537 session_id: self.session_id,
538 origin: InstanceId::local(),
539 authority: AuthorityScope::local(),
540 importance: self.importance,
541 created_at: String::new(),
542 }
543 }
544}
545
546fn store_err(e: surrealdb::Error) -> IjimaError {
547 IjimaError::Store {
548 detail: e.to_string(),
549 }
550}
551
552fn now_millis() -> i64 {
554 std::time::SystemTime::now()
555 .duration_since(std::time::UNIX_EPOCH)
556 .map(|d| d.as_millis() as i64)
557 .unwrap_or(0)
558}
559
560fn hex(bytes: &[u8]) -> String {
562 bytes.iter().map(|b| format!("{b:02x}")).collect()
563}
564
565fn embed_for(embedder: &dyn Embedder, text: &str) -> Result<Option<Vec<f32>>> {
566 Ok(Some(embedder.embed(text)?.0))
567}
568
569fn memory_key(ns: &NamespaceId, id: &MemoryId) -> String {
575 format!("{}:{}", ns.as_str(), id.0)
576}
577
578#[async_trait]
579impl Store for SurrealStore {
580 async fn store_memory(&self, ns: &NamespaceId, memory: Memory) -> Result<MemoryId> {
581 if let Some(existing) = self.check_duplicate(ns, &memory.content).await? {
584 return Err(IjimaError::duplicate(format!(
585 "content already stored as {}",
586 existing.0
587 )));
588 }
589 let (embedding, embed_model) = match &self.embedder {
590 Some(e) => {
591 (
594 embed_for(e.as_ref(), &memory.content)?,
595 Some(e.model_id().to_string()),
596 )
597 }
598 None => (None, None),
599 };
600 let record = MemoryRecord::from_memory(&memory, ns, embedding, embed_model);
601 let _: Option<surrealdb::types::SerdeWrapper<MemoryRecord>> = self
602 .db
603 .create((MEMORIES_TABLE, memory_key(ns, &memory.id)))
604 .content(surrealdb::types::SerdeWrapper(record))
605 .await
606 .map_err(store_err)?;
607 Ok(memory.id)
608 }
609
610 async fn check_duplicate(&self, ns: &NamespaceId, content: &str) -> Result<Option<MemoryId>> {
611 use sha2::{Digest, Sha256};
612 let hash = hex(&Sha256::digest(content.as_bytes()));
613 let mut result = self
614 .db
615 .query(format!(
616 "SELECT memory_id FROM {MEMORIES_TABLE}
617 WHERE namespace = $ns AND content_hash = $hash LIMIT 1"
618 ))
619 .bind(("ns", ns.as_str().to_string()))
620 .bind(("hash", hash))
621 .await
622 .map_err(store_err)?;
623 #[derive(serde::Serialize, Deserialize)]
624 struct IdRow {
625 memory_id: String,
626 }
627 let rows = take_vec::<IdRow>(&mut result)?;
628 Ok(rows.into_iter().next().map(|r| MemoryId(r.memory_id)))
629 }
630
631 async fn recall_memory(&self, ns: &NamespaceId, id: &MemoryId) -> Result<Option<Memory>> {
632 let record: Option<surrealdb::types::SerdeWrapper<MemoryRecord>> = self
633 .db
634 .select((MEMORIES_TABLE, memory_key(ns, id)))
635 .await
636 .map_err(store_err)?;
637 Ok(record.map(|w| w.0.into_memory()))
638 }
639
640 async fn delete_memory(&self, ns: &NamespaceId, id: &MemoryId) -> Result<()> {
641 let _: Option<surrealdb::types::SerdeWrapper<MemoryRecord>> = self
644 .db
645 .delete((MEMORIES_TABLE, memory_key(ns, id)))
646 .await
647 .map_err(store_err)?;
648 Ok(())
649 }
650
651 async fn list_memories(&self, ns: &NamespaceId, limit: usize) -> Result<Vec<Memory>> {
652 let mut result = self
653 .db
654 .query(format!(
655 "SELECT memory_id, content, project, topic, source, harness, session_id, namespace,
656 origin, authority, importance, created_at
657 FROM {MEMORIES_TABLE}
658 WHERE namespace = $ns
659 ORDER BY importance DESC, created_at DESC
660 LIMIT $lim"
661 ))
662 .bind(("ns", ns.as_str().to_string()))
663 .bind(("lim", limit as i64))
664 .await
665 .map_err(store_err)?;
666 let records = take_vec::<MemoryRecord>(&mut result)?;
667 Ok(records.into_iter().map(|r| r.into_memory()).collect())
668 }
669
670 async fn store_stats(&self) -> Result<StoreStats> {
671 #[derive(serde::Serialize, Deserialize)]
674 struct NsRow {
675 namespace: String,
676 }
677 let mut result = self
678 .db
679 .query(format!("SELECT namespace FROM {MEMORIES_TABLE}"))
680 .await
681 .map_err(store_err)?;
682 let rows = take_vec::<NsRow>(&mut result)?;
683 let mut counts: std::collections::BTreeMap<String, usize> =
684 std::collections::BTreeMap::new();
685 for r in &rows {
686 *counts.entry(r.namespace.clone()).or_insert(0) += 1;
687 }
688 let namespaces = counts
689 .into_iter()
690 .map(|(namespace, memories)| NamespaceCount {
691 namespace,
692 memories,
693 })
694 .collect();
695 Ok(StoreStats {
696 total_memories: rows.len(),
697 namespaces,
698 })
699 }
700
701 async fn list_rooms(
702 &self,
703 ns: &NamespaceId,
704 project: Option<&str>,
705 limit: usize,
706 ) -> Result<Vec<Room>> {
707 let counts = self.project_topic_counts(ns).await?;
708 let mut rooms: Vec<Room> = counts
709 .into_iter()
710 .filter(|((p, _), _)| project.is_none_or(|proj| p == proj))
711 .map(|((project, topic), count)| Room {
712 project,
713 topic,
714 count,
715 })
716 .collect();
717 rooms.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.topic.cmp(&b.topic)));
719 rooms.truncate(limit);
720 Ok(rooms)
721 }
722
723 async fn taxonomy(&self, ns: &NamespaceId) -> Result<Vec<ProjectTaxon>> {
724 let counts = self.project_topic_counts(ns).await?;
725 let mut by_project: std::collections::BTreeMap<String, Vec<Room>> =
727 std::collections::BTreeMap::new();
728 for ((project, topic), count) in counts {
729 by_project.entry(project).or_default().push(Room {
730 project: String::new(),
731 topic,
732 count,
733 });
734 }
735 let mut taxons: Vec<ProjectTaxon> = by_project
736 .into_iter()
737 .map(|(project, mut rooms)| {
738 rooms.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.topic.cmp(&b.topic)));
739 for r in &mut rooms {
740 r.project = project.clone();
741 }
742 let total = rooms.iter().map(|r| r.count).sum();
743 ProjectTaxon {
744 project,
745 rooms,
746 total,
747 }
748 })
749 .collect();
750 taxons.sort_by(|a, b| {
752 b.total
753 .cmp(&a.total)
754 .then_with(|| a.project.cmp(&b.project))
755 });
756 Ok(taxons)
757 }
758
759 async fn palace_graph(&self, ns: &NamespaceId) -> Result<PalaceGraph> {
760 let counts = self.project_topic_counts(ns).await?;
761 let mut topics: std::collections::BTreeMap<String, Vec<(String, usize)>> =
763 std::collections::BTreeMap::new();
764 let mut projects: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
765 for ((project, topic), count) in counts {
766 projects.insert(project.clone());
767 topics.entry(topic).or_default().push((project, count));
768 }
769 let mut tunnels = Vec::new();
770 for (topic, mut entries) in topics {
771 if entries.len() < 2 {
772 continue; }
774 entries.sort_by(|a, b| a.0.cmp(&b.0));
775 for i in 0..entries.len() {
777 for j in (i + 1)..entries.len() {
778 tunnels.push(Tunnel {
779 topic: topic.clone(),
780 project_a: entries[i].0.clone(),
781 project_b: entries[j].0.clone(),
782 count_a: entries[i].1,
783 count_b: entries[j].1,
784 });
785 }
786 }
787 }
788 Ok(PalaceGraph {
789 projects: projects.into_iter().collect(),
790 tunnels,
791 })
792 }
793
794 async fn traverse_tunnel(
795 &self,
796 ns: &NamespaceId,
797 topic: &str,
798 project_a: &str,
799 project_b: &str,
800 limit: usize,
801 ) -> Result<TunnelTraversal> {
802 let memories_a = self
803 .project_topic_memories(ns, project_a, topic, limit)
804 .await?;
805 let memories_b = self
806 .project_topic_memories(ns, project_b, topic, limit)
807 .await?;
808 Ok(TunnelTraversal {
809 topic: topic.to_string(),
810 project_a: project_a.to_string(),
811 project_b: project_b.to_string(),
812 memories_a,
813 memories_b,
814 })
815 }
816
817 async fn search_memories(
818 &self,
819 ns: &NamespaceId,
820 embedding: &Embedding,
821 limit: usize,
822 ) -> Result<Vec<SearchHit>> {
823 if self.embedder.is_none() {
824 return Err(IjimaError::Store {
825 detail: "search requires the store to opened with an Embedder".into(),
826 });
827 }
828 let mut result = self
832 .db
833 .query(format!(
834 "SELECT memory_id, content, project, topic, source, harness, session_id, namespace,
835 origin, authority, importance, created_at,
836 vector::similarity::cosine(embedding, $query) AS score
837 FROM {MEMORIES_TABLE}
838 WHERE namespace = $ns AND embedding IS NOT NONE
839 ORDER BY score DESC
840 LIMIT $lim"
841 ))
842 .bind(("ns", ns.as_str().to_string()))
843 .bind(("query", embedding.0.clone()))
844 .bind(("lim", limit as i64))
845 .await
846 .map_err(store_err)?;
847 #[derive(serde::Serialize, Deserialize)]
850 struct ScoredRecord {
851 #[serde(flatten)]
852 rec: MemoryRecord,
853 score: f64,
854 }
855 let rows = take_vec::<ScoredRecord>(&mut result)?;
856 Ok(rows
857 .into_iter()
858 .map(|r| SearchHit {
859 memory: r.rec.into_memory(),
860 similarity: r.score as f32,
861 })
862 .collect())
863 }
864
865 async fn ingest_turn(&self, ns: &NamespaceId, turn: SessionTurn) -> Result<()> {
866 let record = SessionTurnRecord {
867 session_id: turn.session_id.0.clone(),
868 turn_index: turn.turn_index,
869 role: turn.role,
870 content: turn.content,
871 timestamp: turn.timestamp,
872 namespace: ns.as_str().to_string(),
873 };
874 let _: Option<surrealdb::types::SerdeWrapper<SessionTurnRecord>> = self
875 .db
876 .create(TURNS_TABLE)
877 .content(surrealdb::types::SerdeWrapper(record))
878 .await
879 .map_err(store_err)?;
880 Ok(())
881 }
882
883 async fn session_turns(
884 &self,
885 ns: &NamespaceId,
886 session: &SessionId,
887 limit: usize,
888 ) -> Result<Vec<SessionTurn>> {
889 let mut result = self
890 .db
891 .query(format!(
892 "SELECT session_id, turn_index, role, content, timestamp, namespace
893 FROM {TURNS_TABLE}
894 WHERE namespace = $ns AND session_id = $sid
895 ORDER BY turn_index DESC LIMIT $lim",
896 ))
897 .bind(("ns", ns.as_str().to_string()))
898 .bind(("sid", session.0.clone()))
899 .bind(("lim", limit as i64))
900 .await
901 .map_err(store_err)?;
902 let mut records = take_vec::<SessionTurnRecord>(&mut result)?;
903 records.reverse(); Ok(records
905 .into_iter()
906 .map(|r| SessionTurn {
907 session_id: SessionId(r.session_id),
908 turn_index: r.turn_index,
909 role: r.role,
910 content: r.content,
911 timestamp: r.timestamp,
912 })
913 .collect())
914 }
915
916 async fn create_session(&self, ns: &NamespaceId, session: Session) -> Result<SessionId> {
917 let id_str = session.id.0.clone();
918 let record = SessionRecord::from_session(&session, ns);
919 let _: Option<surrealdb::types::SerdeWrapper<SessionRecord>> = self
922 .db
923 .upsert((SESSIONS_TABLE, id_str.clone()))
924 .content(surrealdb::types::SerdeWrapper(record))
925 .await
926 .map_err(store_err)?;
927 Ok(session.id)
928 }
929
930 async fn list_sessions(
931 &self,
932 ns: &NamespaceId,
933 harness: Option<&Harness>,
934 limit: usize,
935 ) -> Result<Vec<Session>> {
936 let mut query = format!(
937 "SELECT session_id, harness, channel, started_at, ended_at, namespace
938 FROM {SESSIONS_TABLE}
939 WHERE namespace = $ns"
940 );
941 if harness.is_some() {
942 query.push_str(" AND harness = $harness");
943 }
944 query.push_str(" ORDER BY started_at DESC LIMIT $lim");
945 let mut q = self
946 .db
947 .query(query)
948 .bind(("ns", ns.as_str().to_string()))
949 .bind(("lim", limit as i64));
950 if let Some(h) = harness {
951 q = q.bind(("harness", h.as_wire_str().to_string()));
952 }
953 let res = q.await.map_err(store_err)?;
954 let mut res = res;
955 let rows = take_vec::<SessionRecord>(&mut res)?;
956 Ok(rows.into_iter().map(SessionRecord::into_session).collect())
957 }
958
959 async fn end_session(
960 &self,
961 ns: &NamespaceId,
962 session: &SessionId,
963 ended_at: String,
964 ) -> Result<()> {
965 let _ = self
968 .db
969 .query(format!(
970 "UPDATE {SESSIONS_TABLE}
971 SET ended_at = $ended
972 WHERE namespace = $ns AND session_id = $sid"
973 ))
974 .bind(("ns", ns.as_str().to_string()))
975 .bind(("sid", session.0.clone()))
976 .bind(("ended", ended_at))
977 .await
978 .map_err(store_err)?;
979 Ok(())
980 }
981
982 async fn enqueue_extraction(
983 &self,
984 ns: &NamespaceId,
985 memory: Memory,
986 confidence: f32,
987 ) -> Result<String> {
988 let now_ms = now_millis();
989 let queue_id = format!("q_{}_{now_ms}", memory.id.0);
990 let record = QueueRecord::from_extraction(&queue_id, &memory, confidence, ns, now_ms);
991 let _: Option<surrealdb::types::SerdeWrapper<QueueRecord>> = self
992 .db
993 .create((QUEUE_TABLE, queue_id.clone()))
994 .content(surrealdb::types::SerdeWrapper(record))
995 .await
996 .map_err(store_err)?;
997 Ok(queue_id)
998 }
999
1000 async fn list_pending(&self, ns: &NamespaceId, limit: usize) -> Result<Vec<QueuedExtraction>> {
1001 let mut result = self
1002 .db
1003 .query(format!(
1004 "SELECT queue_id, memory_id, content, project, topic, importance, session_id, harness, confidence, source_session_id, ts, namespace
1005 FROM {QUEUE_TABLE}
1006 WHERE namespace = $ns
1007 ORDER BY ts DESC LIMIT $lim"
1008 ))
1009 .bind(("ns", ns.as_str().to_string()))
1010 .bind(("lim", limit as i64))
1011 .await
1012 .map_err(store_err)?;
1013 let records = take_vec::<QueueRecord>(&mut result)?;
1014 Ok(records.into_iter().map(QueueRecord::into_queued).collect())
1015 }
1016
1017 async fn accept_extraction(
1018 &self,
1019 ns: &NamespaceId,
1020 queue_id: &str,
1021 ) -> Result<AcceptedExtraction> {
1022 let mut result = self
1025 .db
1026 .query(format!(
1027 "SELECT queue_id, memory_id, content, project, topic, importance, session_id, harness, confidence, source_session_id, ts, namespace
1028 FROM {QUEUE_TABLE}
1029 WHERE namespace = $ns AND queue_id = $qid LIMIT 1"
1030 ))
1031 .bind(("ns", ns.as_str().to_string()))
1032 .bind(("qid", queue_id.to_string()))
1033 .await
1034 .map_err(store_err)?;
1035 let records = take_vec::<QueueRecord>(&mut result)?;
1036 let Some(record) = records.into_iter().next() else {
1037 return Err(IjimaError::not_found(format!("queue entry {queue_id}")));
1038 };
1039 let memory = record.into_memory();
1040 let memory_id = memory.id.clone();
1041 self.store_memory(ns, memory).await?;
1043 let _: Option<surrealdb::types::SerdeWrapper<QueueRecord>> = self
1044 .db
1045 .delete((QUEUE_TABLE, queue_id.to_string()))
1046 .await
1047 .map_err(store_err)?;
1048 Ok(AcceptedExtraction { memory_id })
1049 }
1050
1051 async fn reject_extraction(&self, ns: &NamespaceId, queue_id: &str) -> Result<()> {
1052 let _ = self
1054 .db
1055 .query(format!(
1056 "DELETE FROM {QUEUE_TABLE}
1057 WHERE namespace = $ns AND queue_id = $qid"
1058 ))
1059 .bind(("ns", ns.as_str().to_string()))
1060 .bind(("qid", queue_id.to_string()))
1061 .await
1062 .map_err(store_err)?;
1063 Ok(())
1064 }
1065 async fn write_diary(&self, ns: &NamespaceId, entry: DiaryEntry) -> Result<()> {
1066 let now_ms = std::time::SystemTime::now()
1067 .duration_since(std::time::UNIX_EPOCH)
1068 .map(|d| d.as_millis() as i64)
1069 .unwrap_or(0);
1070 let record = DiaryRecord::from_entry(&entry, ns, now_ms);
1071 let _: Option<surrealdb::types::SerdeWrapper<DiaryRecord>> = self
1072 .db
1073 .create(DIARY_TABLE)
1074 .content(surrealdb::types::SerdeWrapper(record))
1075 .await
1076 .map_err(store_err)?;
1077 Ok(())
1078 }
1079
1080 async fn read_diary(
1081 &self,
1082 ns: &NamespaceId,
1083 agent: &str,
1084 limit: usize,
1085 ) -> Result<Vec<DiaryEntry>> {
1086 let mut result = self
1087 .db
1088 .query(format!(
1089 "SELECT agent, content, topic, timestamp, ts, namespace
1090 FROM {DIARY_TABLE}
1091 WHERE namespace = $ns AND agent = $agent
1092 ORDER BY ts DESC LIMIT $lim"
1093 ))
1094 .bind(("ns", ns.as_str().to_string()))
1095 .bind(("agent", agent.to_string()))
1096 .bind(("lim", limit as i64))
1097 .await
1098 .map_err(store_err)?;
1099 let mut records = take_vec::<DiaryRecord>(&mut result)?;
1100 records.reverse(); Ok(records.into_iter().map(DiaryRecord::into_entry).collect())
1102 }
1103
1104 async fn register_repo(&self, repo: RepoDirectory) -> Result<()> {
1107 let name = repo.name.clone();
1108 let _: Option<surrealdb::types::SerdeWrapper<RepoDirectory>> = self
1109 .db
1110 .upsert((REPO_TABLE, name))
1111 .content(surrealdb::types::SerdeWrapper(repo))
1112 .await
1113 .map_err(store_err)?;
1114 Ok(())
1115 }
1116
1117 async fn list_repos(&self) -> Result<Vec<RepoDirectory>> {
1118 let mut result = self
1119 .db
1120 .query(format!("SELECT * FROM {REPO_TABLE} ORDER BY name"))
1121 .await
1122 .map_err(store_err)?;
1123 let repos = take_vec::<RepoDirectory>(&mut result)?;
1124 Ok(repos)
1125 }
1126
1127 async fn revoke_token(&self, revocation: TokenRevocation) -> Result<()> {
1130 let hash = revocation.token_hash.clone();
1131 let _: Option<surrealdb::types::SerdeWrapper<TokenRevocation>> = self
1132 .db
1133 .upsert((REVOCATIONS_TABLE, hash))
1134 .content(surrealdb::types::SerdeWrapper(revocation))
1135 .await
1136 .map_err(store_err)?;
1137 Ok(())
1138 }
1139
1140 async fn list_revocations(&self) -> Result<Vec<TokenRevocation>> {
1141 let mut result = self
1142 .db
1143 .query(format!(
1144 "SELECT * FROM {REVOCATIONS_TABLE} ORDER BY revoked_at_unix"
1145 ))
1146 .await
1147 .map_err(store_err)?;
1148 let revs = take_vec::<TokenRevocation>(&mut result)?;
1149 Ok(revs)
1150 }
1151
1152 async fn grant_namespace_membership(&self, membership: NamespaceMembership) -> Result<()> {
1153 let key = format!("{}:{}", membership.namespace, membership.principal);
1156 let _: Option<surrealdb::types::SerdeWrapper<NamespaceMembership>> = self
1157 .db
1158 .upsert((NAMESPACE_MEMBERS_TABLE, key))
1159 .content(surrealdb::types::SerdeWrapper(membership))
1160 .await
1161 .map_err(store_err)?;
1162 Ok(())
1163 }
1164
1165 async fn revoke_namespace_membership(&self, ns: &NamespaceId, principal: &str) -> Result<()> {
1166 let key = format!("{}:{}", ns.as_str(), principal);
1167 let _: Option<surrealdb::types::SerdeWrapper<NamespaceMembership>> = self
1168 .db
1169 .delete((NAMESPACE_MEMBERS_TABLE, key))
1170 .await
1171 .map_err(store_err)?;
1172 Ok(())
1173 }
1174
1175 async fn list_namespace_members(&self, ns: &NamespaceId) -> Result<Vec<NamespaceMembership>> {
1176 let mut result = self
1177 .db
1178 .query(format!(
1179 "SELECT * FROM {NAMESPACE_MEMBERS_TABLE} WHERE namespace = $ns \
1180 ORDER BY granted_at_unix"
1181 ))
1182 .bind(("ns", ns.as_str().to_string()))
1183 .await
1184 .map_err(store_err)?;
1185 take_vec::<NamespaceMembership>(&mut result)
1186 }
1187
1188 async fn is_namespace_member(&self, ns: &NamespaceId, principal: &str) -> Result<bool> {
1189 let key = format!("{}:{}", ns.as_str(), principal);
1192 let record: Option<surrealdb::types::SerdeWrapper<NamespaceMembership>> = self
1193 .db
1194 .select((NAMESPACE_MEMBERS_TABLE, key))
1195 .await
1196 .map_err(store_err)?;
1197 Ok(record.is_some())
1198 }
1199}
1200
1201#[derive(Debug, Clone, Serialize, Deserialize)]
1204struct EntityRecord_ {
1205 name: String,
1206 entity_type: String,
1207 namespace: String,
1208}
1209
1210#[derive(Debug, Clone, Serialize, Deserialize)]
1211struct TripleRecord {
1212 triple_id: String,
1214 subject: String,
1215 predicate: String,
1216 object: String,
1217 #[serde(default)]
1218 valid_from: Option<String>,
1219 #[serde(default)]
1220 valid_to: Option<String>,
1221 #[serde(default = "default_record_importance")]
1222 confidence: f32,
1223 namespace: String,
1224 #[serde(default)]
1225 source_memory_id: Option<String>,
1226}
1227
1228impl TripleRecord {
1229 fn into_triple(self) -> Triple {
1230 Triple {
1231 id: self.triple_id,
1232 subject: EntityId(self.subject),
1233 predicate: self.predicate,
1234 object: EntityId(self.object),
1235 valid_from: self.valid_from,
1236 valid_to: self.valid_to,
1237 confidence: self.confidence,
1238 namespace: self.namespace,
1239 source_memory_id: self.source_memory_id,
1240 }
1241 }
1242}
1243
1244#[async_trait::async_trait]
1245impl KnowledgeGraph for SurrealStore {
1246 async fn add_triple(
1247 &self,
1248 ns: &NamespaceId,
1249 subject: EntityId,
1250 predicate: &str,
1251 object: EntityId,
1252 valid_from: Option<&str>,
1253 confidence: f32,
1254 source_memory_id: Option<&str>,
1255 ) -> Result<Triple> {
1256 let ns_str = ns.as_str().to_string();
1257 let subj = subject.0.clone();
1258 let obj = object.0.clone();
1259 for eid in [&subj, &obj] {
1261 let record = EntityRecord_ {
1262 name: eid.clone(),
1263 entity_type: "unknown".into(),
1264 namespace: ns_str.clone(),
1265 };
1266 let _: std::result::Result<Option<surrealdb::types::SerdeWrapper<EntityRecord_>>, _> =
1267 self.db
1268 .create((ENTITIES_TABLE, eid.clone()))
1269 .content(surrealdb::types::SerdeWrapper(record))
1270 .await;
1271 }
1273 let triple_id = format!("{subj}:{predicate}:{obj}");
1277 let record = TripleRecord {
1278 triple_id: triple_id.clone(),
1279 subject: subj.clone(),
1280 predicate: predicate.to_string(),
1281 object: obj.clone(),
1282 valid_from: valid_from.map(str::to_string),
1283 valid_to: None,
1284 confidence,
1285 namespace: ns_str.clone(),
1286 source_memory_id: source_memory_id.map(str::to_string),
1287 };
1288 let _: Option<surrealdb::types::SerdeWrapper<TripleRecord>> = self
1289 .db
1290 .create((TRIPLES_TABLE, triple_id.clone()))
1291 .content(surrealdb::types::SerdeWrapper(record))
1292 .await
1293 .map_err(store_err)?;
1294 Ok(Triple {
1295 id: triple_id,
1296 subject: EntityId(subj),
1297 predicate: predicate.to_string(),
1298 object: EntityId(obj),
1299 valid_from: valid_from.map(str::to_string),
1300 valid_to: None,
1301 confidence,
1302 namespace: ns_str,
1303 source_memory_id: source_memory_id.map(str::to_string),
1304 })
1305 }
1306
1307 async fn query_entity(&self, ns: &NamespaceId, entity: &EntityId) -> Result<EntityRecord> {
1308 let eid = entity.0.clone();
1309 let ent: Option<surrealdb::types::SerdeWrapper<EntityRecord_>> = self
1310 .db
1311 .select((ENTITIES_TABLE, eid.clone()))
1312 .await
1313 .map_err(store_err)?;
1314 let entity_node = ent.and_then(|w| {
1315 let r = w.0;
1316 if r.namespace == ns.as_str() {
1317 Some(Entity {
1318 id: entity.clone(),
1319 name: r.name,
1320 entity_type: r.entity_type,
1321 namespace: r.namespace,
1322 })
1323 } else {
1324 None
1325 }
1326 });
1327 let mut res = self
1328 .db
1329 .query(format!(
1330 "SELECT triple_id, subject, predicate, object, valid_from, valid_to, confidence, namespace, source_memory_id
1331 FROM {TRIPLES_TABLE}
1332 WHERE namespace = $ns AND (subject = $eid OR object = $eid)"
1333 ))
1334 .bind(("ns", ns.as_str().to_string()))
1335 .bind(("eid", eid))
1336 .await
1337 .map_err(store_err)?;
1338 let triples = take_vec::<TripleRecord>(&mut res)?;
1339 let requested = entity.0.as_str();
1340 let mut outgoing = Vec::new();
1341 let mut incoming = Vec::new();
1342 for t in triples {
1343 if t.subject == requested {
1344 outgoing.push(t.into_triple());
1345 } else {
1346 incoming.push(t.into_triple());
1347 }
1348 }
1349 Ok(EntityRecord {
1350 entity: entity_node,
1351 outgoing,
1352 incoming,
1353 })
1354 }
1355
1356 async fn invalidate_triple(&self, ns: &NamespaceId, triple_id: &str) -> Result<()> {
1357 let now = std::time::SystemTime::now()
1358 .duration_since(std::time::UNIX_EPOCH)
1359 .map(|d| d.as_secs().to_string())
1360 .unwrap_or_default();
1361 let _ = self
1362 .db
1363 .query(format!(
1364 "UPDATE {TRIPLES_TABLE} SET valid_to = $now
1365 WHERE triple_id = $tid AND namespace = $ns"
1366 ))
1367 .bind(("now", now))
1368 .bind(("tid", triple_id.to_string()))
1369 .bind(("ns", ns.as_str().to_string()))
1370 .await
1371 .map_err(store_err)?;
1372 Ok(())
1373 }
1374
1375 async fn find_triples(
1376 &self,
1377 ns: &NamespaceId,
1378 subject: Option<&EntityId>,
1379 predicate: Option<&str>,
1380 object: Option<&EntityId>,
1381 ) -> Result<Vec<Triple>> {
1382 let mut conditions = vec!["namespace = $ns".to_string()];
1383 if subject.is_some() {
1384 conditions.push("subject = $subj".into());
1385 }
1386 if predicate.is_some() {
1387 conditions.push("predicate = $pred".into());
1388 }
1389 if object.is_some() {
1390 conditions.push("object = $obj".into());
1391 }
1392 let where_ = conditions.join(" AND ");
1393 let mut q = self
1394 .db
1395 .query(format!(
1396 "SELECT triple_id, subject, predicate, object, valid_from, valid_to, confidence, namespace, source_memory_id
1397 FROM {TRIPLES_TABLE} WHERE {where_} LIMIT 100"
1398 ))
1399 .bind(("ns", ns.as_str().to_string()));
1400 if let Some(s) = subject {
1401 q = q.bind(("subj", s.0.clone()));
1402 }
1403 if let Some(p) = predicate {
1404 q = q.bind(("pred", p.to_string()));
1405 }
1406 if let Some(o) = object {
1407 q = q.bind(("obj", o.0.clone()));
1408 }
1409 let mut res = q.await.map_err(store_err)?;
1410 let triples = take_vec::<TripleRecord>(&mut res)?;
1411 Ok(triples.into_iter().map(TripleRecord::into_triple).collect())
1412 }
1413
1414 async fn kg_timeline(&self, ns: &NamespaceId, limit: usize) -> Result<Vec<Triple>> {
1415 let mut res = self
1416 .db
1417 .query(format!(
1418 "SELECT triple_id, subject, predicate, object, valid_from, valid_to, confidence, namespace, source_memory_id
1419 FROM {TRIPLES_TABLE}
1420 WHERE namespace = $ns
1421 ORDER BY valid_from DESC LIMIT $lim"
1422 ))
1423 .bind(("ns", ns.as_str().to_string()))
1424 .bind(("lim", limit as i64))
1425 .await
1426 .map_err(store_err)?;
1427 let triples = take_vec::<TripleRecord>(&mut res)?;
1428 Ok(triples.into_iter().map(TripleRecord::into_triple).collect())
1429 }
1430
1431 async fn knowledge_stats(&self, ns: &NamespaceId) -> Result<KgStats> {
1432 let mut res = self
1433 .db
1434 .query(format!(
1435 "SELECT namespace FROM {ENTITIES_TABLE} WHERE namespace = $ns;
1436 SELECT namespace FROM {TRIPLES_TABLE} WHERE namespace = $ns;"
1437 ))
1438 .bind(("ns", ns.as_str().to_string()))
1439 .await
1440 .map_err(store_err)?;
1441 #[derive(serde::Serialize, Deserialize)]
1442 struct Row {
1443 #[allow(dead_code)]
1444 namespace: String,
1445 }
1446 let entities = take_vec::<Row>(&mut res)?;
1447 let wrapped: Vec<surrealdb::types::SerdeWrapper<Row>> = res.take(1).map_err(store_err)?;
1448 let triples: Vec<Row> = wrapped.into_iter().map(|w| w.0).collect();
1449 Ok(KgStats {
1450 entities: entities.len(),
1451 triples: triples.len(),
1452 })
1453 }
1454
1455 async fn kg_global_stats(&self) -> Result<KgStats> {
1456 let mut res = self
1457 .db
1458 .query(format!(
1459 "SELECT namespace FROM {ENTITIES_TABLE};
1460 SELECT namespace FROM {TRIPLES_TABLE};"
1461 ))
1462 .await
1463 .map_err(store_err)?;
1464 #[derive(serde::Serialize, Deserialize)]
1465 struct Row {
1466 #[allow(dead_code)]
1467 namespace: String,
1468 }
1469 let entities = take_vec::<Row>(&mut res)?;
1470 let wrapped: Vec<surrealdb::types::SerdeWrapper<Row>> = res.take(1).map_err(store_err)?;
1471 let triples: Vec<Row> = wrapped.into_iter().map(|w| w.0).collect();
1472 Ok(KgStats {
1473 entities: entities.len(),
1474 triples: triples.len(),
1475 })
1476 }
1477}
1478
1479#[cfg(test)]
1480mod tests {
1481 use super::*;
1482 use ijima_core::{NamespaceId, TurnRole};
1483
1484 async fn fresh() -> SurrealStore {
1485 SurrealStore::open_embedded()
1486 .await
1487 .expect("embedded store must open")
1488 }
1489
1490 fn sample_memory(id: &str, content: &str) -> Memory {
1491 Memory {
1492 id: MemoryId(id.into()),
1493 content: content.into(),
1494 project: "ijima".into(),
1495 topic: "test".into(),
1496 source: MemorySource::Explicit,
1497 harness: Harness::Pi,
1498 session_id: Some("sess_1".into()),
1499 origin: InstanceId::local(),
1500 authority: AuthorityScope::local(),
1501 importance: 0.5,
1502 created_at: "0".into(),
1503 }
1504 }
1505
1506 #[tokio::test]
1507 async fn store_then_recall_round_trips() {
1508 let store = fresh().await;
1509 let ns = NamespaceId::new("ns_elliott_private");
1510 let id = store
1511 .store_memory(&ns, sample_memory("mem_1", "decided to use surrealdb"))
1512 .await
1513 .expect("store");
1514 assert_eq!(id.0.as_str(), "mem_1");
1515
1516 let got = store
1517 .recall_memory(&ns, &MemoryId("mem_1".into()))
1518 .await
1519 .expect("recall");
1520 let got = got.expect("must be present");
1521 assert_eq!(got.content, "decided to use surrealdb");
1522 assert_eq!(got.harness, Harness::Pi);
1523 assert_eq!(got.source, MemorySource::Explicit);
1524 assert_eq!(got.origin, InstanceId::local());
1527 assert_eq!(got.authority, AuthorityScope::local());
1528 }
1529
1530 #[tokio::test]
1531 async fn same_memory_id_in_two_namespaces_does_not_collide() {
1532 let store = fresh().await;
1536 let ns_a = NamespaceId::new("ns_import_laptop");
1537 let ns_b = NamespaceId::new("ns_import_desktop");
1538 store
1539 .store_memory(&ns_a, sample_memory("mem_fx1", "laptop copy"))
1540 .await
1541 .expect("store into ns_a");
1542 store
1543 .store_memory(&ns_b, sample_memory("mem_fx1", "desktop copy"))
1544 .await
1545 .expect("same id into ns_b must not collide");
1546 let a = store
1548 .recall_memory(&ns_a, &MemoryId("mem_fx1".into()))
1549 .await
1550 .expect("recall a")
1551 .expect("present a");
1552 let b = store
1553 .recall_memory(&ns_b, &MemoryId("mem_fx1".into()))
1554 .await
1555 .expect("recall b")
1556 .expect("present b");
1557 assert_eq!(a.content, "laptop copy");
1558 assert_eq!(b.content, "desktop copy");
1559 store
1561 .delete_memory(&ns_a, &MemoryId("mem_fx1".into()))
1562 .await
1563 .expect("delete a");
1564 assert!(
1565 store
1566 .recall_memory(&ns_b, &MemoryId("mem_fx1".into()))
1567 .await
1568 .expect("recall b after a delete")
1569 .is_some(),
1570 "ns_b record must survive ns_a deletion"
1571 );
1572 }
1573
1574 #[tokio::test]
1575 async fn list_memories_projects_origin_and_authority() {
1576 let store = fresh().await;
1580 let ns = NamespaceId::new("ns_originproj");
1581 let mut m = sample_memory("m_op", "origin projection probe");
1582 m.origin = InstanceId("elliotthall-laptop".into());
1583 store.store_memory(&ns, m).await.expect("store");
1584 let listed = store.list_memories(&ns, 10).await.expect("list");
1585 let got = listed.first().expect("one row");
1586 assert_eq!(got.origin.0, "elliotthall-laptop");
1587 }
1588
1589 #[tokio::test]
1590 async fn store_rejects_exact_duplicate_content() {
1591 let store = fresh().await;
1592 let ns = NamespaceId::new("ns_dedup");
1593 store
1594 .store_memory(&ns, sample_memory("m1", "identical content"))
1595 .await
1596 .expect("first store");
1597 let result = store
1599 .store_memory(&ns, sample_memory("m2", "identical content"))
1600 .await;
1601 assert!(matches!(result, Err(IjimaError::Duplicate { .. })));
1602
1603 let dup = store
1605 .check_duplicate(&ns, "identical content")
1606 .await
1607 .expect("check");
1608 assert_eq!(dup.as_ref().unwrap().0, "m1");
1609
1610 store
1612 .store_memory(&ns, sample_memory("m3", "different content"))
1613 .await
1614 .expect("distinct content stores");
1615
1616 let other = NamespaceId::new("ns_other");
1618 store
1619 .store_memory(&other, sample_memory("mx", "identical content"))
1620 .await
1621 .expect("same content, different namespace");
1622 }
1623
1624 #[tokio::test]
1625 async fn index_definitions_are_idempotent_and_dedup_still_works() {
1626 let store = fresh().await;
1631 store
1632 .db
1633 .query(SurrealStore::INDEX_DDL)
1634 .await
1635 .expect("re-defining indexes on an already-indexed store");
1636 let ns = NamespaceId::new("ns_idx");
1638 store
1639 .store_memory(&ns, sample_memory("m1", "indexed-content"))
1640 .await
1641 .expect("store");
1642 let dup = store
1643 .check_duplicate(&ns, "indexed-content")
1644 .await
1645 .expect("check");
1646 assert_eq!(dup.as_ref().unwrap().0, "m1");
1647 let other = NamespaceId::new("ns_idx_other");
1649 store
1650 .store_memory(&other, sample_memory("m2", "indexed-content"))
1651 .await
1652 .expect("same content, different ns");
1653 }
1654
1655 #[tokio::test]
1656 async fn namespace_isolation_hides_other_namespace_memories() {
1657 let store = fresh().await;
1658 let alice = NamespaceId::new("ns_alice");
1659 let bob = NamespaceId::new("ns_bob");
1660
1661 store
1662 .store_memory(&alice, sample_memory("mem_a", "alice's secret"))
1663 .await
1664 .unwrap();
1665
1666 let got = store
1667 .recall_memory(&bob, &MemoryId("mem_a".into()))
1668 .await
1669 .expect("recall");
1670 assert!(got.is_none(), "namespace isolation must hide the memory");
1671 }
1672
1673 #[tokio::test]
1674 async fn delete_removes_owned_memory_only() {
1675 let store = fresh().await;
1676 let ns = NamespaceId::new("ns_elliott_private");
1677 store
1678 .store_memory(&ns, sample_memory("mem_1", "bye"))
1679 .await
1680 .unwrap();
1681 store
1682 .delete_memory(&ns, &MemoryId("mem_1".into()))
1683 .await
1684 .expect("delete");
1685 let got = store
1686 .recall_memory(&ns, &MemoryId("mem_1".into()))
1687 .await
1688 .unwrap();
1689 assert!(got.is_none());
1690 }
1691
1692 #[tokio::test]
1693 async fn ingest_then_read_session_turns_chronologically() {
1694 let store = fresh().await;
1695 let ns = NamespaceId::new("ns_elliott_private");
1696 let session = SessionId::new("sess_1");
1697 for (i, content) in ["first", "second", "third"].iter().enumerate() {
1698 store
1699 .ingest_turn(
1700 &ns,
1701 SessionTurn {
1702 session_id: session.clone(),
1703 turn_index: i as u64,
1704 role: if i % 2 == 0 {
1705 TurnRole::User
1706 } else {
1707 TurnRole::Assistant
1708 },
1709 content: (*content).into(),
1710 timestamp: format!("2026-07-05T12:00:0{i}Z"),
1711 },
1712 )
1713 .await
1714 .expect("ingest");
1715 }
1716
1717 let turns = store
1718 .session_turns(&ns, &session, 10)
1719 .await
1720 .expect("read turns");
1721 assert_eq!(turns.len(), 3);
1722 assert_eq!(turns[0].content, "first");
1723 assert_eq!(turns[2].content, "third");
1724 }
1725
1726 #[tokio::test]
1727 async fn session_turns_respect_limit_and_namespace() {
1728 let store = fresh().await;
1729 let ns = NamespaceId::new("ns_elliott_private");
1730 let other = NamespaceId::new("ns_other");
1731 let session = SessionId::new("sess_1");
1732
1733 for i in 0..5 {
1734 store
1735 .ingest_turn(
1736 &ns,
1737 SessionTurn {
1738 session_id: session.clone(),
1739 turn_index: i,
1740 role: TurnRole::User,
1741 content: format!("turn {i}"),
1742 timestamp: format!("2026-07-05T12:00:0{i}Z"),
1743 },
1744 )
1745 .await
1746 .unwrap();
1747 }
1748 store
1749 .ingest_turn(
1750 &other,
1751 SessionTurn {
1752 session_id: session.clone(),
1753 turn_index: 0,
1754 role: TurnRole::User,
1755 content: "intruder".into(),
1756 timestamp: "2026-07-05T12:00:00Z".into(),
1757 },
1758 )
1759 .await
1760 .unwrap();
1761
1762 let last_two = store.session_turns(&ns, &session, 2).await.expect("read");
1763 assert_eq!(last_two.len(), 2);
1764 assert_eq!(last_two[1].content, "turn 4");
1765 assert!(last_two.iter().all(|t| t.content != "intruder"));
1766 }
1767
1768 #[tokio::test]
1769 async fn session_metadata_create_list_end_round_trip() {
1770 let store = fresh().await;
1771 let ns = NamespaceId::new("ns_sessions");
1772
1773 let s1 = Session {
1775 id: SessionId::new("sess_a"),
1776 harness: Harness::Pi,
1777 channel: Some("thread-1".into()),
1778 started_at: "2026-07-05T10:00:00Z".into(),
1779 ended_at: None,
1780 };
1781 let s2 = Session {
1782 id: SessionId::new("sess_b"),
1783 harness: Harness::Sakamoto,
1784 channel: None,
1785 started_at: "2026-07-05T12:00:00Z".into(),
1786 ended_at: None,
1787 };
1788 store
1789 .create_session(&ns, s1.clone())
1790 .await
1791 .expect("create s1");
1792 store
1793 .create_session(&ns, s2.clone())
1794 .await
1795 .expect("create s2");
1796
1797 let all = store.list_sessions(&ns, None, 10).await.expect("list");
1799 assert_eq!(all.len(), 2);
1800 assert_eq!(all[0].id.as_str(), "sess_b");
1801 assert_eq!(all[1].id.as_str(), "sess_a");
1802 assert!(all[0].ended_at.is_none());
1803
1804 let pi_only = store
1806 .list_sessions(&ns, Some(&Harness::Pi), 10)
1807 .await
1808 .expect("list pi");
1809 assert_eq!(pi_only.len(), 1);
1810 assert_eq!(pi_only[0].id.as_str(), "sess_a");
1811
1812 store
1814 .end_session(
1815 &ns,
1816 &SessionId::new("sess_a"),
1817 "2026-07-05T11:30:00Z".into(),
1818 )
1819 .await
1820 .expect("end");
1821 let after_end = store.list_sessions(&ns, None, 10).await.expect("list");
1822 let s1_row = after_end
1823 .iter()
1824 .find(|s| s.id.as_str() == "sess_a")
1825 .unwrap();
1826 assert_eq!(s1_row.ended_at.as_deref(), Some("2026-07-05T11:30:00Z"));
1827
1828 store.create_session(&ns, s1.clone()).await.expect("upsert");
1830 let upserted = store
1831 .list_sessions(&ns, Some(&Harness::Pi), 10)
1832 .await
1833 .expect("list");
1834 assert_eq!(upserted.len(), 1);
1835
1836 let other = NamespaceId::new("ns_other");
1838 let cross = store
1839 .list_sessions(&other, None, 10)
1840 .await
1841 .expect("list other");
1842 assert!(cross.is_empty());
1843 }
1844
1845 #[tokio::test]
1848 async fn diaries_write_appends_and_reads_chronologically() {
1849 let store = fresh().await;
1850 let ns = NamespaceId::new("ns_diaries");
1851
1852 let e1 = DiaryEntry {
1853 agent: "claude".into(),
1854 content: "first entry".into(),
1855 topic: None,
1856 timestamp: String::new(),
1857 };
1858 let e2 = DiaryEntry {
1859 agent: "claude".into(),
1860 content: "second entry".into(),
1861 topic: Some("reflection".into()),
1862 timestamp: "2026-07-05T12:00:00Z".into(),
1863 };
1864 store.write_diary(&ns, e1).await.expect("write1");
1865 store.write_diary(&ns, e2).await.expect("write2");
1866
1867 let entries = store.read_diary(&ns, "claude", 10).await.expect("read");
1868 assert_eq!(entries.len(), 2);
1869 assert_eq!(entries[0].content, "first entry");
1871 assert_eq!(entries[1].content, "second entry");
1872 assert_eq!(entries[1].topic.as_deref(), Some("reflection"));
1873
1874 let other = store.read_diary(&ns, "pi", 10).await.expect("read");
1876 assert!(other.is_empty());
1877
1878 let other_ns = NamespaceId::new("ns_other");
1880 let cross = store
1881 .read_diary(&other_ns, "claude", 10)
1882 .await
1883 .expect("read");
1884 assert!(cross.is_empty());
1885 }
1886
1887 struct TestEmbedder;
1893 impl Embedder for TestEmbedder {
1894 fn dim(&self) -> usize {
1895 4
1896 }
1897 fn embed(&self, text: &str) -> Result<Embedding> {
1898 let words = text.split_whitespace().count() as f32;
1900 let has_surreal = text.contains("surreal") as u32 as f32;
1901 let has_rust = text.contains("rust") as u32 as f32;
1902 let has_memory = text.contains("memory") as u32 as f32;
1903 Ok(Embedding(vec![words, has_surreal, has_rust, has_memory]))
1904 }
1905 }
1906
1907 #[tokio::test]
1908 async fn search_without_embedder_errors() {
1909 let store = fresh().await; let ns = NamespaceId::new("ns_x");
1911 let result = store
1912 .search_memories(&ns, &Embedding(vec![0.0; 4]), 5)
1913 .await;
1914 assert!(matches!(result, Err(IjimaError::Store { .. })));
1915 }
1916
1917 #[tokio::test]
1918 async fn search_finds_nearest_with_embedder() {
1919 let store = SurrealStore::open_embedded_with(Arc::new(TestEmbedder))
1920 .await
1921 .expect("open with embedder");
1922 let ns = NamespaceId::new("ns_elliott_private");
1923
1924 store
1926 .store_memory(&ns, sample_memory("m1", "rust memory store"))
1927 .await
1928 .unwrap();
1929 store
1930 .store_memory(&ns, sample_memory("m2", "surreal db graph"))
1931 .await
1932 .unwrap();
1933 store
1934 .store_memory(&ns, sample_memory("m3", "completely unrelated text"))
1935 .await
1936 .unwrap();
1937
1938 let query = Embedding(vec![3.0, 0.0, 1.0, 1.0]);
1940 let hits = store.search_memories(&ns, &query, 2).await.expect("search");
1941 assert!(!hits.is_empty(), "must find at least one memory");
1942 assert_eq!(hits[0].memory.content, "rust memory store");
1945 assert!(hits[0].similarity > 0.0);
1947 }
1948
1949 #[tokio::test]
1950 async fn search_respects_namespace_isolation() {
1951 let store = SurrealStore::open_embedded_with(Arc::new(TestEmbedder))
1952 .await
1953 .expect("open with embedder");
1954 let alice = NamespaceId::new("ns_alice");
1955 let bob = NamespaceId::new("ns_bob");
1956
1957 store
1958 .store_memory(&alice, sample_memory("a1", "rust memory store"))
1959 .await
1960 .unwrap();
1961
1962 let query = Embedding(vec![3.0, 0.0, 1.0, 1.0]);
1965 let hits = store
1966 .search_memories(&bob, &query, 5)
1967 .await
1968 .expect("search");
1969 assert!(
1970 hits.is_empty(),
1971 "namespace isolation must hide alice's memory"
1972 );
1973 }
1974
1975 #[tokio::test]
1976 async fn list_memories_ranks_by_importance_then_recency() {
1977 let store = fresh().await;
1978 let ns = NamespaceId::new("ns_test");
1979 let mut hi = sample_memory("hi", "important");
1981 hi.importance = 0.9;
1982 hi.created_at = "100".into();
1983 let mut mid = sample_memory("mid", "medium");
1984 mid.importance = 0.5;
1985 mid.created_at = "200".into(); let mut lo = sample_memory("lo", "low");
1987 lo.importance = 0.9;
1988 lo.created_at = "300".into(); store.store_memory(&ns, hi).await.unwrap();
1990 store.store_memory(&ns, mid).await.unwrap();
1991 store.store_memory(&ns, lo).await.unwrap();
1992
1993 let list = store.list_memories(&ns, 10).await.expect("list");
1994 assert_eq!(list.len(), 3);
1997 assert_eq!(list[0].id.0, "lo"); assert_eq!(list[1].id.0, "hi"); assert_eq!(list[2].id.0, "mid"); }
2001
2002 #[tokio::test]
2005 async fn add_triple_then_query_entity() {
2006 let store = fresh().await;
2007 let ns = NamespaceId::new("ns_kg");
2008 let t = store
2009 .add_triple(
2010 &ns,
2011 EntityId::new("Ijima"),
2012 "depends_on",
2013 EntityId::new("SurrealDB"),
2014 Some("100"),
2015 1.0,
2016 None,
2017 )
2018 .await
2019 .expect("add_triple");
2020 assert_eq!(t.subject.as_str(), "Ijima");
2021 assert_eq!(t.predicate, "depends_on");
2022 assert_eq!(t.object.as_str(), "SurrealDB");
2023 assert!(t.valid_to.is_none()); let rec = store
2026 .query_entity(&ns, &EntityId::new("Ijima"))
2027 .await
2028 .expect("query");
2029 assert!(rec.entity.is_some());
2030 assert_eq!(rec.outgoing.len(), 1);
2031 assert_eq!(rec.outgoing[0].object.as_str(), "SurrealDB");
2032 assert!(rec.incoming.is_empty());
2033
2034 let rec = store
2036 .query_entity(&ns, &EntityId::new("SurrealDB"))
2037 .await
2038 .expect("query");
2039 assert_eq!(rec.incoming.len(), 1);
2040 assert!(rec.outgoing.is_empty());
2041 }
2042
2043 #[tokio::test]
2044 async fn invalidate_triple_sets_valid_to() {
2045 let store = fresh().await;
2046 let ns = NamespaceId::new("ns_kg");
2047 store
2048 .add_triple(
2049 &ns,
2050 EntityId::new("a"),
2051 "uses",
2052 EntityId::new("b"),
2053 Some("100"),
2054 1.0,
2055 None,
2056 )
2057 .await
2058 .unwrap();
2059 store
2060 .invalidate_triple(&ns, "a:uses:b")
2061 .await
2062 .expect("invalidate");
2063 let found = store
2064 .find_triples(&ns, None, Some("uses"), None)
2065 .await
2066 .expect("find");
2067 assert_eq!(found.len(), 1);
2068 assert!(found[0].valid_to.is_some(), "valid_to must be set");
2069 }
2070
2071 #[tokio::test]
2072 async fn knowledge_stats_counts_entities_and_triples() {
2073 let store = fresh().await;
2074 let ns = NamespaceId::new("ns_kg");
2075 store
2076 .add_triple(
2077 &ns,
2078 EntityId::new("a"),
2079 "x",
2080 EntityId::new("b"),
2081 None,
2082 1.0,
2083 None,
2084 )
2085 .await
2086 .unwrap();
2087 store
2088 .add_triple(
2089 &ns,
2090 EntityId::new("a"),
2091 "y",
2092 EntityId::new("c"),
2093 None,
2094 1.0,
2095 None,
2096 )
2097 .await
2098 .unwrap();
2099 let stats = store.knowledge_stats(&ns).await.expect("stats");
2100 assert_eq!(stats.entities, 3); assert_eq!(stats.triples, 2);
2102 }
2103
2104 #[tokio::test]
2105 async fn kg_namespace_isolation() {
2106 let store = fresh().await;
2107 let a = NamespaceId::new("ns_a");
2108 let b = NamespaceId::new("ns_b");
2109 store
2110 .add_triple(
2111 &a,
2112 EntityId::new("x"),
2113 "uses",
2114 EntityId::new("y"),
2115 None,
2116 1.0,
2117 None,
2118 )
2119 .await
2120 .unwrap();
2121 let stats = store.knowledge_stats(&b).await.expect("stats");
2123 assert_eq!(stats.triples, 0);
2124 let rec = store
2125 .query_entity(&b, &EntityId::new("x"))
2126 .await
2127 .expect("query");
2128 assert!(rec.outgoing.is_empty());
2129 }
2130
2131 #[tokio::test]
2132 async fn store_stats_counts_across_namespaces() {
2133 let store = fresh().await;
2134 store
2135 .store_memory(&NamespaceId::new("ns_a"), sample_memory("m1", "x"))
2136 .await
2137 .unwrap();
2138 store
2139 .store_memory(&NamespaceId::new("ns_a"), sample_memory("m2", "y"))
2140 .await
2141 .unwrap();
2142 store
2143 .store_memory(&NamespaceId::new("ns_b"), sample_memory("m3", "z"))
2144 .await
2145 .unwrap();
2146 let stats = store.store_stats().await.expect("stats");
2147 assert_eq!(stats.total_memories, 3);
2148 assert_eq!(stats.namespaces.len(), 2);
2149 let ns_a = stats
2150 .namespaces
2151 .iter()
2152 .find(|n| n.namespace == "ns_a")
2153 .expect("ns_a");
2154 assert_eq!(ns_a.memories, 2);
2155 }
2156
2157 fn mem_in(id: &str, project: &str, topic: &str, content: &str) -> Memory {
2159 Memory {
2160 id: MemoryId(id.into()),
2161 content: content.into(),
2162 project: project.into(),
2163 topic: topic.into(),
2164 source: MemorySource::Explicit,
2165 harness: Harness::Pi,
2166 session_id: None,
2167 origin: InstanceId::local(),
2168 authority: AuthorityScope::local(),
2169 importance: 0.5,
2170 created_at: "0".into(),
2171 }
2172 }
2173
2174 #[tokio::test]
2175 async fn palace_organization_rooms_taxonomy_graph_tunnel() {
2176 let store = fresh().await;
2177 let ns = NamespaceId::new("ns_palace");
2178 store
2181 .store_memory(&ns, mem_in("m1", "ijima", "auth", "use schubert"))
2182 .await
2183 .unwrap();
2184 store
2185 .store_memory(&ns, mem_in("m2", "ijima", "auth", "proof tokens"))
2186 .await
2187 .unwrap();
2188 store
2189 .store_memory(&ns, mem_in("m3", "ijima", "store", "surrealdb"))
2190 .await
2191 .unwrap();
2192 store
2193 .store_memory(&ns, mem_in("m4", "karpal", "auth", "workspace manifest"))
2194 .await
2195 .unwrap();
2196 store
2197 .store_memory(&ns, mem_in("m5", "karpal", "build", "ci"))
2198 .await
2199 .unwrap();
2200
2201 let rooms = store.list_rooms(&ns, None, 100).await.expect("rooms");
2203 assert_eq!(rooms.len(), 4);
2204 assert_eq!(rooms[0].count, 2); assert_eq!(rooms[0].project, "ijima");
2206 assert_eq!(rooms[0].topic, "auth");
2207
2208 let karpal_rooms = store
2210 .list_rooms(&ns, Some("karpal"), 100)
2211 .await
2212 .expect("rooms karpal");
2213 assert_eq!(karpal_rooms.len(), 2);
2214 assert!(karpal_rooms.iter().all(|r| r.project == "karpal"));
2215
2216 let taxons = store.taxonomy(&ns).await.expect("taxonomy");
2218 assert_eq!(taxons.len(), 2);
2219 let ijima_t = taxons
2220 .iter()
2221 .find(|t| t.project == "ijima")
2222 .expect("ijima taxon");
2223 assert_eq!(ijima_t.total, 3);
2224 assert_eq!(ijima_t.rooms.len(), 2); let graph = store.palace_graph(&ns).await.expect("graph");
2228 assert_eq!(graph.projects.len(), 2);
2229 assert!(graph.projects.contains(&"ijima".to_string()));
2230 assert!(graph.projects.contains(&"karpal".to_string()));
2231 let auth_tunnel = graph
2232 .tunnels
2233 .iter()
2234 .find(|t| t.topic == "auth")
2235 .expect("auth tunnel");
2236 assert_eq!(auth_tunnel.count_a, 2); assert_eq!(auth_tunnel.count_b, 1); let trav = store
2241 .traverse_tunnel(&ns, "auth", "ijima", "karpal", 10)
2242 .await
2243 .expect("traverse");
2244 assert_eq!(trav.memories_a.len(), 2);
2245 assert_eq!(trav.memories_b.len(), 1);
2246 assert!(trav.memories_a.iter().all(|m| m.project == "ijima"));
2247 assert!(trav.memories_b.iter().all(|m| m.project == "karpal"));
2248 }
2249
2250 #[tokio::test]
2251 async fn mining_queue_enqueue_list_accept_reject() {
2252 let store = fresh().await;
2253 let ns = NamespaceId::new("ns_mining");
2254 let other = NamespaceId::new("ns_other");
2255
2256 let q1 = store
2258 .enqueue_extraction(&ns, sample_memory("m1", "decided to use surrealdb"), 0.6)
2259 .await
2260 .expect("enqueue1");
2261 let q2 = store
2262 .enqueue_extraction(&ns, sample_memory("m2", "see https://example.com"), 0.55)
2263 .await
2264 .expect("enqueue2");
2265
2266 let pending = store.list_pending(&ns, 10).await.expect("list");
2268 assert_eq!(pending.len(), 2);
2269 assert_eq!(pending[0].id, q2); assert_eq!(pending[1].confidence, 0.6);
2271
2272 let cross = store.list_pending(&other, 10).await.expect("list other");
2274 assert!(cross.is_empty());
2275
2276 let accepted = store.accept_extraction(&ns, &q1).await.expect("accept");
2278 assert_eq!(accepted.memory_id.0, "m1");
2279 let after_accept = store.list_pending(&ns, 10).await.expect("list");
2280 assert_eq!(after_accept.len(), 1);
2281 let promoted = store
2283 .recall_memory(&ns, &MemoryId("m1".into()))
2284 .await
2285 .expect("recall");
2286 assert!(promoted.is_some());
2287
2288 store.reject_extraction(&ns, &q2).await.expect("reject");
2290 let after_reject = store.list_pending(&ns, 10).await.expect("list");
2291 assert!(after_reject.is_empty());
2292 let not_promoted = store
2293 .recall_memory(&ns, &MemoryId("m2".into()))
2294 .await
2295 .expect("recall");
2296 assert!(not_promoted.is_none());
2297
2298 let q3 = store
2300 .enqueue_extraction(&ns, sample_memory("m3", "cross test"), 0.5)
2301 .await
2302 .expect("enqueue3");
2303 let cross_accept = store.accept_extraction(&other, &q3).await;
2304 assert!(cross_accept.is_err());
2305 }
2306
2307 #[tokio::test]
2308 async fn persistent_store_survives_reopen() {
2309 let dir = std::env::temp_dir().join(format!(
2310 "ijima-persist-test-{}",
2311 std::time::SystemTime::now()
2312 .duration_since(std::time::UNIX_EPOCH)
2313 .unwrap()
2314 .as_nanos()
2315 ));
2316 let ns = NamespaceId::new("ns_persist");
2317
2318 {
2324 let dir = dir.clone();
2325 let ns = ns.clone();
2326 tokio::spawn(async move {
2327 let store = SurrealStore::open_persistent(&dir).await.expect("open");
2328 store
2329 .store_memory(&ns, sample_memory("mem_p", "survives restart"))
2330 .await
2331 .expect("store");
2332 })
2333 .await
2334 .expect("phase 1 task");
2335 for _ in 0..10 {
2336 tokio::task::yield_now().await;
2337 }
2338 tokio::time::sleep(std::time::Duration::from_millis(200)).await;
2339 }
2340
2341 let store = SurrealStore::open_persistent(&dir).await.expect("reopen");
2343 let got = store
2344 .recall_memory(&ns, &MemoryId("mem_p".into()))
2345 .await
2346 .expect("recall")
2347 .expect("memory must survive reopen");
2348 assert_eq!(got.content, "survives restart");
2349 }
2350
2351 #[tokio::test]
2352 async fn export_to_writes_sql_dump_to_file() {
2353 let dir = std::env::temp_dir().join(format!(
2354 "ijima-export-test-{}",
2355 std::time::SystemTime::now()
2356 .duration_since(std::time::UNIX_EPOCH)
2357 .unwrap()
2358 .as_nanos()
2359 ));
2360 std::fs::create_dir_all(&dir).expect("create dir");
2361 let db_path = dir.join("ijima.db");
2362 let store = crate::SurrealStore::open_persistent(&db_path)
2363 .await
2364 .expect("open");
2365 let ns = NamespaceId::new("ns_export");
2366 store
2367 .store_memory(&ns, sample_memory("mem_x", "export this memory"))
2368 .await
2369 .expect("store");
2370
2371 let out = dir.join("dump.surql");
2372 store.export_to(&out).await.expect("export");
2373 let contents = std::fs::read_to_string(&out).expect("read");
2374 assert!(contents.contains("export this memory"));
2375 let _ = std::fs::remove_dir_all(&dir);
2377 }
2378}
2379
2380#[tokio::test]
2385async fn fresh_store_reads_succeed_before_any_write() {
2386 let store = SurrealStore::open_embedded().await.expect("open");
2387 let revs = store.list_revocations().await.expect("fresh revocations");
2388 assert!(revs.is_empty());
2389 let stats = store.store_stats().await.expect("fresh stats");
2390 assert_eq!(stats.total_memories, 0);
2391 assert!(stats.namespaces.is_empty());
2392}
2393
2394#[tokio::test]
2397async fn namespace_membership_crud_round_trips() {
2398 let store = SurrealStore::open_embedded().await.expect("open");
2399 let ns = NamespaceId::new("ns_ia_shared");
2400 let mk = |principal: &str, at: u64| NamespaceMembership {
2401 namespace: "ns_ia_shared".into(),
2402 principal: principal.into(),
2403 granted_at_unix: at,
2404 granted_by: "root".into(),
2405 };
2406
2407 assert!(
2409 !store
2410 .is_namespace_member(&ns, "elliott")
2411 .await
2412 .expect("check")
2413 );
2414
2415 store
2416 .grant_namespace_membership(mk("elliott", 100))
2417 .await
2418 .expect("grant");
2419 store
2420 .grant_namespace_membership(mk("sara", 200))
2421 .await
2422 .expect("grant");
2423 store
2425 .grant_namespace_membership(mk("elliott", 300))
2426 .await
2427 .expect("re-grant");
2428
2429 assert!(
2430 store
2431 .is_namespace_member(&ns, "elliott")
2432 .await
2433 .expect("check")
2434 );
2435 let members = store.list_namespace_members(&ns).await.expect("list");
2436 assert_eq!(members.len(), 2, "re-grant must not duplicate: {members:?}");
2437 assert_eq!(members[0].principal, "sara");
2440 assert_eq!(members[1].principal, "elliott");
2441 assert_eq!(members[1].granted_at_unix, 300, "re-grant refreshes");
2442
2443 store
2444 .revoke_namespace_membership(&ns, "elliott")
2445 .await
2446 .expect("revoke");
2447 assert!(
2448 !store
2449 .is_namespace_member(&ns, "elliott")
2450 .await
2451 .expect("check")
2452 );
2453 store
2455 .revoke_namespace_membership(&ns, "ghost")
2456 .await
2457 .expect("revoke absent");
2458 assert!(
2460 !store
2461 .is_namespace_member(&NamespaceId::new("ns_kellas_shared"), "sara")
2462 .await
2463 .expect("check")
2464 );
2465}
2466
2467#[tokio::test]
2471async fn revocation_upsert_and_ordering() {
2472 use ijima_core::TokenRevocation;
2473 let store = SurrealStore::open_embedded().await.expect("open");
2474 let mk = |hash: &str, at: u64| TokenRevocation {
2475 token_hash: hash.to_string(),
2476 revoked_at_unix: at,
2477 reason: None,
2478 };
2479 store.revoke_token(mk("bbb", 20)).await.expect("revoke b");
2480 store.revoke_token(mk("aaa", 10)).await.expect("revoke a");
2481 store
2483 .revoke_token(mk("bbb", 20))
2484 .await
2485 .expect("re-revoke b");
2486 let revs = store.list_revocations().await.expect("list");
2487 assert_eq!(revs.len(), 2, "upsert keeps one row per hash");
2488 assert_eq!(revs[0].token_hash, "aaa", "oldest first");
2489 assert_eq!(revs[1].token_hash, "bbb");
2490}