Skip to main content

ijima_server/
backend_surreal.rs

1// Copyright (C) 2026 Industrial Algebra
2// SPDX-License-Identifier: Apache-2.0
3
4//! SurrealDB store backend — Ijima's primary backend (D6).
5//!
6//! Implements [`ijima_core::Store`] over an embedded SurrealDB instance
7//! (`kv-mem` engine). The same impl will back a server-mode deployment
8//! once the engine type is widened (future).
9//!
10//! ## Vector search
11//!
12//! When constructed with [`SurrealStore::open_embedded_with`] (an
13//! [`Embedder`]), each stored memory is embedded at write time and
14//! indexed via **Cosine similarity** ranking (`vector::similarity::cosine`),
15//! brute-force over the namespace (an **HNSW** index is the planned
16//! optimization once the SDK can emit typed `<N>f32` vectors; MTREE is
17//! deprecated in SurrealDB 2.6). `search_memories`
18//!
19//! ## Namespace handling (v0)
20//!
21//! Ijima's [`NamespaceId`](ijima_core::NamespaceId) is stored as a record
22//! field and every query filters by it, so isolation is enforced
23//! per-query under concurrent access. A future promotion maps each
24//! namespace onto a native SurrealDB `NS`/`DB` scope once the
25//! concurrency semantics of per-query `USE NS/DB` are validated.
26
27use 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
41/// The SurrealDB namespace/database the Ijima instance lives in.
42const 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";
50/// Entity nodes (knowledge-graph).
51const ENTITIES_TABLE: &str = "entities";
52/// Triple edges (knowledge-graph).
53const TRIPLES_TABLE: &str = "triples";
54/// Repo directory (global Context Mapper registry).
55const REPO_TABLE: &str = "repo_directory";
56
57/// A SurrealDB-backed [`Store`].
58pub struct SurrealStore {
59    db: Surreal<Db>,
60    /// When present, memories are embedded at write time and
61    /// [`Store::search_memories`] is available. When absent, search
62    /// returns [`IjimaError::Store`].
63    embedder: Option<Arc<dyn Embedder>>,
64}
65
66impl SurrealStore {
67    /// Opens an **in-memory** embedded store (the `Mem` engine), with no
68    /// embedder. Use for tests; data does not survive restart.
69    ///
70    /// # Errors
71    ///
72    /// Returns [`IjimaError::Store`] if SurrealDB cannot initialize.
73    pub async fn open_embedded() -> Result<Self> {
74        Self::open_with_db(new_mem().await?, None).await
75    }
76
77    /// Opens an **in-memory** embedded store with an [`Embedder`].
78    ///
79    /// # Errors
80    ///
81    /// Returns [`IjimaError::Store`] if SurrealDB cannot initialize.
82    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    /// Opens a **persistent** store (the `SurrealKv` engine) at `path`,
87    /// with no embedder. Data survives restart. Creates the directory if
88    /// absent.
89    ///
90    /// # Errors
91    ///
92    /// Returns [`IjimaError::Store`] if SurrealDB cannot initialize or
93    /// the path is unwritable.
94    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    /// Opens a **persistent** store with an [`Embedder`].
99    ///
100    /// # Errors
101    ///
102    /// Returns [`IjimaError::Store`] if SurrealDB cannot initialize or
103    /// the path is unwritable.
104    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        // Define indexes on the hot query columns. SurrealDB is schemaless
119        // by default and does NOT auto-index fields, so without these every
120        // namespace-scoped SELECT and every dedup check is a full table scan
121        // (O(N²) for a bulk import). `IF NOT EXISTS` makes this idempotent
122        // on existing stores.
123        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    /// The `DEFINE INDEX` statements run at store open. Covers the columns
132    /// every query filters on: `namespace` (all scoped reads),
133    /// `(namespace, content_hash)` (dedup), `(namespace, session_id)`
134    /// (session turns), and the knowledge-graph traversal keys.
135    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    /// Fetches all `(project, topic)` cells in `ns` with their memory counts.
149    /// One query, aggregated in Rust (consistent with `store_stats`).
150    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    /// Returns memories in `ns` matching `project` + `topic`, ranked by
176    /// importance desc then recency desc (mirrors `list_memories`).
177    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    /// Exports the entire store as a SurrealDB SQL dump to `path`.
203    /// Requires a persistent backend (SurrealKv); in-memory stores do not
204    /// support Backup.
205    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// ---------- wire records ----------
239
240/// The persisted form of a [`Memory`] plus its owning namespace and an
241/// optional embedding vector.
242#[derive(Debug, Clone, Serialize, Deserialize)]
243struct MemoryRecord {
244    /// The [`MemoryId`] string (mirrors the SurrealDB record id so search
245    /// results can return it without parsing a RecordId).
246    memory_id: String,
247    content: String,
248    /// SHA-256 of `content` (content-hash dedup, Phase 2.2). Stored for
249    /// O(1) exact-duplicate lookup within a namespace. `#[serde(default)]`
250    /// so SELECTs that omit it (list/search projections) still deserialize.
251    #[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    /// Provenance: the authoring instance (ADR provenance-tier). Defaults
264    /// to local for documents written before the field existed.
265    #[serde(default)]
266    origin: InstanceId,
267    /// Provenance: the authority scope (source-of-truth) for the record's
268    /// domain (ADR provenance-tier). Defaults to local.
269    #[serde(default)]
270    authority: AuthorityScope,
271    /// Which embedding model produced `embedding` (D10 provenance), e.g.
272    /// `sentence-transformers/all-MiniLM-L6-v2@main`. Absent when no
273    /// embedder is configured.
274    #[serde(default, skip_serializing_if = "Option::is_none")]
275    embed_model: Option<String>,
276    /// Embedding vector (present when the store was opened with an
277    /// [`Embedder`]). `#[serde(default)]` so namespace-filtered selects
278    /// that omit it still deserialize.
279    #[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/// The persisted form of a [`SessionTurn`] plus its owning namespace.
333#[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/// Stored row for a diary entry (Phase 3.3). `ts` is an internal
344/// epoch-millis field for correct numeric ORDER BY (string timestamps
345/// don't sort lexicographically across formats).
346#[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/// Stored row for a session's metadata (Phase 2.3).
380#[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/// Stored row for a queued mining extraction (ADR M2). `ts` is an internal
416/// epoch-millis for correct ORDER BY.
417#[derive(Debug, Clone, Serialize, Deserialize)]
418struct QueueRecord {
419    queue_id: String,
420    /// Nested memory fields (flattened for storage).
421    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
506/// Current epoch-millis (for internal ORDER BY fields).
507fn 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
514/// Lowercase hex of a byte slice (for content-hash).
515fn 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        // Content-hash dedup (Phase 2.2): reject exact duplicates within
527        // the namespace, returning the existing id so callers can recover.
528        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                // Stamp the model id (D10 provenance) so a future model
537                // swap is detectable and a re-embed pass can be triggered.
538                (
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        // Verify ownership before deleting (isolation).
590        let existing = self.recall_memory(ns, id).await?;
591        if existing.is_none() {
592            return Ok(()); // absent or wrong namespace — nothing to delete
593        }
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        // One query for all namespaces; aggregate in Rust (avoids
622        // SurrealDB aggregate-function ambiguity, fine at v0 scale).
623        #[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        // count desc, then topic for determinism.
668        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        // group by project
676        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        // projects with the most memories first.
701        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        // topic -> set of (project, count)
712        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; // a tunnel needs two distinct projects
723            }
724            entries.sort_by(|a, b| a.0.cmp(&b.0));
725            // every unordered pair of distinct projects on this topic
726            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        // Cosine similarity ranking, brute-force over the namespace.
779        // Correct (no ANN approximation); HNSW is the planned
780        // optimization once the SDK emits typed `<N>f32` vectors.
781        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        // Deserialize each row as a MemoryRecord (missing fields default)
798        // flattened alongside the computed `score`.
799        #[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(); // query was DESC for "last N"; return chronological
854        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        // Upsert: create if absent, replace metadata if present. Session
870        // ids are globally unique (same convention as memory ids).
871        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        // Scoped by namespace + session id (safety: a principal can only
916        // end sessions in their own namespace).
917        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        // Read the queued record (scoped by namespace), promote to the
973        // palace, then delete from the queue.
974        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        // store_memory does content-hash dedup; a duplicate promote is a no-op.
992        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        // Scoped delete: only removes if the namespace matches.
1003        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(); // chronological (DESC → reverse)
1051        Ok(records.into_iter().map(DiaryRecord::into_entry).collect())
1052    }
1053
1054    // ===== Repo directory (global registry — Context Mapper) =====
1055
1056    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// ---------- knowledge graph ----------
1079
1080#[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    /// Deterministic id: `<subject>:<predicate>:<object>`.
1090    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        // Create both entity nodes (idempotent — ignore already-exists).
1137        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            // Ignore errors (record already exists from a prior triple).
1149        }
1150        // Store the triple as a record (keyed by deterministic triple_id).
1151        // Native graph edges (RELATE) are a future optimization for
1152        // multi-hop traversal; v0 queries via subject/object fields.
1153        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        // Provenance fields persist through a store→recall round-trip
1399        // (ADR provenance-tier): origin/authority survive SurrealDB.
1400        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        // Same content, different id — must be rejected.
1413        let result = store
1414            .store_memory(&ns, sample_memory("m2", "identical content"))
1415            .await;
1416        assert!(matches!(result, Err(IjimaError::Duplicate { .. })));
1417
1418        // check_duplicate finds the existing id.
1419        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        // Different content stores fine.
1426        store
1427            .store_memory(&ns, sample_memory("m3", "different content"))
1428            .await
1429            .expect("distinct content stores");
1430
1431        // Same content in a DIFFERENT namespace is not a duplicate.
1432        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        // SurrealDB is schemaless with no auto-indexes; open_with_db defines
1442        // them. Re-running the DDL (as on every store open against an
1443        // existing database) must not error — IF NOT EXISTS guards it. This
1444        // is the fix for the O(N²) full-scan migration hang.
1445        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        // Dedup is still correct after indexing (namespace + content_hash).
1452        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        // Different namespace, same content: not a duplicate (composite index).
1463        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        // Create two sessions (different harnesses).
1589        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        // List all — newest first (s2 has later started_at).
1613        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        // Filter by harness.
1620        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        // End s1.
1628        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        // Upsert (re-create s1) preserves metadata shape.
1644        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        // Namespace isolation: sessions in another namespace are invisible.
1652        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    // ===== Diaries =====
1661
1662    #[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        // chronological: e1 before e2
1685        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        // Agent isolation: different agent sees nothing.
1690        let other = store.read_diary(&ns, "pi", 10).await.expect("read");
1691        assert!(other.is_empty());
1692
1693        // Namespace isolation: different ns sees nothing.
1694        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    // ===== Vector search =====
1703
1704    /// A deterministic test embedder: maps text to a small fixed-dim
1705    /// vector so nearest-neighbour behaviour is predictable. Production
1706    /// uses candle + all-MiniLM-L6-v2 (384-dim).
1707    struct TestEmbedder;
1708    impl Embedder for TestEmbedder {
1709        fn dim(&self) -> usize {
1710            4
1711        }
1712        fn embed(&self, text: &str) -> Result<Embedding> {
1713            // Cosine-relevant signal: word-count + keyword flags.
1714            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; // no embedder
1725        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        // Three memories with distinct keyword signals.
1740        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        // Query close to the "rust memory" vector.
1754        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        // The nearest hit must be the rust/memory one (m1), not the
1758        // unrelated text.
1759        assert_eq!(hits[0].memory.content, "rust memory store");
1760        // Scored search: each hit carries its cosine similarity.
1761        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        // Bob searches with an identical query embedding — must see
1778        // nothing from alice's namespace.
1779        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        // Three memories with varying importance + recency.
1795        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(); // newer but lower importance
1801        let mut lo = sample_memory("lo", "low");
1802        lo.importance = 0.9;
1803        lo.created_at = "300".into(); // same importance as hi, newer
1804        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        // importance DESC first: the two 0.9s before the 0.5.
1810        // Among the 0.9s, created_at DESC: lo (300) before hi (100).
1811        assert_eq!(list.len(), 3);
1812        assert_eq!(list[0].id.0, "lo"); // 0.9, 300
1813        assert_eq!(list[1].id.0, "hi"); // 0.9, 100
1814        assert_eq!(list[2].id.0, "mid"); // 0.5
1815    }
1816
1817    // ===== knowledge graph =====
1818
1819    #[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()); // current
1839
1840        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        // From the object's side, it's incoming.
1850        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); // a, b, c
1916        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        // b sees nothing.
1937        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    /// Builds a memory with explicit project/topic (sample_memory hardcodes them).
1973    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        // Two projects sharing the topic "auth" (a tunnel), plus a
1994        // project-only topic.
1995        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        // list_rooms — all rooms, count desc.
2017        let rooms = store.list_rooms(&ns, None, 100).await.expect("rooms");
2018        assert_eq!(rooms.len(), 4);
2019        assert_eq!(rooms[0].count, 2); // ijima/auth
2020        assert_eq!(rooms[0].project, "ijima");
2021        assert_eq!(rooms[0].topic, "auth");
2022
2023        // list_rooms filtered to karpal.
2024        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        // taxonomy — two projects.
2032        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); // auth, store
2040
2041        // palace_graph — projects + the auth tunnel between ijima & karpal.
2042        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); // ijima
2052        assert_eq!(auth_tunnel.count_b, 1); // karpal
2053
2054        // traverse_tunnel — memories from both sides.
2055        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        // Enqueue two PendingReview extractions.
2072        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        // list_pending sees both (newest first).
2082        let pending = store.list_pending(&ns, 10).await.expect("list");
2083        assert_eq!(pending.len(), 2);
2084        assert_eq!(pending[0].id, q2); // newest first
2085        assert_eq!(pending[1].confidence, 0.6);
2086
2087        // Namespace isolation: other ns sees nothing.
2088        let cross = store.list_pending(&other, 10).await.expect("list other");
2089        assert!(cross.is_empty());
2090
2091        // Accept q1 → promotes to palace + removes from queue.
2092        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        // The promoted memory is now in the palace.
2097        let promoted = store
2098            .recall_memory(&ns, &MemoryId("m1".into()))
2099            .await
2100            .expect("recall");
2101        assert!(promoted.is_some());
2102
2103        // Reject q2 → drops without promoting.
2104        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        // Cross-namespace accept fails (queue entry is in ns, not other).
2114        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        // Write with one instance.
2134        {
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        // A fresh instance pointing at the same path must see the memory.
2143        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        // Clean up.
2177        let _ = std::fs::remove_dir_all(&dir);
2178    }
2179}