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    NamespaceMembership, PalaceGraph, ProjectTaxon, QueuedExtraction, RepoDirectory, Result, Room,
34    SearchHit, Session, SessionId, SessionTurn, Store, StoreStats, TokenRevocation, Triple, Tunnel,
35    TunnelTraversal, embeddings::Embedder, harness::Harness, memory::MemorySource,
36};
37use serde::{Deserialize, Serialize};
38use surrealdb::Surreal;
39use surrealdb::engine::local::{Db, Mem, SurrealKv};
40
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";
50const REVOCATIONS_TABLE: &str = "token_revocations";
51const NAMESPACE_MEMBERS_TABLE: &str = "namespace_members";
52/// Entity nodes (knowledge-graph).
53const ENTITIES_TABLE: &str = "entities";
54/// Triple edges (knowledge-graph).
55const TRIPLES_TABLE: &str = "triples";
56/// Repo directory (global Context Mapper registry).
57const REPO_TABLE: &str = "repo_directory";
58
59/// A SurrealDB-backed [`Store`].
60pub struct SurrealStore {
61    db: Surreal<Db>,
62    /// When present, memories are embedded at write time and
63    /// [`Store::search_memories`] is available. When absent, search
64    /// returns [`IjimaError::Store`].
65    embedder: Option<Arc<dyn Embedder>>,
66}
67
68impl SurrealStore {
69    /// Opens an **in-memory** embedded store (the `Mem` engine), with no
70    /// embedder. Use for tests; data does not survive restart.
71    ///
72    /// # Errors
73    ///
74    /// Returns [`IjimaError::Store`] if SurrealDB cannot initialize.
75    pub async fn open_embedded() -> Result<Self> {
76        Self::open_with_db(new_mem().await?, None).await
77    }
78
79    /// Opens an **in-memory** embedded store with an [`Embedder`].
80    ///
81    /// # Errors
82    ///
83    /// Returns [`IjimaError::Store`] if SurrealDB cannot initialize.
84    pub async fn open_embedded_with(embedder: Arc<dyn Embedder>) -> Result<Self> {
85        Self::open_with_db(new_mem().await?, Some(embedder)).await
86    }
87
88    /// Opens a **persistent** store (the `SurrealKv` engine) at `path`,
89    /// with no embedder. Data survives restart. Creates the directory if
90    /// absent.
91    ///
92    /// # Locking (surrealdb 3)
93    ///
94    /// The engine holds a directory `LOCK` file while open. Dropping the
95    /// store handle does not release it synchronously — the datastore's
96    /// background tasks must wind down first (a task boundary + yield).
97    /// Cross-process restarts are unaffected (the OS releases the lock at
98    /// process exit); same-process sequential reopens of the same path
99    /// need the spawn-and-yield pattern (see the persistence test).
100    ///
101    /// # Errors
102    ///
103    /// Returns [`IjimaError::Store`] if SurrealDB cannot initialize or
104    /// the path is unwritable.
105    pub async fn open_persistent(path: impl AsRef<std::path::Path>) -> Result<Self> {
106        Self::open_with_db(new_surrealkv(&path).await?, None).await
107    }
108
109    /// Opens a **persistent** store with an [`Embedder`].
110    ///
111    /// # Errors
112    ///
113    /// Returns [`IjimaError::Store`] if SurrealDB cannot initialize or
114    /// the path is unwritable.
115    pub async fn open_persistent_with(
116        path: impl AsRef<std::path::Path>,
117        embedder: Arc<dyn Embedder>,
118    ) -> Result<Self> {
119        Self::open_with_db(new_surrealkv(&path).await?, Some(embedder)).await
120    }
121
122    async fn open_with_db(db: Surreal<Db>, embedder: Option<Arc<dyn Embedder>>) -> Result<Self> {
123        db.use_ns(SURREAL_NS)
124            .use_db(SURREAL_DB)
125            .await
126            .map_err(|e| IjimaError::Store {
127                detail: format!("surrealdb use_ns/use_db: {e}"),
128            })?;
129        // Define indexes on the hot query columns. SurrealDB is schemaless
130        // by default and does NOT auto-index fields, so without these every
131        // namespace-scoped SELECT and every dedup check is a full table scan
132        // (O(N²) for a bulk import). `IF NOT EXISTS` makes this idempotent
133        // on existing stores.
134        db.query(Self::INDEX_DDL)
135            .await
136            .map_err(|e| IjimaError::Store {
137                detail: format!("surrealdb define indexes: {e}"),
138            })?;
139        Ok(Self { db, embedder })
140    }
141
142    /// The `DEFINE INDEX` statements run at store open. Covers the columns
143    /// every query filters on: `namespace` (all scoped reads),
144    /// `(namespace, content_hash)` (dedup), `(namespace, session_id)`
145    /// (session turns), and the knowledge-graph traversal keys.
146    ///
147    /// The `DEFINE TABLE` statements matter on surrealdb 3: `SELECT FROM
148    /// <table>` on a table that has never been written is a hard error
149    /// (v2 returned an empty set). Boot-time reads — revocation
150    /// hydration, `/status` counts — hit fresh databases, so every table
151    /// is defined up front (schemaless, `IF NOT EXISTS` idempotent).
152    const INDEX_DDL: &str = r#"
153        DEFINE TABLE IF NOT EXISTS memories          SCHEMALESS;
154        DEFINE TABLE IF NOT EXISTS session_turns     SCHEMALESS;
155        DEFINE TABLE IF NOT EXISTS sessions          SCHEMALESS;
156        DEFINE TABLE IF NOT EXISTS diaries           SCHEMALESS;
157        DEFINE TABLE IF NOT EXISTS mining_queue      SCHEMALESS;
158        DEFINE TABLE IF NOT EXISTS entities          SCHEMALESS;
159        DEFINE TABLE IF NOT EXISTS triples           SCHEMALESS;
160        DEFINE TABLE IF NOT EXISTS token_revocations SCHEMALESS;
161        DEFINE TABLE IF NOT EXISTS namespace_members SCHEMALESS;
162        DEFINE INDEX IF NOT EXISTS nmem_ns      ON TABLE namespace_members FIELDS namespace;
163        DEFINE INDEX IF NOT EXISTS nmem_ns_princ ON TABLE namespace_members FIELDS namespace, principal;
164        DEFINE INDEX IF NOT EXISTS mem_ns        ON TABLE memories       FIELDS namespace;
165        DEFINE INDEX IF NOT EXISTS mem_ns_hash   ON TABLE memories       FIELDS namespace, content_hash;
166        DEFINE INDEX IF NOT EXISTS turns_ns_sess ON TABLE session_turns   FIELDS namespace, session_id;
167        DEFINE INDEX IF NOT EXISTS sess_ns       ON TABLE sessions        FIELDS namespace;
168        DEFINE INDEX IF NOT EXISTS diary_ns_ag   ON TABLE diaries         FIELDS namespace, agent;
169        DEFINE INDEX IF NOT EXISTS queue_ns      ON TABLE mining_queue    FIELDS namespace;
170        DEFINE INDEX IF NOT EXISTS ent_ns        ON TABLE entities        FIELDS namespace;
171        DEFINE INDEX IF NOT EXISTS trip_ns       ON TABLE triples         FIELDS namespace;
172        DEFINE INDEX IF NOT EXISTS trip_subj     ON TABLE triples         FIELDS subject;
173        DEFINE INDEX IF NOT EXISTS trip_obj      ON TABLE triples         FIELDS object;
174    "#;
175
176    /// Fetches all `(project, topic)` cells in `ns` with their memory counts.
177    /// One query, aggregated in Rust (consistent with `store_stats`).
178    async fn project_topic_counts(
179        &self,
180        ns: &NamespaceId,
181    ) -> Result<std::collections::BTreeMap<(String, String), usize>> {
182        #[derive(serde::Serialize, Deserialize)]
183        struct Row {
184            project: String,
185            topic: String,
186        }
187        let mut result = self
188            .db
189            .query(format!(
190                "SELECT project, topic FROM {MEMORIES_TABLE} WHERE namespace = $ns"
191            ))
192            .bind(("ns", ns.as_str().to_string()))
193            .await
194            .map_err(store_err)?;
195        let rows = take_vec::<Row>(&mut result)?;
196        let mut counts = std::collections::BTreeMap::new();
197        for r in rows {
198            *counts.entry((r.project, r.topic)).or_insert(0) += 1;
199        }
200        Ok(counts)
201    }
202
203    /// Returns memories in `ns` matching `project` + `topic`, ranked by
204    /// importance desc then recency desc (mirrors `list_memories`).
205    async fn project_topic_memories(
206        &self,
207        ns: &NamespaceId,
208        project: &str,
209        topic: &str,
210        limit: usize,
211    ) -> Result<Vec<Memory>> {
212        let mut result = self
213            .db
214            .query(format!(
215                "SELECT memory_id, content, project, topic, source, harness, session_id, namespace,
216                        origin, authority, importance, created_at
217                 FROM {MEMORIES_TABLE}
218                 WHERE namespace = $ns AND project = $proj AND topic = $topic
219                 ORDER BY importance DESC, created_at DESC LIMIT $lim"
220            ))
221            .bind(("ns", ns.as_str().to_string()))
222            .bind(("proj", project.to_string()))
223            .bind(("topic", topic.to_string()))
224            .bind(("lim", limit as i64))
225            .await
226            .map_err(store_err)?;
227        let records = take_vec::<MemoryRecord>(&mut result)?;
228        Ok(records.into_iter().map(|r| r.into_memory()).collect())
229    }
230
231    /// Exports the entire store as a SurrealDB SQL dump to `path`.
232    /// Requires a persistent backend (SurrealKv); in-memory stores do not
233    /// support Backup.
234    pub async fn export_to(&self, path: impl AsRef<std::path::Path>) -> Result<()> {
235        self.db
236            .export(path.as_ref())
237            .await
238            .map_err(|e| IjimaError::Store {
239                detail: format!("export: {e}"),
240            })?;
241        Ok(())
242    }
243}
244
245async fn new_mem() -> Result<Surreal<Db>> {
246    Surreal::new::<Mem>(())
247        .await
248        .map_err(|e| IjimaError::Store {
249            detail: format!("surrealdb mem init: {e}"),
250        })
251}
252
253async fn new_surrealkv(path: impl AsRef<std::path::Path>) -> Result<Surreal<Db>> {
254    let path = path.as_ref();
255    if let Some(parent) = path.parent() {
256        std::fs::create_dir_all(parent).map_err(|e| IjimaError::Store {
257            detail: format!("mkdir {}: {e}", parent.display()),
258        })?;
259    }
260    Surreal::new::<SurrealKv>(path.to_string_lossy().to_string())
261        .await
262        .map_err(|e| IjimaError::Store {
263            detail: format!("surrealdb surrealkv init at {}: {e}", path.display()),
264        })
265}
266
267// ---------- surrealdb 3.x serde bridge ----------
268// surrealdb 3's `take`/`content` are typed over `SurrealValue`. Our wire
269// records embed ijima-core provenance types (MemorySource, Harness,
270// InstanceId…) that are serde-only — implementing SurrealValue for them
271// would drag a surrealdb dependency into ijima-core. `SerdeWrapper` is
272// the SDK's official bridge for exactly this case ("types where an
273// implementation of SurrealValue isn't available or practical").
274
275/// Takes query result row 0 as `Vec<T>` through the serde bridge.
276fn take_vec<T>(result: &mut surrealdb::IndexedResults) -> Result<Vec<T>>
277where
278    T: serde::Serialize + serde::de::DeserializeOwned + 'static,
279{
280    let wrapped: Vec<surrealdb::types::SerdeWrapper<T>> = result.take(0).map_err(store_err)?;
281    Ok(wrapped.into_iter().map(|w| w.0).collect())
282}
283
284// ---------- wire records ----------
285
286/// The persisted form of a [`Memory`] plus its owning namespace and an
287/// optional embedding vector.
288#[derive(Debug, Clone, Serialize, Deserialize)]
289struct MemoryRecord {
290    /// The [`MemoryId`] string (mirrors the SurrealDB record id so search
291    /// results can return it without parsing a RecordId).
292    memory_id: String,
293    content: String,
294    /// SHA-256 of `content` (content-hash dedup, Phase 2.2). Stored for
295    /// O(1) exact-duplicate lookup within a namespace. `#[serde(default)]`
296    /// so SELECTs that omit it (list/search projections) still deserialize.
297    #[serde(default)]
298    content_hash: String,
299    project: String,
300    topic: String,
301    source: MemorySource,
302    harness: Harness,
303    session_id: Option<String>,
304    namespace: String,
305    #[serde(default = "default_record_importance")]
306    importance: f32,
307    #[serde(default)]
308    created_at: String,
309    /// Provenance: the authoring instance (ADR provenance-tier). Defaults
310    /// to local for documents written before the field existed.
311    #[serde(default)]
312    origin: InstanceId,
313    /// Provenance: the authority scope (source-of-truth) for the record's
314    /// domain (ADR provenance-tier). Defaults to local.
315    #[serde(default)]
316    authority: AuthorityScope,
317    /// Which embedding model produced `embedding` (D10 provenance), e.g.
318    /// `sentence-transformers/all-MiniLM-L6-v2@main`. Absent when no
319    /// embedder is configured.
320    #[serde(default, skip_serializing_if = "Option::is_none")]
321    embed_model: Option<String>,
322    /// Embedding vector (present when the store was opened with an
323    /// [`Embedder`]). `#[serde(default)]` so namespace-filtered selects
324    /// that omit it still deserialize.
325    #[serde(default, skip_serializing_if = "Option::is_none")]
326    embedding: Option<Vec<f32>>,
327}
328
329fn default_record_importance() -> f32 {
330    0.5
331}
332
333impl MemoryRecord {
334    fn from_memory(
335        memory: &Memory,
336        ns: &NamespaceId,
337        embedding: Option<Vec<f32>>,
338        embed_model: Option<String>,
339    ) -> Self {
340        use sha2::{Digest, Sha256};
341        let content_hash = hex(&Sha256::digest(memory.content.as_bytes()));
342        Self {
343            memory_id: memory.id.0.clone(),
344            content: memory.content.clone(),
345            content_hash,
346            project: memory.project.clone(),
347            topic: memory.topic.clone(),
348            source: memory.source,
349            harness: memory.harness,
350            session_id: memory.session_id.clone(),
351            namespace: ns.as_str().to_string(),
352            importance: memory.importance,
353            created_at: memory.created_at.clone(),
354            origin: memory.origin.clone(),
355            authority: memory.authority.clone(),
356            embed_model,
357            embedding,
358        }
359    }
360
361    fn into_memory(self) -> Memory {
362        Memory {
363            id: MemoryId(self.memory_id),
364            content: self.content,
365            project: self.project,
366            topic: self.topic,
367            source: self.source,
368            harness: self.harness,
369            session_id: self.session_id,
370            origin: self.origin,
371            authority: self.authority,
372            importance: self.importance,
373            created_at: self.created_at,
374        }
375    }
376}
377
378/// The persisted form of a [`SessionTurn`] plus its owning namespace.
379#[derive(Debug, Clone, Serialize, Deserialize)]
380struct SessionTurnRecord {
381    session_id: String,
382    turn_index: u64,
383    role: ijima_core::TurnRole,
384    content: String,
385    timestamp: String,
386    namespace: String,
387}
388
389/// Stored row for a diary entry (Phase 3.3). `ts` is an internal
390/// epoch-millis field for correct numeric ORDER BY (string timestamps
391/// don't sort lexicographically across formats).
392#[derive(Debug, Clone, Serialize, Deserialize)]
393struct DiaryRecord {
394    agent: String,
395    content: String,
396    #[serde(default, skip_serializing_if = "Option::is_none")]
397    topic: Option<String>,
398    timestamp: String,
399    ts: i64,
400    namespace: String,
401}
402
403impl DiaryRecord {
404    fn from_entry(entry: &ijima_core::DiaryEntry, ns: &NamespaceId, now_ms: i64) -> Self {
405        Self {
406            agent: entry.agent.clone(),
407            content: entry.content.clone(),
408            topic: entry.topic.clone(),
409            timestamp: entry.timestamp.clone(),
410            ts: now_ms,
411            namespace: ns.as_str().to_string(),
412        }
413    }
414
415    fn into_entry(self) -> ijima_core::DiaryEntry {
416        ijima_core::DiaryEntry {
417            agent: self.agent,
418            content: self.content,
419            topic: self.topic,
420            timestamp: self.timestamp,
421        }
422    }
423}
424
425/// Stored row for a session's metadata (Phase 2.3).
426#[derive(Debug, Clone, Serialize, Deserialize)]
427struct SessionRecord {
428    session_id: String,
429    harness: String,
430    #[serde(default, skip_serializing_if = "Option::is_none")]
431    channel: Option<String>,
432    started_at: String,
433    #[serde(default, skip_serializing_if = "Option::is_none")]
434    ended_at: Option<String>,
435    namespace: String,
436}
437
438impl SessionRecord {
439    fn from_session(session: &Session, ns: &NamespaceId) -> Self {
440        Self {
441            session_id: session.id.0.clone(),
442            harness: session.harness.as_wire_str().to_string(),
443            channel: session.channel.clone(),
444            started_at: session.started_at.clone(),
445            ended_at: session.ended_at.clone(),
446            namespace: ns.as_str().to_string(),
447        }
448    }
449
450    fn into_session(self) -> Session {
451        Session {
452            id: SessionId(self.session_id),
453            harness: Harness::from_wire_str(&self.harness),
454            channel: self.channel,
455            started_at: self.started_at,
456            ended_at: self.ended_at,
457        }
458    }
459}
460
461/// Stored row for a queued mining extraction (ADR M2). `ts` is an internal
462/// epoch-millis for correct ORDER BY.
463#[derive(Debug, Clone, Serialize, Deserialize)]
464struct QueueRecord {
465    queue_id: String,
466    /// Nested memory fields (flattened for storage).
467    memory_id: String,
468    content: String,
469    project: String,
470    topic: String,
471    #[serde(default)]
472    importance: f32,
473    session_id: Option<String>,
474    harness: String,
475    confidence: f32,
476    source_session_id: String,
477    ts: i64,
478    namespace: String,
479}
480
481impl QueueRecord {
482    fn from_extraction(
483        queue_id: &str,
484        memory: &Memory,
485        confidence: f32,
486        ns: &NamespaceId,
487        now_ms: i64,
488    ) -> Self {
489        Self {
490            queue_id: queue_id.to_string(),
491            memory_id: memory.id.0.clone(),
492            content: memory.content.clone(),
493            project: memory.project.clone(),
494            topic: memory.topic.clone(),
495            importance: memory.importance,
496            session_id: memory.session_id.clone(),
497            harness: memory.harness.as_wire_str().to_string(),
498            confidence,
499            source_session_id: memory.session_id.clone().unwrap_or_default(),
500            ts: now_ms,
501            namespace: ns.as_str().to_string(),
502        }
503    }
504
505    fn into_queued(self) -> QueuedExtraction {
506        let harness = Harness::from_wire_str(&self.harness);
507        let memory = Memory {
508            id: MemoryId(self.memory_id),
509            content: self.content,
510            project: self.project,
511            topic: self.topic,
512            source: MemorySource::Mined,
513            harness,
514            session_id: self.session_id,
515            origin: InstanceId::local(),
516            authority: AuthorityScope::local(),
517            importance: self.importance,
518            created_at: String::new(),
519        };
520        QueuedExtraction {
521            id: self.queue_id,
522            memory,
523            confidence: self.confidence,
524            source_session_id: self.source_session_id,
525            queued_at: self.ts.to_string(),
526        }
527    }
528
529    fn into_memory(self) -> Memory {
530        Memory {
531            id: MemoryId(self.memory_id),
532            content: self.content,
533            project: self.project,
534            topic: self.topic,
535            source: MemorySource::Mined,
536            harness: Harness::from_wire_str(&self.harness),
537            session_id: self.session_id,
538            origin: InstanceId::local(),
539            authority: AuthorityScope::local(),
540            importance: self.importance,
541            created_at: String::new(),
542        }
543    }
544}
545
546fn store_err(e: surrealdb::Error) -> IjimaError {
547    IjimaError::Store {
548        detail: e.to_string(),
549    }
550}
551
552/// Current epoch-millis (for internal ORDER BY fields).
553fn now_millis() -> i64 {
554    std::time::SystemTime::now()
555        .duration_since(std::time::UNIX_EPOCH)
556        .map(|d| d.as_millis() as i64)
557        .unwrap_or(0)
558}
559
560/// Lowercase hex of a byte slice (for content-hash).
561fn hex(bytes: &[u8]) -> String {
562    bytes.iter().map(|b| format!("{b:02x}")).collect()
563}
564
565fn embed_for(embedder: &dyn Embedder, text: &str) -> Result<Option<Vec<f32>>> {
566    Ok(Some(embedder.embed(text)?.0))
567}
568
569/// Surreal record key for a memory: `<namespace>:<memory-id>`. Record
570/// ids are global per table in SurrealDB, so the namespace must be part
571/// of the key — otherwise the same logical id (e.g. the same pi-mempalace
572/// row imported from two workstations, WS2) collides across namespaces.
573/// The logical id stays in `memory_id` for wire responses.
574fn memory_key(ns: &NamespaceId, id: &MemoryId) -> String {
575    format!("{}:{}", ns.as_str(), id.0)
576}
577
578#[async_trait]
579impl Store for SurrealStore {
580    async fn store_memory(&self, ns: &NamespaceId, memory: Memory) -> Result<MemoryId> {
581        // Content-hash dedup (Phase 2.2): reject exact duplicates within
582        // the namespace, returning the existing id so callers can recover.
583        if let Some(existing) = self.check_duplicate(ns, &memory.content).await? {
584            return Err(IjimaError::duplicate(format!(
585                "content already stored as {}",
586                existing.0
587            )));
588        }
589        let (embedding, embed_model) = match &self.embedder {
590            Some(e) => {
591                // Stamp the model id (D10 provenance) so a future model
592                // swap is detectable and a re-embed pass can be triggered.
593                (
594                    embed_for(e.as_ref(), &memory.content)?,
595                    Some(e.model_id().to_string()),
596                )
597            }
598            None => (None, None),
599        };
600        let record = MemoryRecord::from_memory(&memory, ns, embedding, embed_model);
601        let _: Option<surrealdb::types::SerdeWrapper<MemoryRecord>> = self
602            .db
603            .create((MEMORIES_TABLE, memory_key(ns, &memory.id)))
604            .content(surrealdb::types::SerdeWrapper(record))
605            .await
606            .map_err(store_err)?;
607        Ok(memory.id)
608    }
609
610    async fn check_duplicate(&self, ns: &NamespaceId, content: &str) -> Result<Option<MemoryId>> {
611        use sha2::{Digest, Sha256};
612        let hash = hex(&Sha256::digest(content.as_bytes()));
613        let mut result = self
614            .db
615            .query(format!(
616                "SELECT memory_id FROM {MEMORIES_TABLE}
617                 WHERE namespace = $ns AND content_hash = $hash LIMIT 1"
618            ))
619            .bind(("ns", ns.as_str().to_string()))
620            .bind(("hash", hash))
621            .await
622            .map_err(store_err)?;
623        #[derive(serde::Serialize, Deserialize)]
624        struct IdRow {
625            memory_id: String,
626        }
627        let rows = take_vec::<IdRow>(&mut result)?;
628        Ok(rows.into_iter().next().map(|r| MemoryId(r.memory_id)))
629    }
630
631    async fn recall_memory(&self, ns: &NamespaceId, id: &MemoryId) -> Result<Option<Memory>> {
632        let record: Option<surrealdb::types::SerdeWrapper<MemoryRecord>> = self
633            .db
634            .select((MEMORIES_TABLE, memory_key(ns, id)))
635            .await
636            .map_err(store_err)?;
637        Ok(record.map(|w| w.0.into_memory()))
638    }
639
640    async fn delete_memory(&self, ns: &NamespaceId, id: &MemoryId) -> Result<()> {
641        // The namespaced record key is unique, so delete is naturally
642        // isolated (no-op when the record is absent).
643        let _: Option<surrealdb::types::SerdeWrapper<MemoryRecord>> = self
644            .db
645            .delete((MEMORIES_TABLE, memory_key(ns, id)))
646            .await
647            .map_err(store_err)?;
648        Ok(())
649    }
650
651    async fn list_memories(&self, ns: &NamespaceId, limit: usize) -> Result<Vec<Memory>> {
652        let mut result = self
653            .db
654            .query(format!(
655                "SELECT memory_id, content, project, topic, source, harness, session_id, namespace,
656                        origin, authority, importance, created_at
657                 FROM {MEMORIES_TABLE}
658                 WHERE namespace = $ns
659                 ORDER BY importance DESC, created_at DESC
660                 LIMIT $lim"
661            ))
662            .bind(("ns", ns.as_str().to_string()))
663            .bind(("lim", limit as i64))
664            .await
665            .map_err(store_err)?;
666        let records = take_vec::<MemoryRecord>(&mut result)?;
667        Ok(records.into_iter().map(|r| r.into_memory()).collect())
668    }
669
670    async fn store_stats(&self) -> Result<StoreStats> {
671        // One query for all namespaces; aggregate in Rust (avoids
672        // SurrealDB aggregate-function ambiguity, fine at v0 scale).
673        #[derive(serde::Serialize, Deserialize)]
674        struct NsRow {
675            namespace: String,
676        }
677        let mut result = self
678            .db
679            .query(format!("SELECT namespace FROM {MEMORIES_TABLE}"))
680            .await
681            .map_err(store_err)?;
682        let rows = take_vec::<NsRow>(&mut result)?;
683        let mut counts: std::collections::BTreeMap<String, usize> =
684            std::collections::BTreeMap::new();
685        for r in &rows {
686            *counts.entry(r.namespace.clone()).or_insert(0) += 1;
687        }
688        let namespaces = counts
689            .into_iter()
690            .map(|(namespace, memories)| NamespaceCount {
691                namespace,
692                memories,
693            })
694            .collect();
695        Ok(StoreStats {
696            total_memories: rows.len(),
697            namespaces,
698        })
699    }
700
701    async fn list_rooms(
702        &self,
703        ns: &NamespaceId,
704        project: Option<&str>,
705        limit: usize,
706    ) -> Result<Vec<Room>> {
707        let counts = self.project_topic_counts(ns).await?;
708        let mut rooms: Vec<Room> = counts
709            .into_iter()
710            .filter(|((p, _), _)| project.is_none_or(|proj| p == proj))
711            .map(|((project, topic), count)| Room {
712                project,
713                topic,
714                count,
715            })
716            .collect();
717        // count desc, then topic for determinism.
718        rooms.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.topic.cmp(&b.topic)));
719        rooms.truncate(limit);
720        Ok(rooms)
721    }
722
723    async fn taxonomy(&self, ns: &NamespaceId) -> Result<Vec<ProjectTaxon>> {
724        let counts = self.project_topic_counts(ns).await?;
725        // group by project
726        let mut by_project: std::collections::BTreeMap<String, Vec<Room>> =
727            std::collections::BTreeMap::new();
728        for ((project, topic), count) in counts {
729            by_project.entry(project).or_default().push(Room {
730                project: String::new(),
731                topic,
732                count,
733            });
734        }
735        let mut taxons: Vec<ProjectTaxon> = by_project
736            .into_iter()
737            .map(|(project, mut rooms)| {
738                rooms.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.topic.cmp(&b.topic)));
739                for r in &mut rooms {
740                    r.project = project.clone();
741                }
742                let total = rooms.iter().map(|r| r.count).sum();
743                ProjectTaxon {
744                    project,
745                    rooms,
746                    total,
747                }
748            })
749            .collect();
750        // projects with the most memories first.
751        taxons.sort_by(|a, b| {
752            b.total
753                .cmp(&a.total)
754                .then_with(|| a.project.cmp(&b.project))
755        });
756        Ok(taxons)
757    }
758
759    async fn palace_graph(&self, ns: &NamespaceId) -> Result<PalaceGraph> {
760        let counts = self.project_topic_counts(ns).await?;
761        // topic -> set of (project, count)
762        let mut topics: std::collections::BTreeMap<String, Vec<(String, usize)>> =
763            std::collections::BTreeMap::new();
764        let mut projects: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
765        for ((project, topic), count) in counts {
766            projects.insert(project.clone());
767            topics.entry(topic).or_default().push((project, count));
768        }
769        let mut tunnels = Vec::new();
770        for (topic, mut entries) in topics {
771            if entries.len() < 2 {
772                continue; // a tunnel needs two distinct projects
773            }
774            entries.sort_by(|a, b| a.0.cmp(&b.0));
775            // every unordered pair of distinct projects on this topic
776            for i in 0..entries.len() {
777                for j in (i + 1)..entries.len() {
778                    tunnels.push(Tunnel {
779                        topic: topic.clone(),
780                        project_a: entries[i].0.clone(),
781                        project_b: entries[j].0.clone(),
782                        count_a: entries[i].1,
783                        count_b: entries[j].1,
784                    });
785                }
786            }
787        }
788        Ok(PalaceGraph {
789            projects: projects.into_iter().collect(),
790            tunnels,
791        })
792    }
793
794    async fn traverse_tunnel(
795        &self,
796        ns: &NamespaceId,
797        topic: &str,
798        project_a: &str,
799        project_b: &str,
800        limit: usize,
801    ) -> Result<TunnelTraversal> {
802        let memories_a = self
803            .project_topic_memories(ns, project_a, topic, limit)
804            .await?;
805        let memories_b = self
806            .project_topic_memories(ns, project_b, topic, limit)
807            .await?;
808        Ok(TunnelTraversal {
809            topic: topic.to_string(),
810            project_a: project_a.to_string(),
811            project_b: project_b.to_string(),
812            memories_a,
813            memories_b,
814        })
815    }
816
817    async fn search_memories(
818        &self,
819        ns: &NamespaceId,
820        embedding: &Embedding,
821        limit: usize,
822    ) -> Result<Vec<SearchHit>> {
823        if self.embedder.is_none() {
824            return Err(IjimaError::Store {
825                detail: "search requires the store to opened with an Embedder".into(),
826            });
827        }
828        // Cosine similarity ranking, brute-force over the namespace.
829        // Correct (no ANN approximation); HNSW is the planned
830        // optimization once the SDK emits typed `<N>f32` vectors.
831        let mut result = self
832            .db
833            .query(format!(
834                "SELECT memory_id, content, project, topic, source, harness, session_id, namespace,
835                        origin, authority, importance, created_at,
836                        vector::similarity::cosine(embedding, $query) AS score
837                 FROM {MEMORIES_TABLE}
838                 WHERE namespace = $ns AND embedding IS NOT NONE
839                 ORDER BY score DESC
840                 LIMIT $lim"
841            ))
842            .bind(("ns", ns.as_str().to_string()))
843            .bind(("query", embedding.0.clone()))
844            .bind(("lim", limit as i64))
845            .await
846            .map_err(store_err)?;
847        // Deserialize each row as a MemoryRecord (missing fields default)
848        // flattened alongside the computed `score`.
849        #[derive(serde::Serialize, Deserialize)]
850        struct ScoredRecord {
851            #[serde(flatten)]
852            rec: MemoryRecord,
853            score: f64,
854        }
855        let rows = take_vec::<ScoredRecord>(&mut result)?;
856        Ok(rows
857            .into_iter()
858            .map(|r| SearchHit {
859                memory: r.rec.into_memory(),
860                similarity: r.score as f32,
861            })
862            .collect())
863    }
864
865    async fn ingest_turn(&self, ns: &NamespaceId, turn: SessionTurn) -> Result<()> {
866        let record = SessionTurnRecord {
867            session_id: turn.session_id.0.clone(),
868            turn_index: turn.turn_index,
869            role: turn.role,
870            content: turn.content,
871            timestamp: turn.timestamp,
872            namespace: ns.as_str().to_string(),
873        };
874        let _: Option<surrealdb::types::SerdeWrapper<SessionTurnRecord>> = self
875            .db
876            .create(TURNS_TABLE)
877            .content(surrealdb::types::SerdeWrapper(record))
878            .await
879            .map_err(store_err)?;
880        Ok(())
881    }
882
883    async fn session_turns(
884        &self,
885        ns: &NamespaceId,
886        session: &SessionId,
887        limit: usize,
888    ) -> Result<Vec<SessionTurn>> {
889        let mut result = self
890            .db
891            .query(format!(
892                "SELECT session_id, turn_index, role, content, timestamp, namespace
893                 FROM {TURNS_TABLE}
894                 WHERE namespace = $ns AND session_id = $sid
895                 ORDER BY turn_index DESC LIMIT $lim",
896            ))
897            .bind(("ns", ns.as_str().to_string()))
898            .bind(("sid", session.0.clone()))
899            .bind(("lim", limit as i64))
900            .await
901            .map_err(store_err)?;
902        let mut records = take_vec::<SessionTurnRecord>(&mut result)?;
903        records.reverse(); // query was DESC for "last N"; return chronological
904        Ok(records
905            .into_iter()
906            .map(|r| SessionTurn {
907                session_id: SessionId(r.session_id),
908                turn_index: r.turn_index,
909                role: r.role,
910                content: r.content,
911                timestamp: r.timestamp,
912            })
913            .collect())
914    }
915
916    async fn create_session(&self, ns: &NamespaceId, session: Session) -> Result<SessionId> {
917        let id_str = session.id.0.clone();
918        let record = SessionRecord::from_session(&session, ns);
919        // Upsert: create if absent, replace metadata if present. Session
920        // ids are globally unique (same convention as memory ids).
921        let _: Option<surrealdb::types::SerdeWrapper<SessionRecord>> = self
922            .db
923            .upsert((SESSIONS_TABLE, id_str.clone()))
924            .content(surrealdb::types::SerdeWrapper(record))
925            .await
926            .map_err(store_err)?;
927        Ok(session.id)
928    }
929
930    async fn list_sessions(
931        &self,
932        ns: &NamespaceId,
933        harness: Option<&Harness>,
934        limit: usize,
935    ) -> Result<Vec<Session>> {
936        let mut query = format!(
937            "SELECT session_id, harness, channel, started_at, ended_at, namespace
938             FROM {SESSIONS_TABLE}
939             WHERE namespace = $ns"
940        );
941        if harness.is_some() {
942            query.push_str(" AND harness = $harness");
943        }
944        query.push_str(" ORDER BY started_at DESC LIMIT $lim");
945        let mut q = self
946            .db
947            .query(query)
948            .bind(("ns", ns.as_str().to_string()))
949            .bind(("lim", limit as i64));
950        if let Some(h) = harness {
951            q = q.bind(("harness", h.as_wire_str().to_string()));
952        }
953        let res = q.await.map_err(store_err)?;
954        let mut res = res;
955        let rows = take_vec::<SessionRecord>(&mut res)?;
956        Ok(rows.into_iter().map(SessionRecord::into_session).collect())
957    }
958
959    async fn end_session(
960        &self,
961        ns: &NamespaceId,
962        session: &SessionId,
963        ended_at: String,
964    ) -> Result<()> {
965        // Scoped by namespace + session id (safety: a principal can only
966        // end sessions in their own namespace).
967        let _ = self
968            .db
969            .query(format!(
970                "UPDATE {SESSIONS_TABLE}
971                 SET ended_at = $ended
972                 WHERE namespace = $ns AND session_id = $sid"
973            ))
974            .bind(("ns", ns.as_str().to_string()))
975            .bind(("sid", session.0.clone()))
976            .bind(("ended", ended_at))
977            .await
978            .map_err(store_err)?;
979        Ok(())
980    }
981
982    async fn enqueue_extraction(
983        &self,
984        ns: &NamespaceId,
985        memory: Memory,
986        confidence: f32,
987    ) -> Result<String> {
988        let now_ms = now_millis();
989        let queue_id = format!("q_{}_{now_ms}", memory.id.0);
990        let record = QueueRecord::from_extraction(&queue_id, &memory, confidence, ns, now_ms);
991        let _: Option<surrealdb::types::SerdeWrapper<QueueRecord>> = self
992            .db
993            .create((QUEUE_TABLE, queue_id.clone()))
994            .content(surrealdb::types::SerdeWrapper(record))
995            .await
996            .map_err(store_err)?;
997        Ok(queue_id)
998    }
999
1000    async fn list_pending(&self, ns: &NamespaceId, limit: usize) -> Result<Vec<QueuedExtraction>> {
1001        let mut result = self
1002            .db
1003            .query(format!(
1004                "SELECT queue_id, memory_id, content, project, topic, importance, session_id, harness, confidence, source_session_id, ts, namespace
1005                 FROM {QUEUE_TABLE}
1006                 WHERE namespace = $ns
1007                 ORDER BY ts DESC LIMIT $lim"
1008            ))
1009            .bind(("ns", ns.as_str().to_string()))
1010            .bind(("lim", limit as i64))
1011            .await
1012            .map_err(store_err)?;
1013        let records = take_vec::<QueueRecord>(&mut result)?;
1014        Ok(records.into_iter().map(QueueRecord::into_queued).collect())
1015    }
1016
1017    async fn accept_extraction(
1018        &self,
1019        ns: &NamespaceId,
1020        queue_id: &str,
1021    ) -> Result<AcceptedExtraction> {
1022        // Read the queued record (scoped by namespace), promote to the
1023        // palace, then delete from the queue.
1024        let mut result = self
1025            .db
1026            .query(format!(
1027                "SELECT queue_id, memory_id, content, project, topic, importance, session_id, harness, confidence, source_session_id, ts, namespace
1028                 FROM {QUEUE_TABLE}
1029                 WHERE namespace = $ns AND queue_id = $qid LIMIT 1"
1030            ))
1031            .bind(("ns", ns.as_str().to_string()))
1032            .bind(("qid", queue_id.to_string()))
1033            .await
1034            .map_err(store_err)?;
1035        let records = take_vec::<QueueRecord>(&mut result)?;
1036        let Some(record) = records.into_iter().next() else {
1037            return Err(IjimaError::not_found(format!("queue entry {queue_id}")));
1038        };
1039        let memory = record.into_memory();
1040        let memory_id = memory.id.clone();
1041        // store_memory does content-hash dedup; a duplicate promote is a no-op.
1042        self.store_memory(ns, memory).await?;
1043        let _: Option<surrealdb::types::SerdeWrapper<QueueRecord>> = self
1044            .db
1045            .delete((QUEUE_TABLE, queue_id.to_string()))
1046            .await
1047            .map_err(store_err)?;
1048        Ok(AcceptedExtraction { memory_id })
1049    }
1050
1051    async fn reject_extraction(&self, ns: &NamespaceId, queue_id: &str) -> Result<()> {
1052        // Scoped delete: only removes if the namespace matches.
1053        let _ = self
1054            .db
1055            .query(format!(
1056                "DELETE FROM {QUEUE_TABLE}
1057                 WHERE namespace = $ns AND queue_id = $qid"
1058            ))
1059            .bind(("ns", ns.as_str().to_string()))
1060            .bind(("qid", queue_id.to_string()))
1061            .await
1062            .map_err(store_err)?;
1063        Ok(())
1064    }
1065    async fn write_diary(&self, ns: &NamespaceId, entry: DiaryEntry) -> Result<()> {
1066        let now_ms = std::time::SystemTime::now()
1067            .duration_since(std::time::UNIX_EPOCH)
1068            .map(|d| d.as_millis() as i64)
1069            .unwrap_or(0);
1070        let record = DiaryRecord::from_entry(&entry, ns, now_ms);
1071        let _: Option<surrealdb::types::SerdeWrapper<DiaryRecord>> = self
1072            .db
1073            .create(DIARY_TABLE)
1074            .content(surrealdb::types::SerdeWrapper(record))
1075            .await
1076            .map_err(store_err)?;
1077        Ok(())
1078    }
1079
1080    async fn read_diary(
1081        &self,
1082        ns: &NamespaceId,
1083        agent: &str,
1084        limit: usize,
1085    ) -> Result<Vec<DiaryEntry>> {
1086        let mut result = self
1087            .db
1088            .query(format!(
1089                "SELECT agent, content, topic, timestamp, ts, namespace
1090                 FROM {DIARY_TABLE}
1091                 WHERE namespace = $ns AND agent = $agent
1092                 ORDER BY ts DESC LIMIT $lim"
1093            ))
1094            .bind(("ns", ns.as_str().to_string()))
1095            .bind(("agent", agent.to_string()))
1096            .bind(("lim", limit as i64))
1097            .await
1098            .map_err(store_err)?;
1099        let mut records = take_vec::<DiaryRecord>(&mut result)?;
1100        records.reverse(); // chronological (DESC → reverse)
1101        Ok(records.into_iter().map(DiaryRecord::into_entry).collect())
1102    }
1103
1104    // ===== Repo directory (global registry — Context Mapper) =====
1105
1106    async fn register_repo(&self, repo: RepoDirectory) -> Result<()> {
1107        let name = repo.name.clone();
1108        let _: Option<surrealdb::types::SerdeWrapper<RepoDirectory>> = self
1109            .db
1110            .upsert((REPO_TABLE, name))
1111            .content(surrealdb::types::SerdeWrapper(repo))
1112            .await
1113            .map_err(store_err)?;
1114        Ok(())
1115    }
1116
1117    async fn list_repos(&self) -> Result<Vec<RepoDirectory>> {
1118        let mut result = self
1119            .db
1120            .query(format!("SELECT * FROM {REPO_TABLE} ORDER BY name"))
1121            .await
1122            .map_err(store_err)?;
1123        let repos = take_vec::<RepoDirectory>(&mut result)?;
1124        Ok(repos)
1125    }
1126
1127    // ===== Token revocation (WS1b) =====
1128
1129    async fn revoke_token(&self, revocation: TokenRevocation) -> Result<()> {
1130        let hash = revocation.token_hash.clone();
1131        let _: Option<surrealdb::types::SerdeWrapper<TokenRevocation>> = self
1132            .db
1133            .upsert((REVOCATIONS_TABLE, hash))
1134            .content(surrealdb::types::SerdeWrapper(revocation))
1135            .await
1136            .map_err(store_err)?;
1137        Ok(())
1138    }
1139
1140    async fn list_revocations(&self) -> Result<Vec<TokenRevocation>> {
1141        let mut result = self
1142            .db
1143            .query(format!(
1144                "SELECT * FROM {REVOCATIONS_TABLE} ORDER BY revoked_at_unix"
1145            ))
1146            .await
1147            .map_err(store_err)?;
1148        let revs = take_vec::<TokenRevocation>(&mut result)?;
1149        Ok(revs)
1150    }
1151
1152    async fn grant_namespace_membership(&self, membership: NamespaceMembership) -> Result<()> {
1153        // Upsert by the natural key (ns, principal): re-granting refreshes
1154        // granted_at/granted_by without duplicating.
1155        let key = format!("{}:{}", membership.namespace, membership.principal);
1156        let _: Option<surrealdb::types::SerdeWrapper<NamespaceMembership>> = self
1157            .db
1158            .upsert((NAMESPACE_MEMBERS_TABLE, key))
1159            .content(surrealdb::types::SerdeWrapper(membership))
1160            .await
1161            .map_err(store_err)?;
1162        Ok(())
1163    }
1164
1165    async fn revoke_namespace_membership(&self, ns: &NamespaceId, principal: &str) -> Result<()> {
1166        let key = format!("{}:{}", ns.as_str(), principal);
1167        let _: Option<surrealdb::types::SerdeWrapper<NamespaceMembership>> = self
1168            .db
1169            .delete((NAMESPACE_MEMBERS_TABLE, key))
1170            .await
1171            .map_err(store_err)?;
1172        Ok(())
1173    }
1174
1175    async fn list_namespace_members(&self, ns: &NamespaceId) -> Result<Vec<NamespaceMembership>> {
1176        let mut result = self
1177            .db
1178            .query(format!(
1179                "SELECT * FROM {NAMESPACE_MEMBERS_TABLE} WHERE namespace = $ns \
1180                 ORDER BY granted_at_unix"
1181            ))
1182            .bind(("ns", ns.as_str().to_string()))
1183            .await
1184            .map_err(store_err)?;
1185        take_vec::<NamespaceMembership>(&mut result)
1186    }
1187
1188    async fn is_namespace_member(&self, ns: &NamespaceId, principal: &str) -> Result<bool> {
1189        // Direct natural-key lookup (`<ns>:<principal>`) — no query,
1190        // no aggregation semantics, the same key grant/revoke address.
1191        let key = format!("{}:{}", ns.as_str(), principal);
1192        let record: Option<surrealdb::types::SerdeWrapper<NamespaceMembership>> = self
1193            .db
1194            .select((NAMESPACE_MEMBERS_TABLE, key))
1195            .await
1196            .map_err(store_err)?;
1197        Ok(record.is_some())
1198    }
1199}
1200
1201// ---------- knowledge graph ----------
1202
1203#[derive(Debug, Clone, Serialize, Deserialize)]
1204struct EntityRecord_ {
1205    name: String,
1206    entity_type: String,
1207    namespace: String,
1208}
1209
1210#[derive(Debug, Clone, Serialize, Deserialize)]
1211struct TripleRecord {
1212    /// Deterministic id: `<subject>:<predicate>:<object>`.
1213    triple_id: String,
1214    subject: String,
1215    predicate: String,
1216    object: String,
1217    #[serde(default)]
1218    valid_from: Option<String>,
1219    #[serde(default)]
1220    valid_to: Option<String>,
1221    #[serde(default = "default_record_importance")]
1222    confidence: f32,
1223    namespace: String,
1224    #[serde(default)]
1225    source_memory_id: Option<String>,
1226}
1227
1228impl TripleRecord {
1229    fn into_triple(self) -> Triple {
1230        Triple {
1231            id: self.triple_id,
1232            subject: EntityId(self.subject),
1233            predicate: self.predicate,
1234            object: EntityId(self.object),
1235            valid_from: self.valid_from,
1236            valid_to: self.valid_to,
1237            confidence: self.confidence,
1238            namespace: self.namespace,
1239            source_memory_id: self.source_memory_id,
1240        }
1241    }
1242}
1243
1244#[async_trait::async_trait]
1245impl KnowledgeGraph for SurrealStore {
1246    async fn add_triple(
1247        &self,
1248        ns: &NamespaceId,
1249        subject: EntityId,
1250        predicate: &str,
1251        object: EntityId,
1252        valid_from: Option<&str>,
1253        confidence: f32,
1254        source_memory_id: Option<&str>,
1255    ) -> Result<Triple> {
1256        let ns_str = ns.as_str().to_string();
1257        let subj = subject.0.clone();
1258        let obj = object.0.clone();
1259        // Create both entity nodes (idempotent — ignore already-exists).
1260        for eid in [&subj, &obj] {
1261            let record = EntityRecord_ {
1262                name: eid.clone(),
1263                entity_type: "unknown".into(),
1264                namespace: ns_str.clone(),
1265            };
1266            let _: std::result::Result<Option<surrealdb::types::SerdeWrapper<EntityRecord_>>, _> =
1267                self.db
1268                    .create((ENTITIES_TABLE, eid.clone()))
1269                    .content(surrealdb::types::SerdeWrapper(record))
1270                    .await;
1271            // Ignore errors (record already exists from a prior triple).
1272        }
1273        // Store the triple as a record (keyed by deterministic triple_id).
1274        // Native graph edges (RELATE) are a future optimization for
1275        // multi-hop traversal; v0 queries via subject/object fields.
1276        let triple_id = format!("{subj}:{predicate}:{obj}");
1277        let record = TripleRecord {
1278            triple_id: triple_id.clone(),
1279            subject: subj.clone(),
1280            predicate: predicate.to_string(),
1281            object: obj.clone(),
1282            valid_from: valid_from.map(str::to_string),
1283            valid_to: None,
1284            confidence,
1285            namespace: ns_str.clone(),
1286            source_memory_id: source_memory_id.map(str::to_string),
1287        };
1288        let _: Option<surrealdb::types::SerdeWrapper<TripleRecord>> = self
1289            .db
1290            .create((TRIPLES_TABLE, triple_id.clone()))
1291            .content(surrealdb::types::SerdeWrapper(record))
1292            .await
1293            .map_err(store_err)?;
1294        Ok(Triple {
1295            id: triple_id,
1296            subject: EntityId(subj),
1297            predicate: predicate.to_string(),
1298            object: EntityId(obj),
1299            valid_from: valid_from.map(str::to_string),
1300            valid_to: None,
1301            confidence,
1302            namespace: ns_str,
1303            source_memory_id: source_memory_id.map(str::to_string),
1304        })
1305    }
1306
1307    async fn query_entity(&self, ns: &NamespaceId, entity: &EntityId) -> Result<EntityRecord> {
1308        let eid = entity.0.clone();
1309        let ent: Option<surrealdb::types::SerdeWrapper<EntityRecord_>> = self
1310            .db
1311            .select((ENTITIES_TABLE, eid.clone()))
1312            .await
1313            .map_err(store_err)?;
1314        let entity_node = ent.and_then(|w| {
1315            let r = w.0;
1316            if r.namespace == ns.as_str() {
1317                Some(Entity {
1318                    id: entity.clone(),
1319                    name: r.name,
1320                    entity_type: r.entity_type,
1321                    namespace: r.namespace,
1322                })
1323            } else {
1324                None
1325            }
1326        });
1327        let mut res = self
1328            .db
1329            .query(format!(
1330                "SELECT triple_id, subject, predicate, object, valid_from, valid_to, confidence, namespace, source_memory_id
1331                 FROM {TRIPLES_TABLE}
1332                 WHERE namespace = $ns AND (subject = $eid OR object = $eid)"
1333            ))
1334            .bind(("ns", ns.as_str().to_string()))
1335            .bind(("eid", eid))
1336            .await
1337            .map_err(store_err)?;
1338        let triples = take_vec::<TripleRecord>(&mut res)?;
1339        let requested = entity.0.as_str();
1340        let mut outgoing = Vec::new();
1341        let mut incoming = Vec::new();
1342        for t in triples {
1343            if t.subject == requested {
1344                outgoing.push(t.into_triple());
1345            } else {
1346                incoming.push(t.into_triple());
1347            }
1348        }
1349        Ok(EntityRecord {
1350            entity: entity_node,
1351            outgoing,
1352            incoming,
1353        })
1354    }
1355
1356    async fn invalidate_triple(&self, ns: &NamespaceId, triple_id: &str) -> Result<()> {
1357        let now = std::time::SystemTime::now()
1358            .duration_since(std::time::UNIX_EPOCH)
1359            .map(|d| d.as_secs().to_string())
1360            .unwrap_or_default();
1361        let _ = self
1362            .db
1363            .query(format!(
1364                "UPDATE {TRIPLES_TABLE} SET valid_to = $now
1365                 WHERE triple_id = $tid AND namespace = $ns"
1366            ))
1367            .bind(("now", now))
1368            .bind(("tid", triple_id.to_string()))
1369            .bind(("ns", ns.as_str().to_string()))
1370            .await
1371            .map_err(store_err)?;
1372        Ok(())
1373    }
1374
1375    async fn find_triples(
1376        &self,
1377        ns: &NamespaceId,
1378        subject: Option<&EntityId>,
1379        predicate: Option<&str>,
1380        object: Option<&EntityId>,
1381    ) -> Result<Vec<Triple>> {
1382        let mut conditions = vec!["namespace = $ns".to_string()];
1383        if subject.is_some() {
1384            conditions.push("subject = $subj".into());
1385        }
1386        if predicate.is_some() {
1387            conditions.push("predicate = $pred".into());
1388        }
1389        if object.is_some() {
1390            conditions.push("object = $obj".into());
1391        }
1392        let where_ = conditions.join(" AND ");
1393        let mut q = self
1394            .db
1395            .query(format!(
1396                "SELECT triple_id, subject, predicate, object, valid_from, valid_to, confidence, namespace, source_memory_id
1397                 FROM {TRIPLES_TABLE} WHERE {where_} LIMIT 100"
1398            ))
1399            .bind(("ns", ns.as_str().to_string()));
1400        if let Some(s) = subject {
1401            q = q.bind(("subj", s.0.clone()));
1402        }
1403        if let Some(p) = predicate {
1404            q = q.bind(("pred", p.to_string()));
1405        }
1406        if let Some(o) = object {
1407            q = q.bind(("obj", o.0.clone()));
1408        }
1409        let mut res = q.await.map_err(store_err)?;
1410        let triples = take_vec::<TripleRecord>(&mut res)?;
1411        Ok(triples.into_iter().map(TripleRecord::into_triple).collect())
1412    }
1413
1414    async fn kg_timeline(&self, ns: &NamespaceId, limit: usize) -> Result<Vec<Triple>> {
1415        let mut res = self
1416            .db
1417            .query(format!(
1418                "SELECT triple_id, subject, predicate, object, valid_from, valid_to, confidence, namespace, source_memory_id
1419                 FROM {TRIPLES_TABLE}
1420                 WHERE namespace = $ns
1421                 ORDER BY valid_from DESC LIMIT $lim"
1422            ))
1423            .bind(("ns", ns.as_str().to_string()))
1424            .bind(("lim", limit as i64))
1425            .await
1426            .map_err(store_err)?;
1427        let triples = take_vec::<TripleRecord>(&mut res)?;
1428        Ok(triples.into_iter().map(TripleRecord::into_triple).collect())
1429    }
1430
1431    async fn knowledge_stats(&self, ns: &NamespaceId) -> Result<KgStats> {
1432        let mut res = self
1433            .db
1434            .query(format!(
1435                "SELECT namespace FROM {ENTITIES_TABLE} WHERE namespace = $ns;
1436                 SELECT namespace FROM {TRIPLES_TABLE} WHERE namespace = $ns;"
1437            ))
1438            .bind(("ns", ns.as_str().to_string()))
1439            .await
1440            .map_err(store_err)?;
1441        #[derive(serde::Serialize, Deserialize)]
1442        struct Row {
1443            #[allow(dead_code)]
1444            namespace: String,
1445        }
1446        let entities = take_vec::<Row>(&mut res)?;
1447        let wrapped: Vec<surrealdb::types::SerdeWrapper<Row>> = res.take(1).map_err(store_err)?;
1448        let triples: Vec<Row> = wrapped.into_iter().map(|w| w.0).collect();
1449        Ok(KgStats {
1450            entities: entities.len(),
1451            triples: triples.len(),
1452        })
1453    }
1454
1455    async fn kg_global_stats(&self) -> Result<KgStats> {
1456        let mut res = self
1457            .db
1458            .query(format!(
1459                "SELECT namespace FROM {ENTITIES_TABLE};
1460                 SELECT namespace FROM {TRIPLES_TABLE};"
1461            ))
1462            .await
1463            .map_err(store_err)?;
1464        #[derive(serde::Serialize, Deserialize)]
1465        struct Row {
1466            #[allow(dead_code)]
1467            namespace: String,
1468        }
1469        let entities = take_vec::<Row>(&mut res)?;
1470        let wrapped: Vec<surrealdb::types::SerdeWrapper<Row>> = res.take(1).map_err(store_err)?;
1471        let triples: Vec<Row> = wrapped.into_iter().map(|w| w.0).collect();
1472        Ok(KgStats {
1473            entities: entities.len(),
1474            triples: triples.len(),
1475        })
1476    }
1477}
1478
1479#[cfg(test)]
1480mod tests {
1481    use super::*;
1482    use ijima_core::{NamespaceId, TurnRole};
1483
1484    async fn fresh() -> SurrealStore {
1485        SurrealStore::open_embedded()
1486            .await
1487            .expect("embedded store must open")
1488    }
1489
1490    fn sample_memory(id: &str, content: &str) -> Memory {
1491        Memory {
1492            id: MemoryId(id.into()),
1493            content: content.into(),
1494            project: "ijima".into(),
1495            topic: "test".into(),
1496            source: MemorySource::Explicit,
1497            harness: Harness::Pi,
1498            session_id: Some("sess_1".into()),
1499            origin: InstanceId::local(),
1500            authority: AuthorityScope::local(),
1501            importance: 0.5,
1502            created_at: "0".into(),
1503        }
1504    }
1505
1506    #[tokio::test]
1507    async fn store_then_recall_round_trips() {
1508        let store = fresh().await;
1509        let ns = NamespaceId::new("ns_elliott_private");
1510        let id = store
1511            .store_memory(&ns, sample_memory("mem_1", "decided to use surrealdb"))
1512            .await
1513            .expect("store");
1514        assert_eq!(id.0.as_str(), "mem_1");
1515
1516        let got = store
1517            .recall_memory(&ns, &MemoryId("mem_1".into()))
1518            .await
1519            .expect("recall");
1520        let got = got.expect("must be present");
1521        assert_eq!(got.content, "decided to use surrealdb");
1522        assert_eq!(got.harness, Harness::Pi);
1523        assert_eq!(got.source, MemorySource::Explicit);
1524        // Provenance fields persist through a store→recall round-trip
1525        // (ADR provenance-tier): origin/authority survive SurrealDB.
1526        assert_eq!(got.origin, InstanceId::local());
1527        assert_eq!(got.authority, AuthorityScope::local());
1528    }
1529
1530    #[tokio::test]
1531    async fn same_memory_id_in_two_namespaces_does_not_collide() {
1532        // WS2: the same source corpus imported under two sources carries
1533        // the same logical ids into different namespaces. Record keys must
1534        // be namespaced or the second store 500s on `already exists`.
1535        let store = fresh().await;
1536        let ns_a = NamespaceId::new("ns_import_laptop");
1537        let ns_b = NamespaceId::new("ns_import_desktop");
1538        store
1539            .store_memory(&ns_a, sample_memory("mem_fx1", "laptop copy"))
1540            .await
1541            .expect("store into ns_a");
1542        store
1543            .store_memory(&ns_b, sample_memory("mem_fx1", "desktop copy"))
1544            .await
1545            .expect("same id into ns_b must not collide");
1546        // Each namespace recalls its own record under the same logical id.
1547        let a = store
1548            .recall_memory(&ns_a, &MemoryId("mem_fx1".into()))
1549            .await
1550            .expect("recall a")
1551            .expect("present a");
1552        let b = store
1553            .recall_memory(&ns_b, &MemoryId("mem_fx1".into()))
1554            .await
1555            .expect("recall b")
1556            .expect("present b");
1557        assert_eq!(a.content, "laptop copy");
1558        assert_eq!(b.content, "desktop copy");
1559        // Deleting in one namespace leaves the other intact.
1560        store
1561            .delete_memory(&ns_a, &MemoryId("mem_fx1".into()))
1562            .await
1563            .expect("delete a");
1564        assert!(
1565            store
1566                .recall_memory(&ns_b, &MemoryId("mem_fx1".into()))
1567                .await
1568                .expect("recall b after a delete")
1569                .is_some(),
1570            "ns_b record must survive ns_a deletion"
1571        );
1572    }
1573
1574    #[tokio::test]
1575    async fn list_memories_projects_origin_and_authority() {
1576        // The browse SELECT must project origin/authority or the
1577        // deserializer defaults them to `local`, silently erasing import
1578        // provenance (WS2).
1579        let store = fresh().await;
1580        let ns = NamespaceId::new("ns_originproj");
1581        let mut m = sample_memory("m_op", "origin projection probe");
1582        m.origin = InstanceId("elliotthall-laptop".into());
1583        store.store_memory(&ns, m).await.expect("store");
1584        let listed = store.list_memories(&ns, 10).await.expect("list");
1585        let got = listed.first().expect("one row");
1586        assert_eq!(got.origin.0, "elliotthall-laptop");
1587    }
1588
1589    #[tokio::test]
1590    async fn store_rejects_exact_duplicate_content() {
1591        let store = fresh().await;
1592        let ns = NamespaceId::new("ns_dedup");
1593        store
1594            .store_memory(&ns, sample_memory("m1", "identical content"))
1595            .await
1596            .expect("first store");
1597        // Same content, different id — must be rejected.
1598        let result = store
1599            .store_memory(&ns, sample_memory("m2", "identical content"))
1600            .await;
1601        assert!(matches!(result, Err(IjimaError::Duplicate { .. })));
1602
1603        // check_duplicate finds the existing id.
1604        let dup = store
1605            .check_duplicate(&ns, "identical content")
1606            .await
1607            .expect("check");
1608        assert_eq!(dup.as_ref().unwrap().0, "m1");
1609
1610        // Different content stores fine.
1611        store
1612            .store_memory(&ns, sample_memory("m3", "different content"))
1613            .await
1614            .expect("distinct content stores");
1615
1616        // Same content in a DIFFERENT namespace is not a duplicate.
1617        let other = NamespaceId::new("ns_other");
1618        store
1619            .store_memory(&other, sample_memory("mx", "identical content"))
1620            .await
1621            .expect("same content, different namespace");
1622    }
1623
1624    #[tokio::test]
1625    async fn index_definitions_are_idempotent_and_dedup_still_works() {
1626        // SurrealDB is schemaless with no auto-indexes; open_with_db defines
1627        // them. Re-running the DDL (as on every store open against an
1628        // existing database) must not error — IF NOT EXISTS guards it. This
1629        // is the fix for the O(N²) full-scan migration hang.
1630        let store = fresh().await;
1631        store
1632            .db
1633            .query(SurrealStore::INDEX_DDL)
1634            .await
1635            .expect("re-defining indexes on an already-indexed store");
1636        // Dedup is still correct after indexing (namespace + content_hash).
1637        let ns = NamespaceId::new("ns_idx");
1638        store
1639            .store_memory(&ns, sample_memory("m1", "indexed-content"))
1640            .await
1641            .expect("store");
1642        let dup = store
1643            .check_duplicate(&ns, "indexed-content")
1644            .await
1645            .expect("check");
1646        assert_eq!(dup.as_ref().unwrap().0, "m1");
1647        // Different namespace, same content: not a duplicate (composite index).
1648        let other = NamespaceId::new("ns_idx_other");
1649        store
1650            .store_memory(&other, sample_memory("m2", "indexed-content"))
1651            .await
1652            .expect("same content, different ns");
1653    }
1654
1655    #[tokio::test]
1656    async fn namespace_isolation_hides_other_namespace_memories() {
1657        let store = fresh().await;
1658        let alice = NamespaceId::new("ns_alice");
1659        let bob = NamespaceId::new("ns_bob");
1660
1661        store
1662            .store_memory(&alice, sample_memory("mem_a", "alice's secret"))
1663            .await
1664            .unwrap();
1665
1666        let got = store
1667            .recall_memory(&bob, &MemoryId("mem_a".into()))
1668            .await
1669            .expect("recall");
1670        assert!(got.is_none(), "namespace isolation must hide the memory");
1671    }
1672
1673    #[tokio::test]
1674    async fn delete_removes_owned_memory_only() {
1675        let store = fresh().await;
1676        let ns = NamespaceId::new("ns_elliott_private");
1677        store
1678            .store_memory(&ns, sample_memory("mem_1", "bye"))
1679            .await
1680            .unwrap();
1681        store
1682            .delete_memory(&ns, &MemoryId("mem_1".into()))
1683            .await
1684            .expect("delete");
1685        let got = store
1686            .recall_memory(&ns, &MemoryId("mem_1".into()))
1687            .await
1688            .unwrap();
1689        assert!(got.is_none());
1690    }
1691
1692    #[tokio::test]
1693    async fn ingest_then_read_session_turns_chronologically() {
1694        let store = fresh().await;
1695        let ns = NamespaceId::new("ns_elliott_private");
1696        let session = SessionId::new("sess_1");
1697        for (i, content) in ["first", "second", "third"].iter().enumerate() {
1698            store
1699                .ingest_turn(
1700                    &ns,
1701                    SessionTurn {
1702                        session_id: session.clone(),
1703                        turn_index: i as u64,
1704                        role: if i % 2 == 0 {
1705                            TurnRole::User
1706                        } else {
1707                            TurnRole::Assistant
1708                        },
1709                        content: (*content).into(),
1710                        timestamp: format!("2026-07-05T12:00:0{i}Z"),
1711                    },
1712                )
1713                .await
1714                .expect("ingest");
1715        }
1716
1717        let turns = store
1718            .session_turns(&ns, &session, 10)
1719            .await
1720            .expect("read turns");
1721        assert_eq!(turns.len(), 3);
1722        assert_eq!(turns[0].content, "first");
1723        assert_eq!(turns[2].content, "third");
1724    }
1725
1726    #[tokio::test]
1727    async fn session_turns_respect_limit_and_namespace() {
1728        let store = fresh().await;
1729        let ns = NamespaceId::new("ns_elliott_private");
1730        let other = NamespaceId::new("ns_other");
1731        let session = SessionId::new("sess_1");
1732
1733        for i in 0..5 {
1734            store
1735                .ingest_turn(
1736                    &ns,
1737                    SessionTurn {
1738                        session_id: session.clone(),
1739                        turn_index: i,
1740                        role: TurnRole::User,
1741                        content: format!("turn {i}"),
1742                        timestamp: format!("2026-07-05T12:00:0{i}Z"),
1743                    },
1744                )
1745                .await
1746                .unwrap();
1747        }
1748        store
1749            .ingest_turn(
1750                &other,
1751                SessionTurn {
1752                    session_id: session.clone(),
1753                    turn_index: 0,
1754                    role: TurnRole::User,
1755                    content: "intruder".into(),
1756                    timestamp: "2026-07-05T12:00:00Z".into(),
1757                },
1758            )
1759            .await
1760            .unwrap();
1761
1762        let last_two = store.session_turns(&ns, &session, 2).await.expect("read");
1763        assert_eq!(last_two.len(), 2);
1764        assert_eq!(last_two[1].content, "turn 4");
1765        assert!(last_two.iter().all(|t| t.content != "intruder"));
1766    }
1767
1768    #[tokio::test]
1769    async fn session_metadata_create_list_end_round_trip() {
1770        let store = fresh().await;
1771        let ns = NamespaceId::new("ns_sessions");
1772
1773        // Create two sessions (different harnesses).
1774        let s1 = Session {
1775            id: SessionId::new("sess_a"),
1776            harness: Harness::Pi,
1777            channel: Some("thread-1".into()),
1778            started_at: "2026-07-05T10:00:00Z".into(),
1779            ended_at: None,
1780        };
1781        let s2 = Session {
1782            id: SessionId::new("sess_b"),
1783            harness: Harness::Sakamoto,
1784            channel: None,
1785            started_at: "2026-07-05T12:00:00Z".into(),
1786            ended_at: None,
1787        };
1788        store
1789            .create_session(&ns, s1.clone())
1790            .await
1791            .expect("create s1");
1792        store
1793            .create_session(&ns, s2.clone())
1794            .await
1795            .expect("create s2");
1796
1797        // List all — newest first (s2 has later started_at).
1798        let all = store.list_sessions(&ns, None, 10).await.expect("list");
1799        assert_eq!(all.len(), 2);
1800        assert_eq!(all[0].id.as_str(), "sess_b");
1801        assert_eq!(all[1].id.as_str(), "sess_a");
1802        assert!(all[0].ended_at.is_none());
1803
1804        // Filter by harness.
1805        let pi_only = store
1806            .list_sessions(&ns, Some(&Harness::Pi), 10)
1807            .await
1808            .expect("list pi");
1809        assert_eq!(pi_only.len(), 1);
1810        assert_eq!(pi_only[0].id.as_str(), "sess_a");
1811
1812        // End s1.
1813        store
1814            .end_session(
1815                &ns,
1816                &SessionId::new("sess_a"),
1817                "2026-07-05T11:30:00Z".into(),
1818            )
1819            .await
1820            .expect("end");
1821        let after_end = store.list_sessions(&ns, None, 10).await.expect("list");
1822        let s1_row = after_end
1823            .iter()
1824            .find(|s| s.id.as_str() == "sess_a")
1825            .unwrap();
1826        assert_eq!(s1_row.ended_at.as_deref(), Some("2026-07-05T11:30:00Z"));
1827
1828        // Upsert (re-create s1) preserves metadata shape.
1829        store.create_session(&ns, s1.clone()).await.expect("upsert");
1830        let upserted = store
1831            .list_sessions(&ns, Some(&Harness::Pi), 10)
1832            .await
1833            .expect("list");
1834        assert_eq!(upserted.len(), 1);
1835
1836        // Namespace isolation: sessions in another namespace are invisible.
1837        let other = NamespaceId::new("ns_other");
1838        let cross = store
1839            .list_sessions(&other, None, 10)
1840            .await
1841            .expect("list other");
1842        assert!(cross.is_empty());
1843    }
1844
1845    // ===== Diaries =====
1846
1847    #[tokio::test]
1848    async fn diaries_write_appends_and_reads_chronologically() {
1849        let store = fresh().await;
1850        let ns = NamespaceId::new("ns_diaries");
1851
1852        let e1 = DiaryEntry {
1853            agent: "claude".into(),
1854            content: "first entry".into(),
1855            topic: None,
1856            timestamp: String::new(),
1857        };
1858        let e2 = DiaryEntry {
1859            agent: "claude".into(),
1860            content: "second entry".into(),
1861            topic: Some("reflection".into()),
1862            timestamp: "2026-07-05T12:00:00Z".into(),
1863        };
1864        store.write_diary(&ns, e1).await.expect("write1");
1865        store.write_diary(&ns, e2).await.expect("write2");
1866
1867        let entries = store.read_diary(&ns, "claude", 10).await.expect("read");
1868        assert_eq!(entries.len(), 2);
1869        // chronological: e1 before e2
1870        assert_eq!(entries[0].content, "first entry");
1871        assert_eq!(entries[1].content, "second entry");
1872        assert_eq!(entries[1].topic.as_deref(), Some("reflection"));
1873
1874        // Agent isolation: different agent sees nothing.
1875        let other = store.read_diary(&ns, "pi", 10).await.expect("read");
1876        assert!(other.is_empty());
1877
1878        // Namespace isolation: different ns sees nothing.
1879        let other_ns = NamespaceId::new("ns_other");
1880        let cross = store
1881            .read_diary(&other_ns, "claude", 10)
1882            .await
1883            .expect("read");
1884        assert!(cross.is_empty());
1885    }
1886
1887    // ===== Vector search =====
1888
1889    /// A deterministic test embedder: maps text to a small fixed-dim
1890    /// vector so nearest-neighbour behaviour is predictable. Production
1891    /// uses candle + all-MiniLM-L6-v2 (384-dim).
1892    struct TestEmbedder;
1893    impl Embedder for TestEmbedder {
1894        fn dim(&self) -> usize {
1895            4
1896        }
1897        fn embed(&self, text: &str) -> Result<Embedding> {
1898            // Cosine-relevant signal: word-count + keyword flags.
1899            let words = text.split_whitespace().count() as f32;
1900            let has_surreal = text.contains("surreal") as u32 as f32;
1901            let has_rust = text.contains("rust") as u32 as f32;
1902            let has_memory = text.contains("memory") as u32 as f32;
1903            Ok(Embedding(vec![words, has_surreal, has_rust, has_memory]))
1904        }
1905    }
1906
1907    #[tokio::test]
1908    async fn search_without_embedder_errors() {
1909        let store = fresh().await; // no embedder
1910        let ns = NamespaceId::new("ns_x");
1911        let result = store
1912            .search_memories(&ns, &Embedding(vec![0.0; 4]), 5)
1913            .await;
1914        assert!(matches!(result, Err(IjimaError::Store { .. })));
1915    }
1916
1917    #[tokio::test]
1918    async fn search_finds_nearest_with_embedder() {
1919        let store = SurrealStore::open_embedded_with(Arc::new(TestEmbedder))
1920            .await
1921            .expect("open with embedder");
1922        let ns = NamespaceId::new("ns_elliott_private");
1923
1924        // Three memories with distinct keyword signals.
1925        store
1926            .store_memory(&ns, sample_memory("m1", "rust memory store"))
1927            .await
1928            .unwrap();
1929        store
1930            .store_memory(&ns, sample_memory("m2", "surreal db graph"))
1931            .await
1932            .unwrap();
1933        store
1934            .store_memory(&ns, sample_memory("m3", "completely unrelated text"))
1935            .await
1936            .unwrap();
1937
1938        // Query close to the "rust memory" vector.
1939        let query = Embedding(vec![3.0, 0.0, 1.0, 1.0]);
1940        let hits = store.search_memories(&ns, &query, 2).await.expect("search");
1941        assert!(!hits.is_empty(), "must find at least one memory");
1942        // The nearest hit must be the rust/memory one (m1), not the
1943        // unrelated text.
1944        assert_eq!(hits[0].memory.content, "rust memory store");
1945        // Scored search: each hit carries its cosine similarity.
1946        assert!(hits[0].similarity > 0.0);
1947    }
1948
1949    #[tokio::test]
1950    async fn search_respects_namespace_isolation() {
1951        let store = SurrealStore::open_embedded_with(Arc::new(TestEmbedder))
1952            .await
1953            .expect("open with embedder");
1954        let alice = NamespaceId::new("ns_alice");
1955        let bob = NamespaceId::new("ns_bob");
1956
1957        store
1958            .store_memory(&alice, sample_memory("a1", "rust memory store"))
1959            .await
1960            .unwrap();
1961
1962        // Bob searches with an identical query embedding — must see
1963        // nothing from alice's namespace.
1964        let query = Embedding(vec![3.0, 0.0, 1.0, 1.0]);
1965        let hits = store
1966            .search_memories(&bob, &query, 5)
1967            .await
1968            .expect("search");
1969        assert!(
1970            hits.is_empty(),
1971            "namespace isolation must hide alice's memory"
1972        );
1973    }
1974
1975    #[tokio::test]
1976    async fn list_memories_ranks_by_importance_then_recency() {
1977        let store = fresh().await;
1978        let ns = NamespaceId::new("ns_test");
1979        // Three memories with varying importance + recency.
1980        let mut hi = sample_memory("hi", "important");
1981        hi.importance = 0.9;
1982        hi.created_at = "100".into();
1983        let mut mid = sample_memory("mid", "medium");
1984        mid.importance = 0.5;
1985        mid.created_at = "200".into(); // newer but lower importance
1986        let mut lo = sample_memory("lo", "low");
1987        lo.importance = 0.9;
1988        lo.created_at = "300".into(); // same importance as hi, newer
1989        store.store_memory(&ns, hi).await.unwrap();
1990        store.store_memory(&ns, mid).await.unwrap();
1991        store.store_memory(&ns, lo).await.unwrap();
1992
1993        let list = store.list_memories(&ns, 10).await.expect("list");
1994        // importance DESC first: the two 0.9s before the 0.5.
1995        // Among the 0.9s, created_at DESC: lo (300) before hi (100).
1996        assert_eq!(list.len(), 3);
1997        assert_eq!(list[0].id.0, "lo"); // 0.9, 300
1998        assert_eq!(list[1].id.0, "hi"); // 0.9, 100
1999        assert_eq!(list[2].id.0, "mid"); // 0.5
2000    }
2001
2002    // ===== knowledge graph =====
2003
2004    #[tokio::test]
2005    async fn add_triple_then_query_entity() {
2006        let store = fresh().await;
2007        let ns = NamespaceId::new("ns_kg");
2008        let t = store
2009            .add_triple(
2010                &ns,
2011                EntityId::new("Ijima"),
2012                "depends_on",
2013                EntityId::new("SurrealDB"),
2014                Some("100"),
2015                1.0,
2016                None,
2017            )
2018            .await
2019            .expect("add_triple");
2020        assert_eq!(t.subject.as_str(), "Ijima");
2021        assert_eq!(t.predicate, "depends_on");
2022        assert_eq!(t.object.as_str(), "SurrealDB");
2023        assert!(t.valid_to.is_none()); // current
2024
2025        let rec = store
2026            .query_entity(&ns, &EntityId::new("Ijima"))
2027            .await
2028            .expect("query");
2029        assert!(rec.entity.is_some());
2030        assert_eq!(rec.outgoing.len(), 1);
2031        assert_eq!(rec.outgoing[0].object.as_str(), "SurrealDB");
2032        assert!(rec.incoming.is_empty());
2033
2034        // From the object's side, it's incoming.
2035        let rec = store
2036            .query_entity(&ns, &EntityId::new("SurrealDB"))
2037            .await
2038            .expect("query");
2039        assert_eq!(rec.incoming.len(), 1);
2040        assert!(rec.outgoing.is_empty());
2041    }
2042
2043    #[tokio::test]
2044    async fn invalidate_triple_sets_valid_to() {
2045        let store = fresh().await;
2046        let ns = NamespaceId::new("ns_kg");
2047        store
2048            .add_triple(
2049                &ns,
2050                EntityId::new("a"),
2051                "uses",
2052                EntityId::new("b"),
2053                Some("100"),
2054                1.0,
2055                None,
2056            )
2057            .await
2058            .unwrap();
2059        store
2060            .invalidate_triple(&ns, "a:uses:b")
2061            .await
2062            .expect("invalidate");
2063        let found = store
2064            .find_triples(&ns, None, Some("uses"), None)
2065            .await
2066            .expect("find");
2067        assert_eq!(found.len(), 1);
2068        assert!(found[0].valid_to.is_some(), "valid_to must be set");
2069    }
2070
2071    #[tokio::test]
2072    async fn knowledge_stats_counts_entities_and_triples() {
2073        let store = fresh().await;
2074        let ns = NamespaceId::new("ns_kg");
2075        store
2076            .add_triple(
2077                &ns,
2078                EntityId::new("a"),
2079                "x",
2080                EntityId::new("b"),
2081                None,
2082                1.0,
2083                None,
2084            )
2085            .await
2086            .unwrap();
2087        store
2088            .add_triple(
2089                &ns,
2090                EntityId::new("a"),
2091                "y",
2092                EntityId::new("c"),
2093                None,
2094                1.0,
2095                None,
2096            )
2097            .await
2098            .unwrap();
2099        let stats = store.knowledge_stats(&ns).await.expect("stats");
2100        assert_eq!(stats.entities, 3); // a, b, c
2101        assert_eq!(stats.triples, 2);
2102    }
2103
2104    #[tokio::test]
2105    async fn kg_namespace_isolation() {
2106        let store = fresh().await;
2107        let a = NamespaceId::new("ns_a");
2108        let b = NamespaceId::new("ns_b");
2109        store
2110            .add_triple(
2111                &a,
2112                EntityId::new("x"),
2113                "uses",
2114                EntityId::new("y"),
2115                None,
2116                1.0,
2117                None,
2118            )
2119            .await
2120            .unwrap();
2121        // b sees nothing.
2122        let stats = store.knowledge_stats(&b).await.expect("stats");
2123        assert_eq!(stats.triples, 0);
2124        let rec = store
2125            .query_entity(&b, &EntityId::new("x"))
2126            .await
2127            .expect("query");
2128        assert!(rec.outgoing.is_empty());
2129    }
2130
2131    #[tokio::test]
2132    async fn store_stats_counts_across_namespaces() {
2133        let store = fresh().await;
2134        store
2135            .store_memory(&NamespaceId::new("ns_a"), sample_memory("m1", "x"))
2136            .await
2137            .unwrap();
2138        store
2139            .store_memory(&NamespaceId::new("ns_a"), sample_memory("m2", "y"))
2140            .await
2141            .unwrap();
2142        store
2143            .store_memory(&NamespaceId::new("ns_b"), sample_memory("m3", "z"))
2144            .await
2145            .unwrap();
2146        let stats = store.store_stats().await.expect("stats");
2147        assert_eq!(stats.total_memories, 3);
2148        assert_eq!(stats.namespaces.len(), 2);
2149        let ns_a = stats
2150            .namespaces
2151            .iter()
2152            .find(|n| n.namespace == "ns_a")
2153            .expect("ns_a");
2154        assert_eq!(ns_a.memories, 2);
2155    }
2156
2157    /// Builds a memory with explicit project/topic (sample_memory hardcodes them).
2158    fn mem_in(id: &str, project: &str, topic: &str, content: &str) -> Memory {
2159        Memory {
2160            id: MemoryId(id.into()),
2161            content: content.into(),
2162            project: project.into(),
2163            topic: topic.into(),
2164            source: MemorySource::Explicit,
2165            harness: Harness::Pi,
2166            session_id: None,
2167            origin: InstanceId::local(),
2168            authority: AuthorityScope::local(),
2169            importance: 0.5,
2170            created_at: "0".into(),
2171        }
2172    }
2173
2174    #[tokio::test]
2175    async fn palace_organization_rooms_taxonomy_graph_tunnel() {
2176        let store = fresh().await;
2177        let ns = NamespaceId::new("ns_palace");
2178        // Two projects sharing the topic "auth" (a tunnel), plus a
2179        // project-only topic.
2180        store
2181            .store_memory(&ns, mem_in("m1", "ijima", "auth", "use schubert"))
2182            .await
2183            .unwrap();
2184        store
2185            .store_memory(&ns, mem_in("m2", "ijima", "auth", "proof tokens"))
2186            .await
2187            .unwrap();
2188        store
2189            .store_memory(&ns, mem_in("m3", "ijima", "store", "surrealdb"))
2190            .await
2191            .unwrap();
2192        store
2193            .store_memory(&ns, mem_in("m4", "karpal", "auth", "workspace manifest"))
2194            .await
2195            .unwrap();
2196        store
2197            .store_memory(&ns, mem_in("m5", "karpal", "build", "ci"))
2198            .await
2199            .unwrap();
2200
2201        // list_rooms — all rooms, count desc.
2202        let rooms = store.list_rooms(&ns, None, 100).await.expect("rooms");
2203        assert_eq!(rooms.len(), 4);
2204        assert_eq!(rooms[0].count, 2); // ijima/auth
2205        assert_eq!(rooms[0].project, "ijima");
2206        assert_eq!(rooms[0].topic, "auth");
2207
2208        // list_rooms filtered to karpal.
2209        let karpal_rooms = store
2210            .list_rooms(&ns, Some("karpal"), 100)
2211            .await
2212            .expect("rooms karpal");
2213        assert_eq!(karpal_rooms.len(), 2);
2214        assert!(karpal_rooms.iter().all(|r| r.project == "karpal"));
2215
2216        // taxonomy — two projects.
2217        let taxons = store.taxonomy(&ns).await.expect("taxonomy");
2218        assert_eq!(taxons.len(), 2);
2219        let ijima_t = taxons
2220            .iter()
2221            .find(|t| t.project == "ijima")
2222            .expect("ijima taxon");
2223        assert_eq!(ijima_t.total, 3);
2224        assert_eq!(ijima_t.rooms.len(), 2); // auth, store
2225
2226        // palace_graph — projects + the auth tunnel between ijima & karpal.
2227        let graph = store.palace_graph(&ns).await.expect("graph");
2228        assert_eq!(graph.projects.len(), 2);
2229        assert!(graph.projects.contains(&"ijima".to_string()));
2230        assert!(graph.projects.contains(&"karpal".to_string()));
2231        let auth_tunnel = graph
2232            .tunnels
2233            .iter()
2234            .find(|t| t.topic == "auth")
2235            .expect("auth tunnel");
2236        assert_eq!(auth_tunnel.count_a, 2); // ijima
2237        assert_eq!(auth_tunnel.count_b, 1); // karpal
2238
2239        // traverse_tunnel — memories from both sides.
2240        let trav = store
2241            .traverse_tunnel(&ns, "auth", "ijima", "karpal", 10)
2242            .await
2243            .expect("traverse");
2244        assert_eq!(trav.memories_a.len(), 2);
2245        assert_eq!(trav.memories_b.len(), 1);
2246        assert!(trav.memories_a.iter().all(|m| m.project == "ijima"));
2247        assert!(trav.memories_b.iter().all(|m| m.project == "karpal"));
2248    }
2249
2250    #[tokio::test]
2251    async fn mining_queue_enqueue_list_accept_reject() {
2252        let store = fresh().await;
2253        let ns = NamespaceId::new("ns_mining");
2254        let other = NamespaceId::new("ns_other");
2255
2256        // Enqueue two PendingReview extractions.
2257        let q1 = store
2258            .enqueue_extraction(&ns, sample_memory("m1", "decided to use surrealdb"), 0.6)
2259            .await
2260            .expect("enqueue1");
2261        let q2 = store
2262            .enqueue_extraction(&ns, sample_memory("m2", "see https://example.com"), 0.55)
2263            .await
2264            .expect("enqueue2");
2265
2266        // list_pending sees both (newest first).
2267        let pending = store.list_pending(&ns, 10).await.expect("list");
2268        assert_eq!(pending.len(), 2);
2269        assert_eq!(pending[0].id, q2); // newest first
2270        assert_eq!(pending[1].confidence, 0.6);
2271
2272        // Namespace isolation: other ns sees nothing.
2273        let cross = store.list_pending(&other, 10).await.expect("list other");
2274        assert!(cross.is_empty());
2275
2276        // Accept q1 → promotes to palace + removes from queue.
2277        let accepted = store.accept_extraction(&ns, &q1).await.expect("accept");
2278        assert_eq!(accepted.memory_id.0, "m1");
2279        let after_accept = store.list_pending(&ns, 10).await.expect("list");
2280        assert_eq!(after_accept.len(), 1);
2281        // The promoted memory is now in the palace.
2282        let promoted = store
2283            .recall_memory(&ns, &MemoryId("m1".into()))
2284            .await
2285            .expect("recall");
2286        assert!(promoted.is_some());
2287
2288        // Reject q2 → drops without promoting.
2289        store.reject_extraction(&ns, &q2).await.expect("reject");
2290        let after_reject = store.list_pending(&ns, 10).await.expect("list");
2291        assert!(after_reject.is_empty());
2292        let not_promoted = store
2293            .recall_memory(&ns, &MemoryId("m2".into()))
2294            .await
2295            .expect("recall");
2296        assert!(not_promoted.is_none());
2297
2298        // Cross-namespace accept fails (queue entry is in ns, not other).
2299        let q3 = store
2300            .enqueue_extraction(&ns, sample_memory("m3", "cross test"), 0.5)
2301            .await
2302            .expect("enqueue3");
2303        let cross_accept = store.accept_extraction(&other, &q3).await;
2304        assert!(cross_accept.is_err());
2305    }
2306
2307    #[tokio::test]
2308    async fn persistent_store_survives_reopen() {
2309        let dir = std::env::temp_dir().join(format!(
2310            "ijima-persist-test-{}",
2311            std::time::SystemTime::now()
2312                .duration_since(std::time::UNIX_EPOCH)
2313                .unwrap()
2314                .as_nanos()
2315        ));
2316        let ns = NamespaceId::new("ns_persist");
2317
2318        // Phase 1 runs in its own task: surrealdb 3's surrealkv engine
2319        // holds a directory LOCK until the datastore's background tasks
2320        // drop (which needs the task boundary + a yield, not just
2321        // `drop(db)`). Production restarts are cross-process, so the OS
2322        // releases the lock; same-process reopens need this pattern.
2323        {
2324            let dir = dir.clone();
2325            let ns = ns.clone();
2326            tokio::spawn(async move {
2327                let store = SurrealStore::open_persistent(&dir).await.expect("open");
2328                store
2329                    .store_memory(&ns, sample_memory("mem_p", "survives restart"))
2330                    .await
2331                    .expect("store");
2332            })
2333            .await
2334            .expect("phase 1 task");
2335            for _ in 0..10 {
2336                tokio::task::yield_now().await;
2337            }
2338            tokio::time::sleep(std::time::Duration::from_millis(200)).await;
2339        }
2340
2341        // A fresh instance pointing at the same path must see the memory.
2342        let store = SurrealStore::open_persistent(&dir).await.expect("reopen");
2343        let got = store
2344            .recall_memory(&ns, &MemoryId("mem_p".into()))
2345            .await
2346            .expect("recall")
2347            .expect("memory must survive reopen");
2348        assert_eq!(got.content, "survives restart");
2349    }
2350
2351    #[tokio::test]
2352    async fn export_to_writes_sql_dump_to_file() {
2353        let dir = std::env::temp_dir().join(format!(
2354            "ijima-export-test-{}",
2355            std::time::SystemTime::now()
2356                .duration_since(std::time::UNIX_EPOCH)
2357                .unwrap()
2358                .as_nanos()
2359        ));
2360        std::fs::create_dir_all(&dir).expect("create dir");
2361        let db_path = dir.join("ijima.db");
2362        let store = crate::SurrealStore::open_persistent(&db_path)
2363            .await
2364            .expect("open");
2365        let ns = NamespaceId::new("ns_export");
2366        store
2367            .store_memory(&ns, sample_memory("mem_x", "export this memory"))
2368            .await
2369            .expect("store");
2370
2371        let out = dir.join("dump.surql");
2372        store.export_to(&out).await.expect("export");
2373        let contents = std::fs::read_to_string(&out).expect("read");
2374        assert!(contents.contains("export this memory"));
2375        // Clean up.
2376        let _ = std::fs::remove_dir_all(&dir);
2377    }
2378}
2379
2380// surrealdb 3: SELECT FROM a never-written table is a hard error (v2
2381// returned an empty set). Boot-time reads — revocation hydration,
2382// /status counts — hit fresh databases, so INDEX_DDL defines every table
2383// up front. This pins that a fresh store reads cleanly before any write.
2384#[tokio::test]
2385async fn fresh_store_reads_succeed_before_any_write() {
2386    let store = SurrealStore::open_embedded().await.expect("open");
2387    let revs = store.list_revocations().await.expect("fresh revocations");
2388    assert!(revs.is_empty());
2389    let stats = store.store_stats().await.expect("fresh stats");
2390    assert_eq!(stats.total_memories, 0);
2391    assert!(stats.namespaces.is_empty());
2392}
2393
2394// WS3 org walls: membership CRUD contract — idempotent grant, oldest-
2395// first listing, revoke-absent-is-ok, and the hot-path boolean.
2396#[tokio::test]
2397async fn namespace_membership_crud_round_trips() {
2398    let store = SurrealStore::open_embedded().await.expect("open");
2399    let ns = NamespaceId::new("ns_ia_shared");
2400    let mk = |principal: &str, at: u64| NamespaceMembership {
2401        namespace: "ns_ia_shared".into(),
2402        principal: principal.into(),
2403        granted_at_unix: at,
2404        granted_by: "root".into(),
2405    };
2406
2407    // Nobody is a member yet.
2408    assert!(
2409        !store
2410            .is_namespace_member(&ns, "elliott")
2411            .await
2412            .expect("check")
2413    );
2414
2415    store
2416        .grant_namespace_membership(mk("elliott", 100))
2417        .await
2418        .expect("grant");
2419    store
2420        .grant_namespace_membership(mk("sara", 200))
2421        .await
2422        .expect("grant");
2423    // Idempotent re-grant refreshes, never duplicates.
2424    store
2425        .grant_namespace_membership(mk("elliott", 300))
2426        .await
2427        .expect("re-grant");
2428
2429    assert!(
2430        store
2431            .is_namespace_member(&ns, "elliott")
2432            .await
2433            .expect("check")
2434    );
2435    let members = store.list_namespace_members(&ns).await.expect("list");
2436    assert_eq!(members.len(), 2, "re-grant must not duplicate: {members:?}");
2437    // Oldest grant first (elliott's original 100 predates sara's 200;
2438    // the refreshed row keeps key order by granted_at).
2439    assert_eq!(members[0].principal, "sara");
2440    assert_eq!(members[1].principal, "elliott");
2441    assert_eq!(members[1].granted_at_unix, 300, "re-grant refreshes");
2442
2443    store
2444        .revoke_namespace_membership(&ns, "elliott")
2445        .await
2446        .expect("revoke");
2447    assert!(
2448        !store
2449            .is_namespace_member(&ns, "elliott")
2450            .await
2451            .expect("check")
2452    );
2453    // Revoking an absent membership is not an error.
2454    store
2455        .revoke_namespace_membership(&ns, "ghost")
2456        .await
2457        .expect("revoke absent");
2458    // Membership is per-namespace.
2459    assert!(
2460        !store
2461            .is_namespace_member(&NamespaceId::new("ns_kellas_shared"), "sara")
2462            .await
2463            .expect("check")
2464    );
2465}
2466
2467// The revocation round-trip is covered full-stack in api::tests
2468// (token_revocation_kills_the_bearer_immediately); here we pin the store
2469// contract directly: upsert idempotence + oldest-first ordering.
2470#[tokio::test]
2471async fn revocation_upsert_and_ordering() {
2472    use ijima_core::TokenRevocation;
2473    let store = SurrealStore::open_embedded().await.expect("open");
2474    let mk = |hash: &str, at: u64| TokenRevocation {
2475        token_hash: hash.to_string(),
2476        revoked_at_unix: at,
2477        reason: None,
2478    };
2479    store.revoke_token(mk("bbb", 20)).await.expect("revoke b");
2480    store.revoke_token(mk("aaa", 10)).await.expect("revoke a");
2481    // Idempotent re-upsert of the same hash.
2482    store
2483        .revoke_token(mk("bbb", 20))
2484        .await
2485        .expect("re-revoke b");
2486    let revs = store.list_revocations().await.expect("list");
2487    assert_eq!(revs.len(), 2, "upsert keeps one row per hash");
2488    assert_eq!(revs[0].token_hash, "aaa", "oldest first");
2489    assert_eq!(revs[1].token_hash, "bbb");
2490}