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 PalaceGraph, ProjectTaxon, QueuedExtraction, RepoDirectory, Result, Room, SearchHit, Session,
34 SessionId, SessionTurn, Store, StoreStats, Triple, Tunnel, TunnelTraversal,
35 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 ENTITIES_TABLE: &str = "entities";
52const TRIPLES_TABLE: &str = "triples";
54const REPO_TABLE: &str = "repo_directory";
56
57pub struct SurrealStore {
59 db: Surreal<Db>,
60 embedder: Option<Arc<dyn Embedder>>,
64}
65
66impl SurrealStore {
67 pub async fn open_embedded() -> Result<Self> {
74 Self::open_with_db(new_mem().await?, None).await
75 }
76
77 pub async fn open_embedded_with(embedder: Arc<dyn Embedder>) -> Result<Self> {
83 Self::open_with_db(new_mem().await?, Some(embedder)).await
84 }
85
86 pub async fn open_persistent(path: impl AsRef<std::path::Path>) -> Result<Self> {
95 Self::open_with_db(new_surrealkv(&path).await?, None).await
96 }
97
98 pub async fn open_persistent_with(
105 path: impl AsRef<std::path::Path>,
106 embedder: Arc<dyn Embedder>,
107 ) -> Result<Self> {
108 Self::open_with_db(new_surrealkv(&path).await?, Some(embedder)).await
109 }
110
111 async fn open_with_db(db: Surreal<Db>, embedder: Option<Arc<dyn Embedder>>) -> Result<Self> {
112 db.use_ns(SURREAL_NS)
113 .use_db(SURREAL_DB)
114 .await
115 .map_err(|e| IjimaError::Store {
116 detail: format!("surrealdb use_ns/use_db: {e}"),
117 })?;
118 db.query(Self::INDEX_DDL)
124 .await
125 .map_err(|e| IjimaError::Store {
126 detail: format!("surrealdb define indexes: {e}"),
127 })?;
128 Ok(Self { db, embedder })
129 }
130
131 const INDEX_DDL: &str = r#"
136 DEFINE INDEX IF NOT EXISTS mem_ns ON TABLE memories FIELDS namespace;
137 DEFINE INDEX IF NOT EXISTS mem_ns_hash ON TABLE memories FIELDS namespace, content_hash;
138 DEFINE INDEX IF NOT EXISTS turns_ns_sess ON TABLE session_turns FIELDS namespace, session_id;
139 DEFINE INDEX IF NOT EXISTS sess_ns ON TABLE sessions FIELDS namespace;
140 DEFINE INDEX IF NOT EXISTS diary_ns_ag ON TABLE diaries FIELDS namespace, agent;
141 DEFINE INDEX IF NOT EXISTS queue_ns ON TABLE mining_queue FIELDS namespace;
142 DEFINE INDEX IF NOT EXISTS ent_ns ON TABLE entities FIELDS namespace;
143 DEFINE INDEX IF NOT EXISTS trip_ns ON TABLE triples FIELDS namespace;
144 DEFINE INDEX IF NOT EXISTS trip_subj ON TABLE triples FIELDS subject;
145 DEFINE INDEX IF NOT EXISTS trip_obj ON TABLE triples FIELDS object;
146 "#;
147
148 async fn project_topic_counts(
151 &self,
152 ns: &NamespaceId,
153 ) -> Result<std::collections::BTreeMap<(String, String), usize>> {
154 #[derive(Deserialize)]
155 struct Row {
156 project: String,
157 topic: String,
158 }
159 let mut result = self
160 .db
161 .query(format!(
162 "SELECT project, topic FROM {MEMORIES_TABLE} WHERE namespace = $ns"
163 ))
164 .bind(("ns", ns.as_str().to_string()))
165 .await
166 .map_err(store_err)?;
167 let rows: Vec<Row> = result.take(0).map_err(store_err)?;
168 let mut counts = std::collections::BTreeMap::new();
169 for r in rows {
170 *counts.entry((r.project, r.topic)).or_insert(0) += 1;
171 }
172 Ok(counts)
173 }
174
175 async fn project_topic_memories(
178 &self,
179 ns: &NamespaceId,
180 project: &str,
181 topic: &str,
182 limit: usize,
183 ) -> Result<Vec<Memory>> {
184 let mut result = self
185 .db
186 .query(format!(
187 "SELECT memory_id, content, project, topic, source, harness, session_id, namespace, importance, created_at
188 FROM {MEMORIES_TABLE}
189 WHERE namespace = $ns AND project = $proj AND topic = $topic
190 ORDER BY importance DESC, created_at DESC LIMIT $lim"
191 ))
192 .bind(("ns", ns.as_str().to_string()))
193 .bind(("proj", project.to_string()))
194 .bind(("topic", topic.to_string()))
195 .bind(("lim", limit as i64))
196 .await
197 .map_err(store_err)?;
198 let records: Vec<MemoryRecord> = result.take(0).map_err(store_err)?;
199 Ok(records.into_iter().map(|r| r.into_memory()).collect())
200 }
201
202 pub async fn export_to(&self, path: impl AsRef<std::path::Path>) -> Result<()> {
206 self.db
207 .export(path.as_ref())
208 .await
209 .map_err(|e| IjimaError::Store {
210 detail: format!("export: {e}"),
211 })?;
212 Ok(())
213 }
214}
215
216async fn new_mem() -> Result<Surreal<Db>> {
217 Surreal::new::<Mem>(())
218 .await
219 .map_err(|e| IjimaError::Store {
220 detail: format!("surrealdb mem init: {e}"),
221 })
222}
223
224async fn new_surrealkv(path: impl AsRef<std::path::Path>) -> Result<Surreal<Db>> {
225 let path = path.as_ref();
226 if let Some(parent) = path.parent() {
227 std::fs::create_dir_all(parent).map_err(|e| IjimaError::Store {
228 detail: format!("mkdir {}: {e}", parent.display()),
229 })?;
230 }
231 Surreal::new::<SurrealKv>(path.to_string_lossy().to_string())
232 .await
233 .map_err(|e| IjimaError::Store {
234 detail: format!("surrealdb surrealkv init at {}: {e}", path.display()),
235 })
236}
237
238#[derive(Debug, Clone, Serialize, Deserialize)]
243struct MemoryRecord {
244 memory_id: String,
247 content: String,
248 #[serde(default)]
252 content_hash: String,
253 project: String,
254 topic: String,
255 source: MemorySource,
256 harness: Harness,
257 session_id: Option<String>,
258 namespace: String,
259 #[serde(default = "default_record_importance")]
260 importance: f32,
261 #[serde(default)]
262 created_at: String,
263 #[serde(default)]
266 origin: InstanceId,
267 #[serde(default)]
270 authority: AuthorityScope,
271 #[serde(default, skip_serializing_if = "Option::is_none")]
275 embed_model: Option<String>,
276 #[serde(default, skip_serializing_if = "Option::is_none")]
280 embedding: Option<Vec<f32>>,
281}
282
283fn default_record_importance() -> f32 {
284 0.5
285}
286
287impl MemoryRecord {
288 fn from_memory(
289 memory: &Memory,
290 ns: &NamespaceId,
291 embedding: Option<Vec<f32>>,
292 embed_model: Option<String>,
293 ) -> Self {
294 use sha2::{Digest, Sha256};
295 let content_hash = hex(&Sha256::digest(memory.content.as_bytes()));
296 Self {
297 memory_id: memory.id.0.clone(),
298 content: memory.content.clone(),
299 content_hash,
300 project: memory.project.clone(),
301 topic: memory.topic.clone(),
302 source: memory.source,
303 harness: memory.harness,
304 session_id: memory.session_id.clone(),
305 namespace: ns.as_str().to_string(),
306 importance: memory.importance,
307 created_at: memory.created_at.clone(),
308 origin: memory.origin.clone(),
309 authority: memory.authority.clone(),
310 embed_model,
311 embedding,
312 }
313 }
314
315 fn into_memory(self) -> Memory {
316 Memory {
317 id: MemoryId(self.memory_id),
318 content: self.content,
319 project: self.project,
320 topic: self.topic,
321 source: self.source,
322 harness: self.harness,
323 session_id: self.session_id,
324 origin: self.origin,
325 authority: self.authority,
326 importance: self.importance,
327 created_at: self.created_at,
328 }
329 }
330}
331
332#[derive(Debug, Clone, Serialize, Deserialize)]
334struct SessionTurnRecord {
335 session_id: String,
336 turn_index: u64,
337 role: ijima_core::TurnRole,
338 content: String,
339 timestamp: String,
340 namespace: String,
341}
342
343#[derive(Debug, Clone, Serialize, Deserialize)]
347struct DiaryRecord {
348 agent: String,
349 content: String,
350 #[serde(default, skip_serializing_if = "Option::is_none")]
351 topic: Option<String>,
352 timestamp: String,
353 ts: i64,
354 namespace: String,
355}
356
357impl DiaryRecord {
358 fn from_entry(entry: &ijima_core::DiaryEntry, ns: &NamespaceId, now_ms: i64) -> Self {
359 Self {
360 agent: entry.agent.clone(),
361 content: entry.content.clone(),
362 topic: entry.topic.clone(),
363 timestamp: entry.timestamp.clone(),
364 ts: now_ms,
365 namespace: ns.as_str().to_string(),
366 }
367 }
368
369 fn into_entry(self) -> ijima_core::DiaryEntry {
370 ijima_core::DiaryEntry {
371 agent: self.agent,
372 content: self.content,
373 topic: self.topic,
374 timestamp: self.timestamp,
375 }
376 }
377}
378
379#[derive(Debug, Clone, Serialize, Deserialize)]
381struct SessionRecord {
382 session_id: String,
383 harness: String,
384 #[serde(default, skip_serializing_if = "Option::is_none")]
385 channel: Option<String>,
386 started_at: String,
387 #[serde(default, skip_serializing_if = "Option::is_none")]
388 ended_at: Option<String>,
389 namespace: String,
390}
391
392impl SessionRecord {
393 fn from_session(session: &Session, ns: &NamespaceId) -> Self {
394 Self {
395 session_id: session.id.0.clone(),
396 harness: session.harness.as_wire_str().to_string(),
397 channel: session.channel.clone(),
398 started_at: session.started_at.clone(),
399 ended_at: session.ended_at.clone(),
400 namespace: ns.as_str().to_string(),
401 }
402 }
403
404 fn into_session(self) -> Session {
405 Session {
406 id: SessionId(self.session_id),
407 harness: Harness::from_wire_str(&self.harness),
408 channel: self.channel,
409 started_at: self.started_at,
410 ended_at: self.ended_at,
411 }
412 }
413}
414
415#[derive(Debug, Clone, Serialize, Deserialize)]
418struct QueueRecord {
419 queue_id: String,
420 memory_id: String,
422 content: String,
423 project: String,
424 topic: String,
425 #[serde(default)]
426 importance: f32,
427 session_id: Option<String>,
428 harness: String,
429 confidence: f32,
430 source_session_id: String,
431 ts: i64,
432 namespace: String,
433}
434
435impl QueueRecord {
436 fn from_extraction(
437 queue_id: &str,
438 memory: &Memory,
439 confidence: f32,
440 ns: &NamespaceId,
441 now_ms: i64,
442 ) -> Self {
443 Self {
444 queue_id: queue_id.to_string(),
445 memory_id: memory.id.0.clone(),
446 content: memory.content.clone(),
447 project: memory.project.clone(),
448 topic: memory.topic.clone(),
449 importance: memory.importance,
450 session_id: memory.session_id.clone(),
451 harness: memory.harness.as_wire_str().to_string(),
452 confidence,
453 source_session_id: memory.session_id.clone().unwrap_or_default(),
454 ts: now_ms,
455 namespace: ns.as_str().to_string(),
456 }
457 }
458
459 fn into_queued(self) -> QueuedExtraction {
460 let harness = Harness::from_wire_str(&self.harness);
461 let memory = Memory {
462 id: MemoryId(self.memory_id),
463 content: self.content,
464 project: self.project,
465 topic: self.topic,
466 source: MemorySource::Mined,
467 harness,
468 session_id: self.session_id,
469 origin: InstanceId::local(),
470 authority: AuthorityScope::local(),
471 importance: self.importance,
472 created_at: String::new(),
473 };
474 QueuedExtraction {
475 id: self.queue_id,
476 memory,
477 confidence: self.confidence,
478 source_session_id: self.source_session_id,
479 queued_at: self.ts.to_string(),
480 }
481 }
482
483 fn into_memory(self) -> Memory {
484 Memory {
485 id: MemoryId(self.memory_id),
486 content: self.content,
487 project: self.project,
488 topic: self.topic,
489 source: MemorySource::Mined,
490 harness: Harness::from_wire_str(&self.harness),
491 session_id: self.session_id,
492 origin: InstanceId::local(),
493 authority: AuthorityScope::local(),
494 importance: self.importance,
495 created_at: String::new(),
496 }
497 }
498}
499
500fn store_err(e: surrealdb::Error) -> IjimaError {
501 IjimaError::Store {
502 detail: e.to_string(),
503 }
504}
505
506fn now_millis() -> i64 {
508 std::time::SystemTime::now()
509 .duration_since(std::time::UNIX_EPOCH)
510 .map(|d| d.as_millis() as i64)
511 .unwrap_or(0)
512}
513
514fn hex(bytes: &[u8]) -> String {
516 bytes.iter().map(|b| format!("{b:02x}")).collect()
517}
518
519fn embed_for(embedder: &dyn Embedder, text: &str) -> Result<Option<Vec<f32>>> {
520 Ok(Some(embedder.embed(text)?.0))
521}
522
523#[async_trait]
524impl Store for SurrealStore {
525 async fn store_memory(&self, ns: &NamespaceId, memory: Memory) -> Result<MemoryId> {
526 if let Some(existing) = self.check_duplicate(ns, &memory.content).await? {
529 return Err(IjimaError::duplicate(format!(
530 "content already stored as {}",
531 existing.0
532 )));
533 }
534 let (embedding, embed_model) = match &self.embedder {
535 Some(e) => {
536 (
539 embed_for(e.as_ref(), &memory.content)?,
540 Some(e.model_id().to_string()),
541 )
542 }
543 None => (None, None),
544 };
545 let id_str = memory.id.0.clone();
546 let record = MemoryRecord::from_memory(&memory, ns, embedding, embed_model);
547 let _: Option<MemoryRecord> = self
548 .db
549 .create((MEMORIES_TABLE, id_str.clone()))
550 .content(record)
551 .await
552 .map_err(store_err)?;
553 Ok(memory.id)
554 }
555
556 async fn check_duplicate(&self, ns: &NamespaceId, content: &str) -> Result<Option<MemoryId>> {
557 use sha2::{Digest, Sha256};
558 let hash = hex(&Sha256::digest(content.as_bytes()));
559 let mut result = self
560 .db
561 .query(format!(
562 "SELECT memory_id FROM {MEMORIES_TABLE}
563 WHERE namespace = $ns AND content_hash = $hash LIMIT 1"
564 ))
565 .bind(("ns", ns.as_str().to_string()))
566 .bind(("hash", hash))
567 .await
568 .map_err(store_err)?;
569 #[derive(Deserialize)]
570 struct IdRow {
571 memory_id: String,
572 }
573 let rows: Vec<IdRow> = result.take(0).map_err(store_err)?;
574 Ok(rows.into_iter().next().map(|r| MemoryId(r.memory_id)))
575 }
576
577 async fn recall_memory(&self, ns: &NamespaceId, id: &MemoryId) -> Result<Option<Memory>> {
578 let record: Option<MemoryRecord> = self
579 .db
580 .select((MEMORIES_TABLE, id.0.clone()))
581 .await
582 .map_err(store_err)?;
583 Ok(record
584 .filter(|r| r.namespace == ns.as_str())
585 .map(|r| r.into_memory()))
586 }
587
588 async fn delete_memory(&self, ns: &NamespaceId, id: &MemoryId) -> Result<()> {
589 let existing = self.recall_memory(ns, id).await?;
591 if existing.is_none() {
592 return Ok(()); }
594 let _: Option<MemoryRecord> = self
595 .db
596 .delete((MEMORIES_TABLE, id.0.clone()))
597 .await
598 .map_err(store_err)?;
599 Ok(())
600 }
601
602 async fn list_memories(&self, ns: &NamespaceId, limit: usize) -> Result<Vec<Memory>> {
603 let mut result = self
604 .db
605 .query(format!(
606 "SELECT memory_id, content, project, topic, source, harness, session_id, namespace, importance, created_at
607 FROM {MEMORIES_TABLE}
608 WHERE namespace = $ns
609 ORDER BY importance DESC, created_at DESC
610 LIMIT $lim"
611 ))
612 .bind(("ns", ns.as_str().to_string()))
613 .bind(("lim", limit as i64))
614 .await
615 .map_err(store_err)?;
616 let records: Vec<MemoryRecord> = result.take(0).map_err(store_err)?;
617 Ok(records.into_iter().map(|r| r.into_memory()).collect())
618 }
619
620 async fn store_stats(&self) -> Result<StoreStats> {
621 #[derive(Deserialize)]
624 struct NsRow {
625 namespace: String,
626 }
627 let mut result = self
628 .db
629 .query(format!("SELECT namespace FROM {MEMORIES_TABLE}"))
630 .await
631 .map_err(store_err)?;
632 let rows: Vec<NsRow> = result.take(0).map_err(store_err)?;
633 let mut counts: std::collections::BTreeMap<String, usize> =
634 std::collections::BTreeMap::new();
635 for r in &rows {
636 *counts.entry(r.namespace.clone()).or_insert(0) += 1;
637 }
638 let namespaces = counts
639 .into_iter()
640 .map(|(namespace, memories)| NamespaceCount {
641 namespace,
642 memories,
643 })
644 .collect();
645 Ok(StoreStats {
646 total_memories: rows.len(),
647 namespaces,
648 })
649 }
650
651 async fn list_rooms(
652 &self,
653 ns: &NamespaceId,
654 project: Option<&str>,
655 limit: usize,
656 ) -> Result<Vec<Room>> {
657 let counts = self.project_topic_counts(ns).await?;
658 let mut rooms: Vec<Room> = counts
659 .into_iter()
660 .filter(|((p, _), _)| project.is_none_or(|proj| p == proj))
661 .map(|((project, topic), count)| Room {
662 project,
663 topic,
664 count,
665 })
666 .collect();
667 rooms.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.topic.cmp(&b.topic)));
669 rooms.truncate(limit);
670 Ok(rooms)
671 }
672
673 async fn taxonomy(&self, ns: &NamespaceId) -> Result<Vec<ProjectTaxon>> {
674 let counts = self.project_topic_counts(ns).await?;
675 let mut by_project: std::collections::BTreeMap<String, Vec<Room>> =
677 std::collections::BTreeMap::new();
678 for ((project, topic), count) in counts {
679 by_project.entry(project).or_default().push(Room {
680 project: String::new(),
681 topic,
682 count,
683 });
684 }
685 let mut taxons: Vec<ProjectTaxon> = by_project
686 .into_iter()
687 .map(|(project, mut rooms)| {
688 rooms.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.topic.cmp(&b.topic)));
689 for r in &mut rooms {
690 r.project = project.clone();
691 }
692 let total = rooms.iter().map(|r| r.count).sum();
693 ProjectTaxon {
694 project,
695 rooms,
696 total,
697 }
698 })
699 .collect();
700 taxons.sort_by(|a, b| {
702 b.total
703 .cmp(&a.total)
704 .then_with(|| a.project.cmp(&b.project))
705 });
706 Ok(taxons)
707 }
708
709 async fn palace_graph(&self, ns: &NamespaceId) -> Result<PalaceGraph> {
710 let counts = self.project_topic_counts(ns).await?;
711 let mut topics: std::collections::BTreeMap<String, Vec<(String, usize)>> =
713 std::collections::BTreeMap::new();
714 let mut projects: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
715 for ((project, topic), count) in counts {
716 projects.insert(project.clone());
717 topics.entry(topic).or_default().push((project, count));
718 }
719 let mut tunnels = Vec::new();
720 for (topic, mut entries) in topics {
721 if entries.len() < 2 {
722 continue; }
724 entries.sort_by(|a, b| a.0.cmp(&b.0));
725 for i in 0..entries.len() {
727 for j in (i + 1)..entries.len() {
728 tunnels.push(Tunnel {
729 topic: topic.clone(),
730 project_a: entries[i].0.clone(),
731 project_b: entries[j].0.clone(),
732 count_a: entries[i].1,
733 count_b: entries[j].1,
734 });
735 }
736 }
737 }
738 Ok(PalaceGraph {
739 projects: projects.into_iter().collect(),
740 tunnels,
741 })
742 }
743
744 async fn traverse_tunnel(
745 &self,
746 ns: &NamespaceId,
747 topic: &str,
748 project_a: &str,
749 project_b: &str,
750 limit: usize,
751 ) -> Result<TunnelTraversal> {
752 let memories_a = self
753 .project_topic_memories(ns, project_a, topic, limit)
754 .await?;
755 let memories_b = self
756 .project_topic_memories(ns, project_b, topic, limit)
757 .await?;
758 Ok(TunnelTraversal {
759 topic: topic.to_string(),
760 project_a: project_a.to_string(),
761 project_b: project_b.to_string(),
762 memories_a,
763 memories_b,
764 })
765 }
766
767 async fn search_memories(
768 &self,
769 ns: &NamespaceId,
770 embedding: &Embedding,
771 limit: usize,
772 ) -> Result<Vec<SearchHit>> {
773 if self.embedder.is_none() {
774 return Err(IjimaError::Store {
775 detail: "search requires the store to opened with an Embedder".into(),
776 });
777 }
778 let mut result = self
782 .db
783 .query(format!(
784 "SELECT memory_id, content, project, topic, source, harness, session_id, namespace,
785 origin, authority, importance, created_at,
786 vector::similarity::cosine(embedding, $query) AS score
787 FROM {MEMORIES_TABLE}
788 WHERE namespace = $ns AND embedding IS NOT NONE
789 ORDER BY score DESC
790 LIMIT $lim"
791 ))
792 .bind(("ns", ns.as_str().to_string()))
793 .bind(("query", embedding.0.clone()))
794 .bind(("lim", limit as i64))
795 .await
796 .map_err(store_err)?;
797 #[derive(Deserialize)]
800 struct ScoredRecord {
801 #[serde(flatten)]
802 rec: MemoryRecord,
803 score: f64,
804 }
805 let rows: Vec<ScoredRecord> = result.take(0).map_err(store_err)?;
806 Ok(rows
807 .into_iter()
808 .map(|r| SearchHit {
809 memory: r.rec.into_memory(),
810 similarity: r.score as f32,
811 })
812 .collect())
813 }
814
815 async fn ingest_turn(&self, ns: &NamespaceId, turn: SessionTurn) -> Result<()> {
816 let record = SessionTurnRecord {
817 session_id: turn.session_id.0.clone(),
818 turn_index: turn.turn_index,
819 role: turn.role,
820 content: turn.content,
821 timestamp: turn.timestamp,
822 namespace: ns.as_str().to_string(),
823 };
824 let _: Option<SessionTurnRecord> = self
825 .db
826 .create(TURNS_TABLE)
827 .content(record)
828 .await
829 .map_err(store_err)?;
830 Ok(())
831 }
832
833 async fn session_turns(
834 &self,
835 ns: &NamespaceId,
836 session: &SessionId,
837 limit: usize,
838 ) -> Result<Vec<SessionTurn>> {
839 let mut result = self
840 .db
841 .query(format!(
842 "SELECT session_id, turn_index, role, content, timestamp, namespace
843 FROM {TURNS_TABLE}
844 WHERE namespace = $ns AND session_id = $sid
845 ORDER BY turn_index DESC LIMIT $lim",
846 ))
847 .bind(("ns", ns.as_str().to_string()))
848 .bind(("sid", session.0.clone()))
849 .bind(("lim", limit as i64))
850 .await
851 .map_err(store_err)?;
852 let mut records: Vec<SessionTurnRecord> = result.take(0).map_err(store_err)?;
853 records.reverse(); Ok(records
855 .into_iter()
856 .map(|r| SessionTurn {
857 session_id: SessionId(r.session_id),
858 turn_index: r.turn_index,
859 role: r.role,
860 content: r.content,
861 timestamp: r.timestamp,
862 })
863 .collect())
864 }
865
866 async fn create_session(&self, ns: &NamespaceId, session: Session) -> Result<SessionId> {
867 let id_str = session.id.0.clone();
868 let record = SessionRecord::from_session(&session, ns);
869 let _: Option<SessionRecord> = self
872 .db
873 .upsert((SESSIONS_TABLE, id_str.clone()))
874 .content(record)
875 .await
876 .map_err(store_err)?;
877 Ok(session.id)
878 }
879
880 async fn list_sessions(
881 &self,
882 ns: &NamespaceId,
883 harness: Option<&Harness>,
884 limit: usize,
885 ) -> Result<Vec<Session>> {
886 let mut query = format!(
887 "SELECT session_id, harness, channel, started_at, ended_at, namespace
888 FROM {SESSIONS_TABLE}
889 WHERE namespace = $ns"
890 );
891 if harness.is_some() {
892 query.push_str(" AND harness = $harness");
893 }
894 query.push_str(" ORDER BY started_at DESC LIMIT $lim");
895 let mut q = self
896 .db
897 .query(query)
898 .bind(("ns", ns.as_str().to_string()))
899 .bind(("lim", limit as i64));
900 if let Some(h) = harness {
901 q = q.bind(("harness", h.as_wire_str().to_string()));
902 }
903 let res = q.await.map_err(store_err)?;
904 let mut res = res;
905 let rows: Vec<SessionRecord> = res.take(0).map_err(store_err)?;
906 Ok(rows.into_iter().map(SessionRecord::into_session).collect())
907 }
908
909 async fn end_session(
910 &self,
911 ns: &NamespaceId,
912 session: &SessionId,
913 ended_at: String,
914 ) -> Result<()> {
915 let _ = self
918 .db
919 .query(format!(
920 "UPDATE {SESSIONS_TABLE}
921 SET ended_at = $ended
922 WHERE namespace = $ns AND session_id = $sid"
923 ))
924 .bind(("ns", ns.as_str().to_string()))
925 .bind(("sid", session.0.clone()))
926 .bind(("ended", ended_at))
927 .await
928 .map_err(store_err)?;
929 Ok(())
930 }
931
932 async fn enqueue_extraction(
933 &self,
934 ns: &NamespaceId,
935 memory: Memory,
936 confidence: f32,
937 ) -> Result<String> {
938 let now_ms = now_millis();
939 let queue_id = format!("q_{}_{now_ms}", memory.id.0);
940 let record = QueueRecord::from_extraction(&queue_id, &memory, confidence, ns, now_ms);
941 let _: Option<QueueRecord> = self
942 .db
943 .create((QUEUE_TABLE, queue_id.clone()))
944 .content(record)
945 .await
946 .map_err(store_err)?;
947 Ok(queue_id)
948 }
949
950 async fn list_pending(&self, ns: &NamespaceId, limit: usize) -> Result<Vec<QueuedExtraction>> {
951 let mut result = self
952 .db
953 .query(format!(
954 "SELECT queue_id, memory_id, content, project, topic, importance, session_id, harness, confidence, source_session_id, ts, namespace
955 FROM {QUEUE_TABLE}
956 WHERE namespace = $ns
957 ORDER BY ts DESC LIMIT $lim"
958 ))
959 .bind(("ns", ns.as_str().to_string()))
960 .bind(("lim", limit as i64))
961 .await
962 .map_err(store_err)?;
963 let records: Vec<QueueRecord> = result.take(0).map_err(store_err)?;
964 Ok(records.into_iter().map(QueueRecord::into_queued).collect())
965 }
966
967 async fn accept_extraction(
968 &self,
969 ns: &NamespaceId,
970 queue_id: &str,
971 ) -> Result<AcceptedExtraction> {
972 let mut result = self
975 .db
976 .query(format!(
977 "SELECT queue_id, memory_id, content, project, topic, importance, session_id, harness, confidence, source_session_id, ts, namespace
978 FROM {QUEUE_TABLE}
979 WHERE namespace = $ns AND queue_id = $qid LIMIT 1"
980 ))
981 .bind(("ns", ns.as_str().to_string()))
982 .bind(("qid", queue_id.to_string()))
983 .await
984 .map_err(store_err)?;
985 let records: Vec<QueueRecord> = result.take(0).map_err(store_err)?;
986 let Some(record) = records.into_iter().next() else {
987 return Err(IjimaError::not_found(format!("queue entry {queue_id}")));
988 };
989 let memory = record.into_memory();
990 let memory_id = memory.id.clone();
991 self.store_memory(ns, memory).await?;
993 let _: Option<QueueRecord> = self
994 .db
995 .delete((QUEUE_TABLE, queue_id.to_string()))
996 .await
997 .map_err(store_err)?;
998 Ok(AcceptedExtraction { memory_id })
999 }
1000
1001 async fn reject_extraction(&self, ns: &NamespaceId, queue_id: &str) -> Result<()> {
1002 let _ = self
1004 .db
1005 .query(format!(
1006 "DELETE FROM {QUEUE_TABLE}
1007 WHERE namespace = $ns AND queue_id = $qid"
1008 ))
1009 .bind(("ns", ns.as_str().to_string()))
1010 .bind(("qid", queue_id.to_string()))
1011 .await
1012 .map_err(store_err)?;
1013 Ok(())
1014 }
1015 async fn write_diary(&self, ns: &NamespaceId, entry: DiaryEntry) -> Result<()> {
1016 let now_ms = std::time::SystemTime::now()
1017 .duration_since(std::time::UNIX_EPOCH)
1018 .map(|d| d.as_millis() as i64)
1019 .unwrap_or(0);
1020 let record = DiaryRecord::from_entry(&entry, ns, now_ms);
1021 let _: Option<DiaryRecord> = self
1022 .db
1023 .create(DIARY_TABLE)
1024 .content(record)
1025 .await
1026 .map_err(store_err)?;
1027 Ok(())
1028 }
1029
1030 async fn read_diary(
1031 &self,
1032 ns: &NamespaceId,
1033 agent: &str,
1034 limit: usize,
1035 ) -> Result<Vec<DiaryEntry>> {
1036 let mut result = self
1037 .db
1038 .query(format!(
1039 "SELECT agent, content, topic, timestamp, ts, namespace
1040 FROM {DIARY_TABLE}
1041 WHERE namespace = $ns AND agent = $agent
1042 ORDER BY ts DESC LIMIT $lim"
1043 ))
1044 .bind(("ns", ns.as_str().to_string()))
1045 .bind(("agent", agent.to_string()))
1046 .bind(("lim", limit as i64))
1047 .await
1048 .map_err(store_err)?;
1049 let mut records: Vec<DiaryRecord> = result.take(0).map_err(store_err)?;
1050 records.reverse(); Ok(records.into_iter().map(DiaryRecord::into_entry).collect())
1052 }
1053
1054 async fn register_repo(&self, repo: RepoDirectory) -> Result<()> {
1057 let name = repo.name.clone();
1058 let _: Option<RepoDirectory> = self
1059 .db
1060 .upsert((REPO_TABLE, name))
1061 .content(repo)
1062 .await
1063 .map_err(store_err)?;
1064 Ok(())
1065 }
1066
1067 async fn list_repos(&self) -> Result<Vec<RepoDirectory>> {
1068 let mut result = self
1069 .db
1070 .query(format!("SELECT * FROM {REPO_TABLE} ORDER BY name"))
1071 .await
1072 .map_err(store_err)?;
1073 let repos: Vec<RepoDirectory> = result.take(0).map_err(store_err)?;
1074 Ok(repos)
1075 }
1076}
1077
1078#[derive(Debug, Clone, Serialize, Deserialize)]
1081struct EntityRecord_ {
1082 name: String,
1083 entity_type: String,
1084 namespace: String,
1085}
1086
1087#[derive(Debug, Clone, Serialize, Deserialize)]
1088struct TripleRecord {
1089 triple_id: String,
1091 subject: String,
1092 predicate: String,
1093 object: String,
1094 #[serde(default)]
1095 valid_from: Option<String>,
1096 #[serde(default)]
1097 valid_to: Option<String>,
1098 #[serde(default = "default_record_importance")]
1099 confidence: f32,
1100 namespace: String,
1101 #[serde(default)]
1102 source_memory_id: Option<String>,
1103}
1104
1105impl TripleRecord {
1106 fn into_triple(self) -> Triple {
1107 Triple {
1108 id: self.triple_id,
1109 subject: EntityId(self.subject),
1110 predicate: self.predicate,
1111 object: EntityId(self.object),
1112 valid_from: self.valid_from,
1113 valid_to: self.valid_to,
1114 confidence: self.confidence,
1115 namespace: self.namespace,
1116 source_memory_id: self.source_memory_id,
1117 }
1118 }
1119}
1120
1121#[async_trait::async_trait]
1122impl KnowledgeGraph for SurrealStore {
1123 async fn add_triple(
1124 &self,
1125 ns: &NamespaceId,
1126 subject: EntityId,
1127 predicate: &str,
1128 object: EntityId,
1129 valid_from: Option<&str>,
1130 confidence: f32,
1131 source_memory_id: Option<&str>,
1132 ) -> Result<Triple> {
1133 let ns_str = ns.as_str().to_string();
1134 let subj = subject.0.clone();
1135 let obj = object.0.clone();
1136 for eid in [&subj, &obj] {
1138 let record = EntityRecord_ {
1139 name: eid.clone(),
1140 entity_type: "unknown".into(),
1141 namespace: ns_str.clone(),
1142 };
1143 let _: std::result::Result<Option<EntityRecord_>, _> = self
1144 .db
1145 .create((ENTITIES_TABLE, eid.clone()))
1146 .content(record)
1147 .await;
1148 }
1150 let triple_id = format!("{subj}:{predicate}:{obj}");
1154 let record = TripleRecord {
1155 triple_id: triple_id.clone(),
1156 subject: subj.clone(),
1157 predicate: predicate.to_string(),
1158 object: obj.clone(),
1159 valid_from: valid_from.map(str::to_string),
1160 valid_to: None,
1161 confidence,
1162 namespace: ns_str.clone(),
1163 source_memory_id: source_memory_id.map(str::to_string),
1164 };
1165 let _: Option<TripleRecord> = self
1166 .db
1167 .create((TRIPLES_TABLE, triple_id.clone()))
1168 .content(record)
1169 .await
1170 .map_err(store_err)?;
1171 Ok(Triple {
1172 id: triple_id,
1173 subject: EntityId(subj),
1174 predicate: predicate.to_string(),
1175 object: EntityId(obj),
1176 valid_from: valid_from.map(str::to_string),
1177 valid_to: None,
1178 confidence,
1179 namespace: ns_str,
1180 source_memory_id: source_memory_id.map(str::to_string),
1181 })
1182 }
1183
1184 async fn query_entity(&self, ns: &NamespaceId, entity: &EntityId) -> Result<EntityRecord> {
1185 let eid = entity.0.clone();
1186 let ent: Option<EntityRecord_> = self
1187 .db
1188 .select((ENTITIES_TABLE, eid.clone()))
1189 .await
1190 .map_err(store_err)?;
1191 let entity_node = ent.and_then(|r| {
1192 if r.namespace == ns.as_str() {
1193 Some(Entity {
1194 id: entity.clone(),
1195 name: r.name,
1196 entity_type: r.entity_type,
1197 namespace: r.namespace,
1198 })
1199 } else {
1200 None
1201 }
1202 });
1203 let mut res = self
1204 .db
1205 .query(format!(
1206 "SELECT triple_id, subject, predicate, object, valid_from, valid_to, confidence, namespace, source_memory_id
1207 FROM {TRIPLES_TABLE}
1208 WHERE namespace = $ns AND (subject = $eid OR object = $eid)"
1209 ))
1210 .bind(("ns", ns.as_str().to_string()))
1211 .bind(("eid", eid))
1212 .await
1213 .map_err(store_err)?;
1214 let triples: Vec<TripleRecord> = res.take(0).map_err(store_err)?;
1215 let requested = entity.0.as_str();
1216 let mut outgoing = Vec::new();
1217 let mut incoming = Vec::new();
1218 for t in triples {
1219 if t.subject == requested {
1220 outgoing.push(t.into_triple());
1221 } else {
1222 incoming.push(t.into_triple());
1223 }
1224 }
1225 Ok(EntityRecord {
1226 entity: entity_node,
1227 outgoing,
1228 incoming,
1229 })
1230 }
1231
1232 async fn invalidate_triple(&self, ns: &NamespaceId, triple_id: &str) -> Result<()> {
1233 let now = std::time::SystemTime::now()
1234 .duration_since(std::time::UNIX_EPOCH)
1235 .map(|d| d.as_secs().to_string())
1236 .unwrap_or_default();
1237 let _ = self
1238 .db
1239 .query(format!(
1240 "UPDATE {TRIPLES_TABLE} SET valid_to = $now
1241 WHERE triple_id = $tid AND namespace = $ns"
1242 ))
1243 .bind(("now", now))
1244 .bind(("tid", triple_id.to_string()))
1245 .bind(("ns", ns.as_str().to_string()))
1246 .await
1247 .map_err(store_err)?;
1248 Ok(())
1249 }
1250
1251 async fn find_triples(
1252 &self,
1253 ns: &NamespaceId,
1254 subject: Option<&EntityId>,
1255 predicate: Option<&str>,
1256 object: Option<&EntityId>,
1257 ) -> Result<Vec<Triple>> {
1258 let mut conditions = vec!["namespace = $ns".to_string()];
1259 if subject.is_some() {
1260 conditions.push("subject = $subj".into());
1261 }
1262 if predicate.is_some() {
1263 conditions.push("predicate = $pred".into());
1264 }
1265 if object.is_some() {
1266 conditions.push("object = $obj".into());
1267 }
1268 let where_ = conditions.join(" AND ");
1269 let mut q = self
1270 .db
1271 .query(format!(
1272 "SELECT triple_id, subject, predicate, object, valid_from, valid_to, confidence, namespace, source_memory_id
1273 FROM {TRIPLES_TABLE} WHERE {where_} LIMIT 100"
1274 ))
1275 .bind(("ns", ns.as_str().to_string()));
1276 if let Some(s) = subject {
1277 q = q.bind(("subj", s.0.clone()));
1278 }
1279 if let Some(p) = predicate {
1280 q = q.bind(("pred", p.to_string()));
1281 }
1282 if let Some(o) = object {
1283 q = q.bind(("obj", o.0.clone()));
1284 }
1285 let mut res = q.await.map_err(store_err)?;
1286 let triples: Vec<TripleRecord> = res.take(0).map_err(store_err)?;
1287 Ok(triples.into_iter().map(TripleRecord::into_triple).collect())
1288 }
1289
1290 async fn kg_timeline(&self, ns: &NamespaceId, limit: usize) -> Result<Vec<Triple>> {
1291 let mut res = self
1292 .db
1293 .query(format!(
1294 "SELECT triple_id, subject, predicate, object, valid_from, valid_to, confidence, namespace, source_memory_id
1295 FROM {TRIPLES_TABLE}
1296 WHERE namespace = $ns
1297 ORDER BY valid_from DESC LIMIT $lim"
1298 ))
1299 .bind(("ns", ns.as_str().to_string()))
1300 .bind(("lim", limit as i64))
1301 .await
1302 .map_err(store_err)?;
1303 let triples: Vec<TripleRecord> = res.take(0).map_err(store_err)?;
1304 Ok(triples.into_iter().map(TripleRecord::into_triple).collect())
1305 }
1306
1307 async fn knowledge_stats(&self, ns: &NamespaceId) -> Result<KgStats> {
1308 let mut res = self
1309 .db
1310 .query(format!(
1311 "SELECT namespace FROM {ENTITIES_TABLE} WHERE namespace = $ns;
1312 SELECT namespace FROM {TRIPLES_TABLE} WHERE namespace = $ns;"
1313 ))
1314 .bind(("ns", ns.as_str().to_string()))
1315 .await
1316 .map_err(store_err)?;
1317 #[derive(Deserialize)]
1318 struct Row {
1319 #[allow(dead_code)]
1320 namespace: String,
1321 }
1322 let entities: Vec<Row> = res.take(0).map_err(store_err)?;
1323 let triples: Vec<Row> = res.take(1).map_err(store_err)?;
1324 Ok(KgStats {
1325 entities: entities.len(),
1326 triples: triples.len(),
1327 })
1328 }
1329
1330 async fn kg_global_stats(&self) -> Result<KgStats> {
1331 let mut res = self
1332 .db
1333 .query(format!(
1334 "SELECT namespace FROM {ENTITIES_TABLE};
1335 SELECT namespace FROM {TRIPLES_TABLE};"
1336 ))
1337 .await
1338 .map_err(store_err)?;
1339 #[derive(Deserialize)]
1340 struct Row {
1341 #[allow(dead_code)]
1342 namespace: String,
1343 }
1344 let entities: Vec<Row> = res.take(0).map_err(store_err)?;
1345 let triples: Vec<Row> = res.take(1).map_err(store_err)?;
1346 Ok(KgStats {
1347 entities: entities.len(),
1348 triples: triples.len(),
1349 })
1350 }
1351}
1352
1353#[cfg(test)]
1354mod tests {
1355 use super::*;
1356 use ijima_core::{NamespaceId, TurnRole};
1357
1358 async fn fresh() -> SurrealStore {
1359 SurrealStore::open_embedded()
1360 .await
1361 .expect("embedded store must open")
1362 }
1363
1364 fn sample_memory(id: &str, content: &str) -> Memory {
1365 Memory {
1366 id: MemoryId(id.into()),
1367 content: content.into(),
1368 project: "ijima".into(),
1369 topic: "test".into(),
1370 source: MemorySource::Explicit,
1371 harness: Harness::Pi,
1372 session_id: Some("sess_1".into()),
1373 origin: InstanceId::local(),
1374 authority: AuthorityScope::local(),
1375 importance: 0.5,
1376 created_at: "0".into(),
1377 }
1378 }
1379
1380 #[tokio::test]
1381 async fn store_then_recall_round_trips() {
1382 let store = fresh().await;
1383 let ns = NamespaceId::new("ns_elliott_private");
1384 let id = store
1385 .store_memory(&ns, sample_memory("mem_1", "decided to use surrealdb"))
1386 .await
1387 .expect("store");
1388 assert_eq!(id.0.as_str(), "mem_1");
1389
1390 let got = store
1391 .recall_memory(&ns, &MemoryId("mem_1".into()))
1392 .await
1393 .expect("recall");
1394 let got = got.expect("must be present");
1395 assert_eq!(got.content, "decided to use surrealdb");
1396 assert_eq!(got.harness, Harness::Pi);
1397 assert_eq!(got.source, MemorySource::Explicit);
1398 assert_eq!(got.origin, InstanceId::local());
1401 assert_eq!(got.authority, AuthorityScope::local());
1402 }
1403
1404 #[tokio::test]
1405 async fn store_rejects_exact_duplicate_content() {
1406 let store = fresh().await;
1407 let ns = NamespaceId::new("ns_dedup");
1408 store
1409 .store_memory(&ns, sample_memory("m1", "identical content"))
1410 .await
1411 .expect("first store");
1412 let result = store
1414 .store_memory(&ns, sample_memory("m2", "identical content"))
1415 .await;
1416 assert!(matches!(result, Err(IjimaError::Duplicate { .. })));
1417
1418 let dup = store
1420 .check_duplicate(&ns, "identical content")
1421 .await
1422 .expect("check");
1423 assert_eq!(dup.as_ref().unwrap().0, "m1");
1424
1425 store
1427 .store_memory(&ns, sample_memory("m3", "different content"))
1428 .await
1429 .expect("distinct content stores");
1430
1431 let other = NamespaceId::new("ns_other");
1433 store
1434 .store_memory(&other, sample_memory("mx", "identical content"))
1435 .await
1436 .expect("same content, different namespace");
1437 }
1438
1439 #[tokio::test]
1440 async fn index_definitions_are_idempotent_and_dedup_still_works() {
1441 let store = fresh().await;
1446 store
1447 .db
1448 .query(SurrealStore::INDEX_DDL)
1449 .await
1450 .expect("re-defining indexes on an already-indexed store");
1451 let ns = NamespaceId::new("ns_idx");
1453 store
1454 .store_memory(&ns, sample_memory("m1", "indexed-content"))
1455 .await
1456 .expect("store");
1457 let dup = store
1458 .check_duplicate(&ns, "indexed-content")
1459 .await
1460 .expect("check");
1461 assert_eq!(dup.as_ref().unwrap().0, "m1");
1462 let other = NamespaceId::new("ns_idx_other");
1464 store
1465 .store_memory(&other, sample_memory("m2", "indexed-content"))
1466 .await
1467 .expect("same content, different ns");
1468 }
1469
1470 #[tokio::test]
1471 async fn namespace_isolation_hides_other_namespace_memories() {
1472 let store = fresh().await;
1473 let alice = NamespaceId::new("ns_alice");
1474 let bob = NamespaceId::new("ns_bob");
1475
1476 store
1477 .store_memory(&alice, sample_memory("mem_a", "alice's secret"))
1478 .await
1479 .unwrap();
1480
1481 let got = store
1482 .recall_memory(&bob, &MemoryId("mem_a".into()))
1483 .await
1484 .expect("recall");
1485 assert!(got.is_none(), "namespace isolation must hide the memory");
1486 }
1487
1488 #[tokio::test]
1489 async fn delete_removes_owned_memory_only() {
1490 let store = fresh().await;
1491 let ns = NamespaceId::new("ns_elliott_private");
1492 store
1493 .store_memory(&ns, sample_memory("mem_1", "bye"))
1494 .await
1495 .unwrap();
1496 store
1497 .delete_memory(&ns, &MemoryId("mem_1".into()))
1498 .await
1499 .expect("delete");
1500 let got = store
1501 .recall_memory(&ns, &MemoryId("mem_1".into()))
1502 .await
1503 .unwrap();
1504 assert!(got.is_none());
1505 }
1506
1507 #[tokio::test]
1508 async fn ingest_then_read_session_turns_chronologically() {
1509 let store = fresh().await;
1510 let ns = NamespaceId::new("ns_elliott_private");
1511 let session = SessionId::new("sess_1");
1512 for (i, content) in ["first", "second", "third"].iter().enumerate() {
1513 store
1514 .ingest_turn(
1515 &ns,
1516 SessionTurn {
1517 session_id: session.clone(),
1518 turn_index: i as u64,
1519 role: if i % 2 == 0 {
1520 TurnRole::User
1521 } else {
1522 TurnRole::Assistant
1523 },
1524 content: (*content).into(),
1525 timestamp: format!("2026-07-05T12:00:0{i}Z"),
1526 },
1527 )
1528 .await
1529 .expect("ingest");
1530 }
1531
1532 let turns = store
1533 .session_turns(&ns, &session, 10)
1534 .await
1535 .expect("read turns");
1536 assert_eq!(turns.len(), 3);
1537 assert_eq!(turns[0].content, "first");
1538 assert_eq!(turns[2].content, "third");
1539 }
1540
1541 #[tokio::test]
1542 async fn session_turns_respect_limit_and_namespace() {
1543 let store = fresh().await;
1544 let ns = NamespaceId::new("ns_elliott_private");
1545 let other = NamespaceId::new("ns_other");
1546 let session = SessionId::new("sess_1");
1547
1548 for i in 0..5 {
1549 store
1550 .ingest_turn(
1551 &ns,
1552 SessionTurn {
1553 session_id: session.clone(),
1554 turn_index: i,
1555 role: TurnRole::User,
1556 content: format!("turn {i}"),
1557 timestamp: format!("2026-07-05T12:00:0{i}Z"),
1558 },
1559 )
1560 .await
1561 .unwrap();
1562 }
1563 store
1564 .ingest_turn(
1565 &other,
1566 SessionTurn {
1567 session_id: session.clone(),
1568 turn_index: 0,
1569 role: TurnRole::User,
1570 content: "intruder".into(),
1571 timestamp: "2026-07-05T12:00:00Z".into(),
1572 },
1573 )
1574 .await
1575 .unwrap();
1576
1577 let last_two = store.session_turns(&ns, &session, 2).await.expect("read");
1578 assert_eq!(last_two.len(), 2);
1579 assert_eq!(last_two[1].content, "turn 4");
1580 assert!(last_two.iter().all(|t| t.content != "intruder"));
1581 }
1582
1583 #[tokio::test]
1584 async fn session_metadata_create_list_end_round_trip() {
1585 let store = fresh().await;
1586 let ns = NamespaceId::new("ns_sessions");
1587
1588 let s1 = Session {
1590 id: SessionId::new("sess_a"),
1591 harness: Harness::Pi,
1592 channel: Some("thread-1".into()),
1593 started_at: "2026-07-05T10:00:00Z".into(),
1594 ended_at: None,
1595 };
1596 let s2 = Session {
1597 id: SessionId::new("sess_b"),
1598 harness: Harness::Sakamoto,
1599 channel: None,
1600 started_at: "2026-07-05T12:00:00Z".into(),
1601 ended_at: None,
1602 };
1603 store
1604 .create_session(&ns, s1.clone())
1605 .await
1606 .expect("create s1");
1607 store
1608 .create_session(&ns, s2.clone())
1609 .await
1610 .expect("create s2");
1611
1612 let all = store.list_sessions(&ns, None, 10).await.expect("list");
1614 assert_eq!(all.len(), 2);
1615 assert_eq!(all[0].id.as_str(), "sess_b");
1616 assert_eq!(all[1].id.as_str(), "sess_a");
1617 assert!(all[0].ended_at.is_none());
1618
1619 let pi_only = store
1621 .list_sessions(&ns, Some(&Harness::Pi), 10)
1622 .await
1623 .expect("list pi");
1624 assert_eq!(pi_only.len(), 1);
1625 assert_eq!(pi_only[0].id.as_str(), "sess_a");
1626
1627 store
1629 .end_session(
1630 &ns,
1631 &SessionId::new("sess_a"),
1632 "2026-07-05T11:30:00Z".into(),
1633 )
1634 .await
1635 .expect("end");
1636 let after_end = store.list_sessions(&ns, None, 10).await.expect("list");
1637 let s1_row = after_end
1638 .iter()
1639 .find(|s| s.id.as_str() == "sess_a")
1640 .unwrap();
1641 assert_eq!(s1_row.ended_at.as_deref(), Some("2026-07-05T11:30:00Z"));
1642
1643 store.create_session(&ns, s1.clone()).await.expect("upsert");
1645 let upserted = store
1646 .list_sessions(&ns, Some(&Harness::Pi), 10)
1647 .await
1648 .expect("list");
1649 assert_eq!(upserted.len(), 1);
1650
1651 let other = NamespaceId::new("ns_other");
1653 let cross = store
1654 .list_sessions(&other, None, 10)
1655 .await
1656 .expect("list other");
1657 assert!(cross.is_empty());
1658 }
1659
1660 #[tokio::test]
1663 async fn diaries_write_appends_and_reads_chronologically() {
1664 let store = fresh().await;
1665 let ns = NamespaceId::new("ns_diaries");
1666
1667 let e1 = DiaryEntry {
1668 agent: "claude".into(),
1669 content: "first entry".into(),
1670 topic: None,
1671 timestamp: String::new(),
1672 };
1673 let e2 = DiaryEntry {
1674 agent: "claude".into(),
1675 content: "second entry".into(),
1676 topic: Some("reflection".into()),
1677 timestamp: "2026-07-05T12:00:00Z".into(),
1678 };
1679 store.write_diary(&ns, e1).await.expect("write1");
1680 store.write_diary(&ns, e2).await.expect("write2");
1681
1682 let entries = store.read_diary(&ns, "claude", 10).await.expect("read");
1683 assert_eq!(entries.len(), 2);
1684 assert_eq!(entries[0].content, "first entry");
1686 assert_eq!(entries[1].content, "second entry");
1687 assert_eq!(entries[1].topic.as_deref(), Some("reflection"));
1688
1689 let other = store.read_diary(&ns, "pi", 10).await.expect("read");
1691 assert!(other.is_empty());
1692
1693 let other_ns = NamespaceId::new("ns_other");
1695 let cross = store
1696 .read_diary(&other_ns, "claude", 10)
1697 .await
1698 .expect("read");
1699 assert!(cross.is_empty());
1700 }
1701
1702 struct TestEmbedder;
1708 impl Embedder for TestEmbedder {
1709 fn dim(&self) -> usize {
1710 4
1711 }
1712 fn embed(&self, text: &str) -> Result<Embedding> {
1713 let words = text.split_whitespace().count() as f32;
1715 let has_surreal = text.contains("surreal") as u32 as f32;
1716 let has_rust = text.contains("rust") as u32 as f32;
1717 let has_memory = text.contains("memory") as u32 as f32;
1718 Ok(Embedding(vec![words, has_surreal, has_rust, has_memory]))
1719 }
1720 }
1721
1722 #[tokio::test]
1723 async fn search_without_embedder_errors() {
1724 let store = fresh().await; let ns = NamespaceId::new("ns_x");
1726 let result = store
1727 .search_memories(&ns, &Embedding(vec![0.0; 4]), 5)
1728 .await;
1729 assert!(matches!(result, Err(IjimaError::Store { .. })));
1730 }
1731
1732 #[tokio::test]
1733 async fn search_finds_nearest_with_embedder() {
1734 let store = SurrealStore::open_embedded_with(Arc::new(TestEmbedder))
1735 .await
1736 .expect("open with embedder");
1737 let ns = NamespaceId::new("ns_elliott_private");
1738
1739 store
1741 .store_memory(&ns, sample_memory("m1", "rust memory store"))
1742 .await
1743 .unwrap();
1744 store
1745 .store_memory(&ns, sample_memory("m2", "surreal db graph"))
1746 .await
1747 .unwrap();
1748 store
1749 .store_memory(&ns, sample_memory("m3", "completely unrelated text"))
1750 .await
1751 .unwrap();
1752
1753 let query = Embedding(vec![3.0, 0.0, 1.0, 1.0]);
1755 let hits = store.search_memories(&ns, &query, 2).await.expect("search");
1756 assert!(!hits.is_empty(), "must find at least one memory");
1757 assert_eq!(hits[0].memory.content, "rust memory store");
1760 assert!(hits[0].similarity > 0.0);
1762 }
1763
1764 #[tokio::test]
1765 async fn search_respects_namespace_isolation() {
1766 let store = SurrealStore::open_embedded_with(Arc::new(TestEmbedder))
1767 .await
1768 .expect("open with embedder");
1769 let alice = NamespaceId::new("ns_alice");
1770 let bob = NamespaceId::new("ns_bob");
1771
1772 store
1773 .store_memory(&alice, sample_memory("a1", "rust memory store"))
1774 .await
1775 .unwrap();
1776
1777 let query = Embedding(vec![3.0, 0.0, 1.0, 1.0]);
1780 let hits = store
1781 .search_memories(&bob, &query, 5)
1782 .await
1783 .expect("search");
1784 assert!(
1785 hits.is_empty(),
1786 "namespace isolation must hide alice's memory"
1787 );
1788 }
1789
1790 #[tokio::test]
1791 async fn list_memories_ranks_by_importance_then_recency() {
1792 let store = fresh().await;
1793 let ns = NamespaceId::new("ns_test");
1794 let mut hi = sample_memory("hi", "important");
1796 hi.importance = 0.9;
1797 hi.created_at = "100".into();
1798 let mut mid = sample_memory("mid", "medium");
1799 mid.importance = 0.5;
1800 mid.created_at = "200".into(); let mut lo = sample_memory("lo", "low");
1802 lo.importance = 0.9;
1803 lo.created_at = "300".into(); store.store_memory(&ns, hi).await.unwrap();
1805 store.store_memory(&ns, mid).await.unwrap();
1806 store.store_memory(&ns, lo).await.unwrap();
1807
1808 let list = store.list_memories(&ns, 10).await.expect("list");
1809 assert_eq!(list.len(), 3);
1812 assert_eq!(list[0].id.0, "lo"); assert_eq!(list[1].id.0, "hi"); assert_eq!(list[2].id.0, "mid"); }
1816
1817 #[tokio::test]
1820 async fn add_triple_then_query_entity() {
1821 let store = fresh().await;
1822 let ns = NamespaceId::new("ns_kg");
1823 let t = store
1824 .add_triple(
1825 &ns,
1826 EntityId::new("Ijima"),
1827 "depends_on",
1828 EntityId::new("SurrealDB"),
1829 Some("100"),
1830 1.0,
1831 None,
1832 )
1833 .await
1834 .expect("add_triple");
1835 assert_eq!(t.subject.as_str(), "Ijima");
1836 assert_eq!(t.predicate, "depends_on");
1837 assert_eq!(t.object.as_str(), "SurrealDB");
1838 assert!(t.valid_to.is_none()); let rec = store
1841 .query_entity(&ns, &EntityId::new("Ijima"))
1842 .await
1843 .expect("query");
1844 assert!(rec.entity.is_some());
1845 assert_eq!(rec.outgoing.len(), 1);
1846 assert_eq!(rec.outgoing[0].object.as_str(), "SurrealDB");
1847 assert!(rec.incoming.is_empty());
1848
1849 let rec = store
1851 .query_entity(&ns, &EntityId::new("SurrealDB"))
1852 .await
1853 .expect("query");
1854 assert_eq!(rec.incoming.len(), 1);
1855 assert!(rec.outgoing.is_empty());
1856 }
1857
1858 #[tokio::test]
1859 async fn invalidate_triple_sets_valid_to() {
1860 let store = fresh().await;
1861 let ns = NamespaceId::new("ns_kg");
1862 store
1863 .add_triple(
1864 &ns,
1865 EntityId::new("a"),
1866 "uses",
1867 EntityId::new("b"),
1868 Some("100"),
1869 1.0,
1870 None,
1871 )
1872 .await
1873 .unwrap();
1874 store
1875 .invalidate_triple(&ns, "a:uses:b")
1876 .await
1877 .expect("invalidate");
1878 let found = store
1879 .find_triples(&ns, None, Some("uses"), None)
1880 .await
1881 .expect("find");
1882 assert_eq!(found.len(), 1);
1883 assert!(found[0].valid_to.is_some(), "valid_to must be set");
1884 }
1885
1886 #[tokio::test]
1887 async fn knowledge_stats_counts_entities_and_triples() {
1888 let store = fresh().await;
1889 let ns = NamespaceId::new("ns_kg");
1890 store
1891 .add_triple(
1892 &ns,
1893 EntityId::new("a"),
1894 "x",
1895 EntityId::new("b"),
1896 None,
1897 1.0,
1898 None,
1899 )
1900 .await
1901 .unwrap();
1902 store
1903 .add_triple(
1904 &ns,
1905 EntityId::new("a"),
1906 "y",
1907 EntityId::new("c"),
1908 None,
1909 1.0,
1910 None,
1911 )
1912 .await
1913 .unwrap();
1914 let stats = store.knowledge_stats(&ns).await.expect("stats");
1915 assert_eq!(stats.entities, 3); assert_eq!(stats.triples, 2);
1917 }
1918
1919 #[tokio::test]
1920 async fn kg_namespace_isolation() {
1921 let store = fresh().await;
1922 let a = NamespaceId::new("ns_a");
1923 let b = NamespaceId::new("ns_b");
1924 store
1925 .add_triple(
1926 &a,
1927 EntityId::new("x"),
1928 "uses",
1929 EntityId::new("y"),
1930 None,
1931 1.0,
1932 None,
1933 )
1934 .await
1935 .unwrap();
1936 let stats = store.knowledge_stats(&b).await.expect("stats");
1938 assert_eq!(stats.triples, 0);
1939 let rec = store
1940 .query_entity(&b, &EntityId::new("x"))
1941 .await
1942 .expect("query");
1943 assert!(rec.outgoing.is_empty());
1944 }
1945
1946 #[tokio::test]
1947 async fn store_stats_counts_across_namespaces() {
1948 let store = fresh().await;
1949 store
1950 .store_memory(&NamespaceId::new("ns_a"), sample_memory("m1", "x"))
1951 .await
1952 .unwrap();
1953 store
1954 .store_memory(&NamespaceId::new("ns_a"), sample_memory("m2", "y"))
1955 .await
1956 .unwrap();
1957 store
1958 .store_memory(&NamespaceId::new("ns_b"), sample_memory("m3", "z"))
1959 .await
1960 .unwrap();
1961 let stats = store.store_stats().await.expect("stats");
1962 assert_eq!(stats.total_memories, 3);
1963 assert_eq!(stats.namespaces.len(), 2);
1964 let ns_a = stats
1965 .namespaces
1966 .iter()
1967 .find(|n| n.namespace == "ns_a")
1968 .expect("ns_a");
1969 assert_eq!(ns_a.memories, 2);
1970 }
1971
1972 fn mem_in(id: &str, project: &str, topic: &str, content: &str) -> Memory {
1974 Memory {
1975 id: MemoryId(id.into()),
1976 content: content.into(),
1977 project: project.into(),
1978 topic: topic.into(),
1979 source: MemorySource::Explicit,
1980 harness: Harness::Pi,
1981 session_id: None,
1982 origin: InstanceId::local(),
1983 authority: AuthorityScope::local(),
1984 importance: 0.5,
1985 created_at: "0".into(),
1986 }
1987 }
1988
1989 #[tokio::test]
1990 async fn palace_organization_rooms_taxonomy_graph_tunnel() {
1991 let store = fresh().await;
1992 let ns = NamespaceId::new("ns_palace");
1993 store
1996 .store_memory(&ns, mem_in("m1", "ijima", "auth", "use schubert"))
1997 .await
1998 .unwrap();
1999 store
2000 .store_memory(&ns, mem_in("m2", "ijima", "auth", "proof tokens"))
2001 .await
2002 .unwrap();
2003 store
2004 .store_memory(&ns, mem_in("m3", "ijima", "store", "surrealdb"))
2005 .await
2006 .unwrap();
2007 store
2008 .store_memory(&ns, mem_in("m4", "karpal", "auth", "workspace manifest"))
2009 .await
2010 .unwrap();
2011 store
2012 .store_memory(&ns, mem_in("m5", "karpal", "build", "ci"))
2013 .await
2014 .unwrap();
2015
2016 let rooms = store.list_rooms(&ns, None, 100).await.expect("rooms");
2018 assert_eq!(rooms.len(), 4);
2019 assert_eq!(rooms[0].count, 2); assert_eq!(rooms[0].project, "ijima");
2021 assert_eq!(rooms[0].topic, "auth");
2022
2023 let karpal_rooms = store
2025 .list_rooms(&ns, Some("karpal"), 100)
2026 .await
2027 .expect("rooms karpal");
2028 assert_eq!(karpal_rooms.len(), 2);
2029 assert!(karpal_rooms.iter().all(|r| r.project == "karpal"));
2030
2031 let taxons = store.taxonomy(&ns).await.expect("taxonomy");
2033 assert_eq!(taxons.len(), 2);
2034 let ijima_t = taxons
2035 .iter()
2036 .find(|t| t.project == "ijima")
2037 .expect("ijima taxon");
2038 assert_eq!(ijima_t.total, 3);
2039 assert_eq!(ijima_t.rooms.len(), 2); let graph = store.palace_graph(&ns).await.expect("graph");
2043 assert_eq!(graph.projects.len(), 2);
2044 assert!(graph.projects.contains(&"ijima".to_string()));
2045 assert!(graph.projects.contains(&"karpal".to_string()));
2046 let auth_tunnel = graph
2047 .tunnels
2048 .iter()
2049 .find(|t| t.topic == "auth")
2050 .expect("auth tunnel");
2051 assert_eq!(auth_tunnel.count_a, 2); assert_eq!(auth_tunnel.count_b, 1); let trav = store
2056 .traverse_tunnel(&ns, "auth", "ijima", "karpal", 10)
2057 .await
2058 .expect("traverse");
2059 assert_eq!(trav.memories_a.len(), 2);
2060 assert_eq!(trav.memories_b.len(), 1);
2061 assert!(trav.memories_a.iter().all(|m| m.project == "ijima"));
2062 assert!(trav.memories_b.iter().all(|m| m.project == "karpal"));
2063 }
2064
2065 #[tokio::test]
2066 async fn mining_queue_enqueue_list_accept_reject() {
2067 let store = fresh().await;
2068 let ns = NamespaceId::new("ns_mining");
2069 let other = NamespaceId::new("ns_other");
2070
2071 let q1 = store
2073 .enqueue_extraction(&ns, sample_memory("m1", "decided to use surrealdb"), 0.6)
2074 .await
2075 .expect("enqueue1");
2076 let q2 = store
2077 .enqueue_extraction(&ns, sample_memory("m2", "see https://example.com"), 0.55)
2078 .await
2079 .expect("enqueue2");
2080
2081 let pending = store.list_pending(&ns, 10).await.expect("list");
2083 assert_eq!(pending.len(), 2);
2084 assert_eq!(pending[0].id, q2); assert_eq!(pending[1].confidence, 0.6);
2086
2087 let cross = store.list_pending(&other, 10).await.expect("list other");
2089 assert!(cross.is_empty());
2090
2091 let accepted = store.accept_extraction(&ns, &q1).await.expect("accept");
2093 assert_eq!(accepted.memory_id.0, "m1");
2094 let after_accept = store.list_pending(&ns, 10).await.expect("list");
2095 assert_eq!(after_accept.len(), 1);
2096 let promoted = store
2098 .recall_memory(&ns, &MemoryId("m1".into()))
2099 .await
2100 .expect("recall");
2101 assert!(promoted.is_some());
2102
2103 store.reject_extraction(&ns, &q2).await.expect("reject");
2105 let after_reject = store.list_pending(&ns, 10).await.expect("list");
2106 assert!(after_reject.is_empty());
2107 let not_promoted = store
2108 .recall_memory(&ns, &MemoryId("m2".into()))
2109 .await
2110 .expect("recall");
2111 assert!(not_promoted.is_none());
2112
2113 let q3 = store
2115 .enqueue_extraction(&ns, sample_memory("m3", "cross test"), 0.5)
2116 .await
2117 .expect("enqueue3");
2118 let cross_accept = store.accept_extraction(&other, &q3).await;
2119 assert!(cross_accept.is_err());
2120 }
2121
2122 #[tokio::test]
2123 async fn persistent_store_survives_reopen() {
2124 let dir = std::env::temp_dir().join(format!(
2125 "ijima-persist-test-{}",
2126 std::time::SystemTime::now()
2127 .duration_since(std::time::UNIX_EPOCH)
2128 .unwrap()
2129 .as_nanos()
2130 ));
2131 let ns = NamespaceId::new("ns_persist");
2132
2133 {
2135 let store = SurrealStore::open_persistent(&dir).await.expect("open");
2136 store
2137 .store_memory(&ns, sample_memory("mem_p", "survives restart"))
2138 .await
2139 .expect("store");
2140 }
2141
2142 let store = SurrealStore::open_persistent(&dir).await.expect("reopen");
2144 let got = store
2145 .recall_memory(&ns, &MemoryId("mem_p".into()))
2146 .await
2147 .expect("recall")
2148 .expect("memory must survive reopen");
2149 assert_eq!(got.content, "survives restart");
2150 }
2151
2152 #[tokio::test]
2153 async fn export_to_writes_sql_dump_to_file() {
2154 let dir = std::env::temp_dir().join(format!(
2155 "ijima-export-test-{}",
2156 std::time::SystemTime::now()
2157 .duration_since(std::time::UNIX_EPOCH)
2158 .unwrap()
2159 .as_nanos()
2160 ));
2161 std::fs::create_dir_all(&dir).expect("create dir");
2162 let db_path = dir.join("ijima.db");
2163 let store = crate::SurrealStore::open_persistent(&db_path)
2164 .await
2165 .expect("open");
2166 let ns = NamespaceId::new("ns_export");
2167 store
2168 .store_memory(&ns, sample_memory("mem_x", "export this memory"))
2169 .await
2170 .expect("store");
2171
2172 let out = dir.join("dump.surql");
2173 store.export_to(&out).await.expect("export");
2174 let contents = std::fs::read_to_string(&out).expect("read");
2175 assert!(contents.contains("export this memory"));
2176 let _ = std::fs::remove_dir_all(&dir);
2178 }
2179}