Skip to main content

khive_runtime/
runtime.rs

1//! KhiveRuntime — composable handle to all storage capabilities.
2//!
3//! `RuntimeConfig`, `BackendId`, `NamespaceToken`, and embedding model helpers
4//! live in `super::config` and are re-exported from here.
5
6use std::collections::{HashMap, HashSet};
7use std::future::Future;
8use std::sync::{Arc, OnceLock, RwLock};
9
10use khive_db::StorageBackend;
11#[cfg(test)]
12use khive_gate::AllowAllGate;
13use khive_gate::GateRequest;
14use khive_storage::types::{SqlStatement, SqlValue};
15use khive_storage::{
16    AttachmentStore, EntityStore, Event, EventStore, GraphStore, NoteStore, SqlAccess, VectorStore,
17};
18use khive_types::{EdgeEndpointRule, EventKind, Namespace, SubstrateKind};
19use lattice_embed::{EmbeddingModel, EmbeddingService};
20
21use crate::config::{
22    build_embedder_registry, parse_embedding_model_alias, register_configured_embedding_models,
23    sanitize_key, vec_model_key,
24};
25use crate::error::{RuntimeError, RuntimeResult};
26use crate::pack::KindHook;
27
28#[cfg(all(test, target_os = "macos"))]
29const IN_PROCESS_TEST_NOFILE_LIMIT: libc::rlim_t = 4096;
30#[cfg(all(test, target_os = "macos"))]
31static IN_PROCESS_TEST_NOFILE_INIT: std::sync::Once = std::sync::Once::new();
32
33#[cfg(all(test, target_os = "macos"))]
34fn ensure_in_process_test_nofile_limit() {
35    IN_PROCESS_TEST_NOFILE_INIT.call_once(|| {
36        let mut limits = libc::rlimit {
37            rlim_cur: 0,
38            rlim_max: 0,
39        };
40        // SAFETY: `limits` is writable, and only this test binary's soft
41        // limit may change; the inherited hard limit is preserved.
42        assert_eq!(unsafe { libc::getrlimit(libc::RLIMIT_NOFILE, &mut limits) }, 0);
43        assert!(
44            limits.rlim_max >= IN_PROCESS_TEST_NOFILE_LIMIT,
45            "in-process SQLite tests require a hard open-file limit of at least {IN_PROCESS_TEST_NOFILE_LIMIT}"
46        );
47        if limits.rlim_cur < IN_PROCESS_TEST_NOFILE_LIMIT {
48            limits.rlim_cur = IN_PROCESS_TEST_NOFILE_LIMIT;
49            // SAFETY: the new soft limit does not exceed the observed hard
50            // limit, which is left unchanged.
51            assert_eq!(unsafe { libc::setrlimit(libc::RLIMIT_NOFILE, &limits) }, 0);
52        }
53    });
54}
55
56tokio::task_local! {
57    static REQUEST_EMBEDDER_EXCLUSIONS: Arc<HashSet<String>>;
58}
59
60/// Run one request with daemon-only embedding models excluded from registry access.
61pub fn scope_request_embedder_exclusions<F>(
62    excluded: Vec<String>,
63    future: F,
64) -> impl Future<Output = F::Output>
65where
66    F: Future,
67{
68    REQUEST_EMBEDDER_EXCLUSIONS.scope(Arc::new(excluded.into_iter().collect()), future)
69}
70
71/// Carry the current request's embedder exclusions into a spawned task.
72pub fn inherit_request_embedder_scope<F>(future: F) -> impl Future<Output = F::Output>
73where
74    F: Future,
75{
76    let exclusions = REQUEST_EMBEDDER_EXCLUSIONS.try_with(Arc::clone).ok();
77    async move {
78        match exclusions {
79            Some(exclusions) => REQUEST_EMBEDDER_EXCLUSIONS.scope(exclusions, future).await,
80            None => future.await,
81        }
82    }
83}
84
85fn request_excludes_embedder(name: &str) -> bool {
86    let canonical = parse_embedding_model_alias(name)
87        .map(|model| model.to_string())
88        .unwrap_or_else(|| name.to_string());
89    REQUEST_EMBEDDER_EXCLUSIONS
90        .try_with(|excluded| excluded.contains(&canonical))
91        .unwrap_or(false)
92}
93
94/// Callback type for pack-installed entity-type validators.
95///
96/// Receives `(kind, entity_type)` and returns the normalised type string,
97/// or `RuntimeError::InvalidInput` if the type is not registered for that kind.
98/// When `entity_type` is `None`, the implementation must return `Ok(None)`.
99pub type EntityTypeValidatorFn =
100    Arc<dyn Fn(&str, Option<&str>) -> Result<Option<String>, RuntimeError> + Send + Sync>;
101
102/// Pack-aggregated entity-kind update hooks: `(entity kind, hook)` for every
103/// kind whose owning pack both declares it and registers a `KindHook`.
104///
105/// Named rather than written inline because the runtime stores it behind an
106/// `Arc<RwLock<..>>` and passes it across the transport boundary, so the bare
107/// form appears three times and reads as noise at each one.
108pub type EntityKindHooks = Vec<(String, Arc<dyn KindHook>)>;
109
110/// Callback type for a pack-installed note-mutation hook.
111///
112/// Invoked by `update_note` (when the note's text/embedding actually
113/// changed) and `delete_note` (soft or hard) with `(note_kind, note_id)`,
114/// after the mutation has been durably applied. Returns a boxed future so
115/// the hook can await async cache-invalidation work (e.g.
116/// `khive-pack-memory`'s ANN warm-cache generation bump) without
117/// `khive-runtime` depending on any pack crate: dependencies point the
118/// other way, so the runtime exposes an extension point and the pack
119/// installs into it, same shape as `EntityTypeValidatorFn`, just async.
120pub type NoteMutationHookFn = Arc<
121    dyn Fn(String, uuid::Uuid) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>
122        + Send
123        + Sync,
124>;
125
126/// Callback type for a pack-installed note-write validator.
127///
128/// The pack that owns a note kind carrying derivable identity installs one so
129/// that the identity is a function of the authorization token rather than of
130/// caller input, on every write path including direct callers that bypass the
131/// handler layer — same rationale as [`EntityTypeValidatorFn`], which exists
132/// for exactly that reason on the entity side.
133///
134/// Kinds the installing pack does not own must be returned unchanged: the slot
135/// is single-occupancy (like `note_mutation_hook`), so a validator that
136/// rewrote foreign kinds would silently govern every other pack's notes.
137pub type NoteWriteValidatorFn = Arc<
138    dyn Fn(&str, &str, Option<serde_json::Value>) -> Result<Option<serde_json::Value>, RuntimeError>
139        + Send
140        + Sync,
141>;
142
143/// Immutable identity for a non-text vector store owned by a pack consumer.
144///
145/// This does not register an [`crate::EmbedderProvider`]. It gives a pack that
146/// performs its own governed inference a narrow path to a namespace-scoped
147/// Khive vector table while keeping model-key and dimension validation at the
148/// runtime boundary.
149#[derive(Clone, Debug, Eq, PartialEq)]
150pub struct NamedVectorIdentity {
151    model_key: String,
152    model_name: String,
153    dimensions: usize,
154}
155
156struct CachedNamedVectorStores {
157    identity: NamedVectorIdentity,
158    by_namespace: HashMap<String, Arc<dyn VectorStore>>,
159}
160
161fn check_cached_named_vector_identity(
162    cached: &NamedVectorIdentity,
163    requested: &NamedVectorIdentity,
164) -> RuntimeResult<()> {
165    if cached.dimensions() != requested.dimensions() {
166        return Err(RuntimeError::InvalidInput(format!(
167            "named vector model_key {:?} is already bound to {} dimensions, expected {}",
168            requested.model_key(),
169            cached.dimensions(),
170            requested.dimensions()
171        )));
172    }
173    if cached.model_name() != requested.model_name() {
174        return Err(RuntimeError::InvalidInput(format!(
175            "named vector model_key {:?} is already bound to a different active model identity",
176            requested.model_key()
177        )));
178    }
179    Ok(())
180}
181
182impl NamedVectorIdentity {
183    const MAX_MODEL_KEY_BYTES: usize = 128;
184    const MAX_MODEL_NAME_BYTES: usize = 512;
185
186    /// Validate and construct a named vector identity.
187    pub fn new(
188        model_key: impl Into<String>,
189        model_name: impl Into<String>,
190        dimensions: usize,
191    ) -> RuntimeResult<Self> {
192        let model_key = model_key.into();
193        let model_name = model_name.into();
194        if model_key.is_empty()
195            || model_key.len() > Self::MAX_MODEL_KEY_BYTES
196            || !model_key
197                .chars()
198                .all(|c| c.is_ascii_alphanumeric() || c == '_')
199        {
200            return Err(RuntimeError::InvalidInput(format!(
201                "named vector model_key must be 1..={} bytes of ASCII alphanumeric/underscore",
202                Self::MAX_MODEL_KEY_BYTES
203            )));
204        }
205        if model_name.trim().is_empty()
206            || model_name.trim() != model_name
207            || model_name.len() > Self::MAX_MODEL_NAME_BYTES
208        {
209            return Err(RuntimeError::InvalidInput(format!(
210                "named vector model_name must be 1..={} bytes with no surrounding whitespace",
211                Self::MAX_MODEL_NAME_BYTES
212            )));
213        }
214        if !(1..=8192).contains(&dimensions) {
215            return Err(RuntimeError::InvalidInput(format!(
216                "named vector dimensions must be in 1..=8192, got {dimensions}"
217            )));
218        }
219        Ok(Self {
220            model_key,
221            model_name,
222            dimensions,
223        })
224    }
225
226    pub fn model_key(&self) -> &str {
227        &self.model_key
228    }
229
230    pub fn model_name(&self) -> &str {
231        &self.model_name
232    }
233
234    pub fn dimensions(&self) -> usize {
235        self.dimensions
236    }
237}
238
239pub use crate::config::{
240    assert_captured_db_anchor_consistent, assert_db_anchor_consistent, expand_tilde,
241    parse_pack_list, resolve_db_anchor, resolve_project_actor_id, runtime_config_from_khive_config,
242    BackendId, BackendIdError, NamespaceToken, RuntimeConfig,
243};
244
245// ---- KhiveRuntime ----
246
247/// Composable runtime handle used by the MCP server.
248///
249/// Wraps a `StorageBackend` and provides namespace-scoped accessor methods
250/// Snapshot of the main runtime's embedder wiring (registry handle, default
251/// name, and the config model fields the default-resolution path reads).
252/// Carried by secondary-pack runtimes; consumed by [`KhiveRuntime::core`].
253#[derive(Clone)]
254struct CoreEmbedderState {
255    registry: Arc<std::sync::RwLock<crate::embedder_registry::EmbedderRegistry>>,
256    default_embedder_name: Arc<str>,
257    embedding_model: Option<EmbeddingModel>,
258    additional_embedding_models: Vec<EmbeddingModel>,
259}
260
261/// for each storage capability, plus a lazily-loaded embedder.
262#[derive(Clone)]
263pub struct KhiveRuntime {
264    backend: Arc<StorageBackend>,
265    /// Successful named-vector bindings and their namespace-scoped stores.
266    /// Shared by runtime clones so repeated reads do not enter the writer or
267    /// rescan the vector table after the first verified binding.
268    named_vector_stores: Arc<RwLock<HashMap<String, CachedNamedVectorStores>>>,
269    /// The main backend's cache, used when a secondary runtime creates a
270    /// `core()` handle. It must never reuse a secondary backend's store.
271    core_named_vector_stores: Option<Arc<RwLock<HashMap<String, CachedNamedVectorStores>>>>,
272    /// When `Some`, holds the main backend so that `core()` can return a
273    /// main-bound runtime handle without constructing a new connection.
274    /// `None` when this runtime is already bound to the main backend.
275    core_backend: Option<Arc<StorageBackend>>,
276    config: RuntimeConfig,
277    /// ADR-118 exact-leg policy, sampled once at runtime construction.
278    /// Request-time memory/knowledge serving must never re-read the process
279    /// environment because tests and embedded runtimes share one process.
280    ann_fresh_tail_enabled: bool,
281    /// Pack-extensible embedder registry.
282    ///
283    /// Shared across clones via `Arc<RwLock<_>>` so that
284    /// [`register_embedder`](Self::register_embedder) after clone is visible
285    /// to all handles. Built-in lattice models are pre-registered during
286    /// construction; packs may add more via [`PackRuntime::register_embedders`].
287    embedder_registry: Arc<std::sync::RwLock<crate::embedder_registry::EmbedderRegistry>>,
288    default_embedder_name: Arc<str>,
289    /// The MAIN runtime's embedder wiring, carried by secondary-backend
290    /// runtimes so that `core()`-routed writes embed with the main runtime's
291    /// models even when this pack's own registry is empty (`no_embed`).
292    /// `None` on the main runtime, and on secondaries wired before the boot
293    /// path calls [`with_core_embedders_from`](Self::with_core_embedders_from)
294    /// — `core()` then falls back to this runtime's own embedder state.
295    core_embedders: Option<CoreEmbedderState>,
296    /// Pack-extensible edge endpoint rules. Shared across clones
297    /// via `Arc<RwLock<_>>`; installed once by the transport after the
298    /// `VerbRegistry` is built. Empty until installed
299    edge_rules: Arc<RwLock<Vec<EdgeEndpointRule>>>,
300    /// Pack-aggregated valid entity and note kind strings.
301    ///
302    /// Installed by the transport layer after building the `VerbRegistry`.
303    /// When non-empty, `create_entity`, `create_note_inner`, and `import_kg`
304    /// reject kinds not in these sets. When empty (no packs loaded, e.g.
305    /// bare runtime in unit tests), kind validation is skipped — the pack
306    /// handler layer is the primary enforcement point.
307    valid_entity_kinds: Arc<RwLock<Vec<String>>>,
308    valid_note_kinds: Arc<RwLock<Vec<String>>>,
309    /// Pack-installed entity-type validator.
310    ///
311    /// When `Some`, `create_many` calls this function to validate and normalise
312    /// each `(kind, entity_type)` pair before writing. When `None` (bare runtime
313    /// without packs), entity-type validation is skipped — the pack handler layer
314    /// is the primary enforcement point, same as for `valid_entity_kinds`.
315    entity_type_validator: Arc<RwLock<Option<EntityTypeValidatorFn>>>,
316    /// Pack-installed note-mutation hook.
317    ///
318    /// When `Some`, `update_note` (on text change) and `delete_note` (soft
319    /// or hard) call this after the mutation is durably applied, so a pack
320    /// that caches derived state keyed by note content (e.g. `khive-pack-memory`'s
321    /// warm ANN index) can invalidate/advance its own generation counter even
322    /// when the mutation arrived through a different pack's verb (e.g. KG's
323    /// `update`/`delete` on a `kind="memory"` note) that has no dependency on
324    /// the reacting pack. `None` when no pack installs one (bare runtime, or
325    /// no pack cares about note-mutation notifications) — the call becomes a
326    /// no-op check of an `Option`.
327    note_mutation_hook: Arc<RwLock<Option<NoteMutationHookFn>>>,
328    /// Pack-installed note-write validator.
329    ///
330    /// When `Some`, every runtime note-materialisation site that accepts
331    /// caller-supplied `properties` routes them through this function before
332    /// the `Note` is built, so a pack-owned identity property is derived from
333    /// the authorization token instead of trusted from caller input. `None`
334    /// on a bare runtime (no packs) — the properties pass through unchanged.
335    note_write_validator: Arc<RwLock<Option<NoteWriteValidatorFn>>>,
336    /// Pack-installed entity-kind update-validation hooks (issue #2943).
337    ///
338    /// Every `(entity kind, hook)` pair for which an owning pack declares
339    /// the entity kind and registers a `KindHook`, aggregated once by
340    /// `VerbRegistry::entity_kind_hooks` and installed by the transport
341    /// after the registry is built — same timing and rationale as
342    /// `entity_type_validator`: `khive-runtime` does not hold a
343    /// `VerbRegistry`, so this is the extension point that lets
344    /// `prepare_guarded_entity_update` reach a pack's `KindHook` on the
345    /// generic entity `update` path, the counterpart to
346    /// `prepare_note_update_hook` on the note side. Empty until installed
347    /// (bare runtime, or no pack registers an entity-kind hook), which
348    /// leaves the dispatch a no-op.
349    entity_kind_hooks: Arc<RwLock<EntityKindHooks>>,
350    /// Pack-owned note kinds — every note kind declared by a pack other than
351    /// the generic-CRUD pack, installed by the transport from the registry
352    /// (see `VerbRegistry::pack_owned_note_kinds`). Records of these kinds are
353    /// maintained by their owning pack's own verbs, so `update`'s `properties`
354    /// patch is refused on them at the runtime layer and their owned identity
355    /// properties survive a `merge` unchanged. Empty until installed (bare
356    /// runtime), which leaves both rules inert.
357    pack_owned_note_kinds: Arc<RwLock<Vec<String>>>,
358    /// The immutable runtime-owned store/hydration pair (ADR-160 D3).
359    ///
360    /// Boot resolves one store and constructs exactly one shared hydrator,
361    /// then installs that same `Arc` on every runtime handle. The one-shot
362    /// slot rejects replacement so pack runtimes cannot silently split the
363    /// aggregate byte budget after startup. Bare runtimes leave it unset.
364    blob_hydrator: Arc<OnceLock<Arc<crate::blob::BlobHydrator>>>,
365    /// Pack-registered custom fusion executors (ADR-012), keyed by the name
366    /// carried in `FusionStrategy::Custom { name, .. }`.
367    ///
368    /// Unlike `entity_type_validator`/`note_mutation_hook` (single-occupancy —
369    /// one pack owns the slot), multiple packs each register their own named
370    /// strategy under this shared map, so it is keyed rather than a bare
371    /// `Option`. Empty until a pack calls
372    /// [`register_fusion_strategy`](Self::register_fusion_strategy); an
373    /// unregistered `Custom` name at dispatch time is
374    /// `RuntimeError::UnknownFusionStrategy`, never a silent fallback.
375    fusion_executors: Arc<RwLock<HashMap<String, Arc<dyn crate::fusion::FusionExecutor>>>>,
376}
377
378impl KhiveRuntime {
379    /// Create a new runtime with the given config.
380    ///
381    /// The config's `db_path` is used to open or create the SQLite backend.
382    /// This direct constructor is intended for fresh/current single-backend
383    /// databases and tests. Production and multi-backend hosts must use the
384    /// async khive-mcp/kkernel builders so secondary inventory and any
385    /// application-assisted V21 cutover complete before serving. The
386    /// [`from_backend`](Self::from_backend) seam is likewise only for an
387    /// already-prepared backend.
388    pub fn new(config: RuntimeConfig) -> RuntimeResult<Self> {
389        Self::new_with_file_backend(config, |path| StorageBackend::sqlite(path))
390    }
391
392    /// Construct a fixture runtime with a small concurrent reader pool.
393    #[cfg(any(test, feature = "test-internals"))]
394    pub fn new_for_test(config: RuntimeConfig) -> RuntimeResult<Self> {
395        Self::new_with_file_backend(config, |path| StorageBackend::sqlite_for_test(path))
396    }
397
398    fn new_with_file_backend(
399        config: RuntimeConfig,
400        open_file: impl FnOnce(&std::path::Path) -> Result<StorageBackend, khive_db::SqliteError>,
401    ) -> RuntimeResult<Self> {
402        #[cfg(all(test, target_os = "macos"))]
403        ensure_in_process_test_nofile_limit();
404        let backend = match &config.db_path {
405            Some(path) => {
406                if let Some(parent) = path.parent() {
407                    std::fs::create_dir_all(parent).ok();
408                }
409                open_file(path)?
410            }
411            None => StorageBackend::memory()?,
412        };
413        // Writable backends migrate before handlers touch the DB. A detected
414        // read-only snapshot is validated at the current schema version without
415        // attempting migration DDL.
416        let schema_version = backend.prepare_core_schema()?;
417        if schema_version < khive_db::migrations::ATTACHMENT_CUTOVER_VERSION {
418            return Err(khive_db::SqliteError::InvalidData(
419                "database requires the host application-assisted V21 attachment cutover; \
420                 start through khive-mcp/kkernel boot instead of constructing KhiveRuntime \
421                 directly"
422                    .into(),
423            )
424            .into());
425        }
426        if !backend.is_read_only() {
427            register_configured_embedding_models(&backend, &config)?;
428        }
429        Ok(Self::assemble_from_backend(Arc::new(backend), config))
430    }
431
432    /// Open a runtime for read-only inspection (no model registration, no DB creation).
433    ///
434    /// File-backed databases are opened with SQLite read-only/query-only flags
435    /// and must already be at this build's current schema version. No migrations
436    /// or configured-model registration writes are attempted. A `None` path
437    /// retains the historical ephemeral in-memory behavior for tests.
438    pub fn new_readonly(config: RuntimeConfig) -> RuntimeResult<Self> {
439        Self::new_readonly_with_file_backend(config, |path| StorageBackend::sqlite_read_only(path))
440    }
441
442    /// Construct a read-only fixture runtime with a small reader pool.
443    #[cfg(any(test, feature = "test-internals"))]
444    pub fn new_readonly_for_test(config: RuntimeConfig) -> RuntimeResult<Self> {
445        Self::new_readonly_with_file_backend(config, |path| {
446            StorageBackend::sqlite_read_only_for_test(path)
447        })
448    }
449
450    fn new_readonly_with_file_backend(
451        config: RuntimeConfig,
452        open_file: impl FnOnce(&std::path::Path) -> Result<StorageBackend, khive_db::SqliteError>,
453    ) -> RuntimeResult<Self> {
454        #[cfg(all(test, target_os = "macos"))]
455        ensure_in_process_test_nofile_limit();
456        let backend = match &config.db_path {
457            Some(path) => open_file(path)?,
458            None => StorageBackend::memory()?,
459        };
460        backend.prepare_core_schema()?;
461        Ok(Self::assemble_from_backend(Arc::new(backend), config))
462    }
463
464    /// Construct a runtime from an already-opened backend.
465    ///
466    /// This is a low-level, infallible assembly seam for already-prepared
467    /// multi-backend deployments. It does not inspect or migrate the V21
468    /// attachment-cutover state. Production hosts must first run the async
469    /// kkernel/khive-mcp coordinator and must not expose a server over a
470    /// pending or incomplete backend. Prefer [`Self::from_prepared_backend`]
471    /// when constructing one fallible host runtime.
472    ///
473    /// The returned runtime has `db_path = None` and `embedding_model = None`; all
474    /// storage access is through the provided `backend`. Set `backend_id` and
475    /// `default_namespace` via the config builder pattern if non-defaults are needed.
476    pub fn from_backend(backend: Arc<StorageBackend>, config: RuntimeConfig) -> Self {
477        if !backend.is_read_only() {
478            if let Err(err) = register_configured_embedding_models(&backend, &config) {
479                tracing::warn!(error = %err, "failed to register configured embedding models");
480            }
481        }
482        Self::assemble_from_backend(backend, config)
483    }
484
485    /// Construct a single-backend runtime after a host boot coordinator has
486    /// completed schema preparation and any application-assisted cutover.
487    ///
488    /// Unlike [`Self::from_backend`], configured embedding-model registration
489    /// is fallible here, preserving [`Self::new`]'s single-backend startup
490    /// semantics. This method never runs migrations itself.
491    pub fn from_prepared_backend(
492        backend: Arc<StorageBackend>,
493        config: RuntimeConfig,
494    ) -> RuntimeResult<Self> {
495        if backend.attachment_cutover_status()?
496            != khive_db::migrations::AttachmentCutoverStatus::Complete
497        {
498            return Err(khive_db::SqliteError::InvalidData(
499                "from_prepared_backend requires a complete V21 attachment cutover".into(),
500            )
501            .into());
502        }
503        if !backend.is_read_only() {
504            register_configured_embedding_models(&backend, &config)?;
505        }
506        Ok(Self::assemble_from_backend(backend, config))
507    }
508
509    fn assemble_from_backend(backend: Arc<StorageBackend>, config: RuntimeConfig) -> Self {
510        if config.backend_id.as_str() == BackendId::MAIN {
511            backend.pool().main_pool_generation();
512        }
513        let ann_fresh_tail_enabled = crate::config::ann_fresh_tail_enabled_from_env();
514        let (registry, default_embedder_name) = build_embedder_registry(&config);
515        Self {
516            backend,
517            named_vector_stores: Arc::new(RwLock::new(HashMap::new())),
518            core_named_vector_stores: None,
519            core_backend: None,
520            config,
521            ann_fresh_tail_enabled,
522            embedder_registry: Arc::new(std::sync::RwLock::new(registry)),
523            default_embedder_name,
524            core_embedders: None,
525            edge_rules: Arc::new(RwLock::new(Vec::new())),
526            valid_entity_kinds: Arc::new(RwLock::new(Vec::new())),
527            valid_note_kinds: Arc::new(RwLock::new(Vec::new())),
528            entity_type_validator: Arc::new(RwLock::new(None)),
529            note_mutation_hook: Arc::new(RwLock::new(None)),
530            note_write_validator: Arc::new(RwLock::new(None)),
531            entity_kind_hooks: Arc::new(RwLock::new(Vec::new())),
532            pack_owned_note_kinds: Arc::new(RwLock::new(Vec::new())),
533            blob_hydrator: Arc::new(OnceLock::new()),
534            fusion_executors: Arc::new(RwLock::new(HashMap::new())),
535        }
536    }
537
538    /// Wire this runtime as a secondary-backend runtime pointing at `core`.
539    ///
540    /// After this call, `self.core()` returns a handle to `core` rather than
541    /// cloning `self`. The caller (the boot path, not pack code) is responsible
542    /// for passing the correct main backend.
543    ///
544    /// Panics in debug builds if `self.config.backend_id == BackendId::MAIN`,
545    /// because the main runtime does not need a core pointer.
546    pub fn with_core_backend(mut self, core: Arc<StorageBackend>) -> Self {
547        debug_assert_ne!(
548            self.config.backend_id.as_str(),
549            BackendId::MAIN,
550            "with_core_backend must not be called on the main runtime"
551        );
552        core.pool().main_pool_generation();
553        if self.core_named_vector_stores.is_none() {
554            self.core_named_vector_stores = Some(Arc::new(RwLock::new(HashMap::new())));
555        }
556        self.core_backend = Some(core);
557        self
558    }
559
560    /// Carry the main runtime's embedder wiring for `core()`-routed writes.
561    ///
562    /// Boot-path companion to [`with_core_backend`](Self::with_core_backend).
563    /// Without it, `core()` shares this pack runtime's own embedder registry —
564    /// which under `[packs.<name>] no_embed = true` is empty, so core-routed
565    /// concept writes would silently skip embedding on the shared graph.
566    pub fn with_core_embedders_from(mut self, main: &KhiveRuntime) -> Self {
567        debug_assert!(
568            main.core_backend.is_none(),
569            "with_core_embedders_from takes the MAIN runtime"
570        );
571        self.core_embedders = Some(CoreEmbedderState {
572            registry: main.embedder_registry.clone(),
573            default_embedder_name: main.default_embedder_name.clone(),
574            embedding_model: main.config.embedding_model,
575            additional_embedding_models: main.config.additional_embedding_models.clone(),
576        });
577        self.core_named_vector_stores = Some(main.named_vector_stores.clone());
578        if Arc::ptr_eq(&self.backend, &main.backend) {
579            self.named_vector_stores = main.named_vector_stores.clone();
580        }
581        self
582    }
583
584    /// Return a runtime handle bound to the main (shared-graph) backend.
585    ///
586    /// When `self` is already the main runtime (`core_backend` is `None`),
587    /// this returns a clone of `self` — no new backend reference is acquired.
588    ///
589    /// When `self` is a secondary-backend runtime (`core_backend` is `Some`),
590    /// this returns a new `KhiveRuntime` backed by the main
591    /// `Arc<StorageBackend>` and sharing all registry state (`embedder_registry`,
592    /// `edge_rules`, `valid_entity_kinds`, `valid_note_kinds`,
593    /// `entity_type_validator`, `note_mutation_hook`, `entity_kind_hooks`) with `self`.
594    /// No database I/O occurs; no embedding models are reloaded.
595    ///
596    /// Use `core()` for notes and entities that must reside in the shared graph
597    /// so that `memory.recall`, cross-pack search, and `annotates` edges work.
598    /// Use `self` (or `self.sql()`) for pack-auxiliary bulk tables.
599    ///
600    /// Handlers that call `core()` more than once per request or loop should bind
601    /// `let core = self.core();` once and reuse it, since each call clones
602    /// `RuntimeConfig` (a heap-allocated struct containing `Vec<String>` fields).
603    pub fn core(&self) -> KhiveRuntime {
604        match &self.core_backend {
605            // A main-assigned pack runtime has no core pointer, but may still
606            // carry main's embedder wiring: with `no_embed` its OWN registry
607            // is empty, and core-routed concept writes must embed regardless
608            // of which backend the pack was assigned to.
609            None => match &self.core_embedders {
610                None => self.clone(),
611                Some(core_embedders) => {
612                    let mut core = self.clone();
613                    core.config.embedding_model = core_embedders.embedding_model;
614                    core.config.additional_embedding_models =
615                        core_embedders.additional_embedding_models.clone();
616                    core.embedder_registry = core_embedders.registry.clone();
617                    core.default_embedder_name = core_embedders.default_embedder_name.clone();
618                    core.core_embedders = None;
619                    core
620                }
621            },
622            Some(main_arc) => {
623                let mut core_config = self.config.clone();
624                core_config.backend_id = BackendId::main();
625                // Core-routed writes embed with the MAIN runtime's wiring when
626                // the boot path supplied it (see `with_core_embedders_from`);
627                // both the registry handle and the config model fields must
628                // come from main, since default-model resolution reads
629                // `config.embedding_model` (`resolve_embedding_model`).
630                let (embedder_registry, default_embedder_name) = match &self.core_embedders {
631                    Some(core_embedders) => {
632                        core_config.embedding_model = core_embedders.embedding_model;
633                        core_config.additional_embedding_models =
634                            core_embedders.additional_embedding_models.clone();
635                        (
636                            core_embedders.registry.clone(),
637                            core_embedders.default_embedder_name.clone(),
638                        )
639                    }
640                    None => (
641                        self.embedder_registry.clone(),
642                        self.default_embedder_name.clone(),
643                    ),
644                };
645                KhiveRuntime {
646                    backend: main_arc.clone(),
647                    named_vector_stores: self
648                        .core_named_vector_stores
649                        .clone()
650                        .unwrap_or_else(|| Arc::new(RwLock::new(HashMap::new()))),
651                    core_named_vector_stores: None,
652                    core_backend: None,
653                    config: core_config,
654                    ann_fresh_tail_enabled: self.ann_fresh_tail_enabled,
655                    embedder_registry,
656                    default_embedder_name,
657                    core_embedders: None,
658                    edge_rules: self.edge_rules.clone(),
659                    valid_entity_kinds: self.valid_entity_kinds.clone(),
660                    valid_note_kinds: self.valid_note_kinds.clone(),
661                    entity_type_validator: self.entity_type_validator.clone(),
662                    note_mutation_hook: self.note_mutation_hook.clone(),
663                    note_write_validator: self.note_write_validator.clone(),
664                    entity_kind_hooks: self.entity_kind_hooks.clone(),
665                    pack_owned_note_kinds: self.pack_owned_note_kinds.clone(),
666                    blob_hydrator: self.blob_hydrator.clone(),
667                    fusion_executors: self.fusion_executors.clone(),
668                }
669            }
670        }
671    }
672
673    /// Create an in-memory runtime (for tests and ephemeral use).
674    pub fn memory() -> RuntimeResult<Self> {
675        Self::new(RuntimeConfig {
676            db_path: None,
677            packs: vec!["kg".to_string()],
678            brain_profile: None,
679            actor_id: None,
680            ..RuntimeConfig::no_embeddings()
681        })
682    }
683
684    /// Return the [`BackendId`] for this runtime's backend.
685    ///
686    /// Used by `SubstrateCoordinator` in `kkernel`
687    /// to identify which backend owns a given node, and to detect cross-backend merges.
688    pub fn backend_id(&self) -> &BackendId {
689        &self.config.backend_id
690    }
691
692    /// Return the extra-visible namespaces assembled at config load.
693    ///
694    /// OSS dispatch uses this set to widen the default multi-record read scope
695    /// to `['local'] ∪ visible_namespaces`. Writes are unchanged: always
696    /// pinned to `'local'`. This set is also available as gate/cloud policy
697    /// input.
698    pub fn visible_namespaces(&self) -> &[Namespace] {
699        &self.config.visible_namespaces
700    }
701
702    /// Return a reference to the runtime config.
703    pub fn config(&self) -> &RuntimeConfig {
704        &self.config
705    }
706
707    /// Whether this runtime selects the vector arm for a hybrid search —
708    /// true exactly when a default embedding model is configured. Single
709    /// source of truth for the policy every fan-out and single-backend
710    /// dispatch path uses to report `arm_participation`/`vector_selected`.
711    pub fn vector_arm_selected(&self) -> bool {
712        self.config.embedding_model.is_some()
713    }
714
715    /// Return the immutable ADR-118 fresh-tail serving policy captured when
716    /// this runtime was constructed.
717    pub fn ann_fresh_tail_enabled(&self) -> bool {
718        self.ann_fresh_tail_enabled
719    }
720
721    /// Override ADR-118's fresh-tail serving policy for this runtime instance.
722    ///
723    /// This is primarily useful for embedded runtimes and deterministic tests:
724    /// it avoids mutating process-global environment state. Clones and `core()`
725    /// handles preserve the chosen value.
726    pub fn with_ann_fresh_tail_enabled(mut self, enabled: bool) -> Self {
727        self.ann_fresh_tail_enabled = enabled;
728        self
729    }
730
731    /// Return a reference to the underlying storage backend.
732    ///
733    /// This is an embedder/infrastructure surface (connection pools, schema
734    /// plans, diagnostics). Stores obtained from it are NOT wrapped by the
735    /// message-evidence policy that [`Self::notes`] enforces: an embedder
736    /// holding the backend already holds root-equivalent access to the
737    /// database file, so the policy boundary sits at the typed accessors
738    /// pack code uses, not here. Pack code must not take note stores from
739    /// this surface.
740    pub fn backend(&self) -> &StorageBackend {
741        &self.backend
742    }
743
744    /// Whether this runtime's bound backend is explicitly or filesystem-mode
745    /// detected read-only.
746    pub fn is_read_only(&self) -> bool {
747        self.backend.is_read_only()
748    }
749
750    /// Return the directory containing the backend's database file, or `None`
751    /// for an in-memory backend.
752    pub fn backend_data_dir(&self) -> Option<std::path::PathBuf> {
753        self.backend.data_dir()
754    }
755
756    /// Root directory for this database's ANN segment tree (`<db-file>.ann/`
757    /// beside the file), or `None` for an in-memory backend. Scoped to the
758    /// database file itself so two databases sharing a parent directory can
759    /// never adopt each other's segments.
760    pub fn backend_ann_root(&self) -> Option<std::path::PathBuf> {
761        self.backend.ann_root()
762    }
763
764    /// Writer-contention, graph-edge integrity, and WAL/checkpoint diagnostics
765    /// (ADR-091/ADR-135 operator surface): pooled writer and audit-failure
766    /// counters, build identity, duplicate edge-ID and list-ledger counts,
767    /// checkpoint counters, a PASSIVE checkpoint probe, WAL file size, and
768    /// explicitly qualified WAL-pin census. Not write-free: the
769    /// PASSIVE probe may backfill WAL frames into the database (normal
770    /// checkpoint I/O). It never changes logical state, escalates to TRUNCATE,
771    /// creates a missing database file, or deletes sidecar evidence — see
772    /// `khive_db::diagnostics` for the narrowings that make those claims hold.
773    ///
774    /// Always targets the *main* backend via [`Self::core`], regardless of
775    /// which backend this runtime handle is bound to, so a report never
776    /// describes a database this handle is not the canonical owner of.
777    pub async fn db_diagnostics(&self) -> RuntimeResult<khive_db::diagnostics::DbDiagnostics> {
778        // No `VerbRegistry` handle is reachable from a bare `KhiveRuntime`
779        // (the audit-batch seam is owned by whichever registry was built
780        // over this runtime's `EventStore`, not by the runtime itself), so
781        // the batch-health fields report unavailable with a reason here.
782        // Callers that hold the registry — e.g. the `db_diagnostics` verb
783        // handler — use `Self::db_diagnostics_with_audit_metrics` with
784        // `VerbRegistry::audit_batch_metrics()` instead.
785        self.db_diagnostics_with_audit_metrics(None).await
786    }
787
788    /// As [`Self::db_diagnostics`], but with the caller supplying the
789    /// ADR-133 audit-batch health counters from whichever `VerbRegistry`
790    /// owns the seam over this runtime's `EventStore` (typically
791    /// `VerbRegistry::audit_batch_metrics()`). `None` behaves identically to
792    /// [`Self::db_diagnostics`].
793    pub async fn db_diagnostics_with_audit_metrics(
794        &self,
795        runtime_audit_batch_metrics: Option<khive_db::diagnostics::RuntimeAuditBatchMetrics>,
796    ) -> RuntimeResult<khive_db::diagnostics::DbDiagnostics> {
797        let pool = self.core().backend.pool_arc();
798        // Match housekeeping's compiled legacy-record fallback (ADR-091
799        // Amendment 6), independent of checkpoint or local sweep overrides.
800        let legacy_sweep_interval = khive_db::SessionSweepConfig::default().interval;
801        let build_hash = crate::build_info::BUILD_INFO
802            .is_stamped()
803            .then_some(crate::build_info::BUILD_INFO.source_revision);
804        let build = khive_db::diagnostics::BuildIdentity::from_env(
805            crate::build_info::PACKAGE_VERSION,
806            build_hash,
807        );
808
809        let mut report = khive_db::diagnostics::collect_with_runtime_audit_metrics_interruptibly(
810            pool,
811            build,
812            legacy_sweep_interval,
813            crate::pack::audit_append_failure_count(),
814            runtime_audit_batch_metrics,
815        )
816        .await
817        .map_err(RuntimeError::from)?;
818        report.writer_contention.audit_obligation_append_failures =
819            Some(crate::pack::audit_obligation_append_failure_count());
820        report
821            .writer_contention
822            .audit_obligation_append_failures_unavailable_reason = None;
823        Ok(report)
824    }
825
826    // ---- Store accessors (token-scoped) ----
827
828    /// Get an EntityStore scoped to the token's namespace.
829    pub fn entities(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn EntityStore>> {
830        Ok(self
831            .backend
832            .entities_for_namespace(token.namespace().as_str())?)
833    }
834
835    /// Get a GraphStore scoped to the token's namespace.
836    pub fn graph(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn GraphStore>> {
837        Ok(self
838            .backend
839            .graph_for_namespace(token.namespace().as_str())?)
840    }
841
842    /// Get a NoteStore scoped to the token's namespace.
843    ///
844    /// Wrapped in `note_store_guard::PolicyEnforcingNoteStore`, which
845    /// refuses any insert/upsert of a `kind = "message"` note carrying
846    /// `quarantined` / `channel_kind` / `channel_slug` — the transport-owned
847    /// evidence `comm.health` trusts at face value — and refuses patching
848    /// those keys through the property-mutation seams on any note kind, so
849    /// the guard cannot be sidestepped by inserting a clean message note and
850    /// patching the evidence onto it afterward. Full-row writes also preserve
851    /// existing channel-health coordinates while allowing heartbeat metadata to
852    /// change. The trusted channel-ingest path does not go through this accessor; see
853    /// `Self::raw_notes` and [`Self::try_create_note_as_trusted_ingest`].
854    pub fn notes(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn NoteStore>> {
855        Ok(crate::note_store_guard::PolicyEnforcingNoteStore::wrap(
856            self.raw_notes(token)?,
857        ))
858    }
859
860    /// Get the unwrapped, policy-free NoteStore scoped to the token's namespace.
861    ///
862    /// Bypasses `note_store_guard::PolicyEnforcingNoteStore`. Callers
863    /// within this crate that have already enforced the reserved-transport-
864    /// property policy themselves (namely `try_create_note_impl`, which
865    /// applies it conditionally based on whether the caller presented a
866    /// [`crate::pack::ChannelIngestCapability`]) use this to reach storage
867    /// directly rather than run a redundant, less-informed check. Not exposed
868    /// outside this crate — every other caller must use [`Self::notes`].
869    pub(crate) fn raw_notes(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn NoteStore>> {
870        Ok(self
871            .backend
872            .notes_for_namespace(token.namespace().as_str())?)
873    }
874
875    /// Return the role-keyed attachment substrate on the canonical main backend.
876    ///
877    /// Attachment rows are the process-shared BlobStore's sole SQL liveness
878    /// authority. A runtime bound directly to a secondary pack backend must call
879    /// [`Self::core`] first; accepting a secondary mutation here would create a
880    /// reference that the main-database GC sweep cannot see or fence.
881    pub fn attachments(&self) -> RuntimeResult<Arc<dyn AttachmentStore>> {
882        if self.config.backend_id.as_str() != BackendId::MAIN {
883            return Err(RuntimeError::InvalidInput(format!(
884                "attachments are owned by the canonical main backend; runtime backend {:?} must route through KhiveRuntime::core()",
885                self.config.backend_id.as_str()
886            )));
887        }
888        Ok(self.backend.attachments()?)
889    }
890
891    /// Get an EventStore scoped to the token's namespace.
892    ///
893    /// When the events-daemon split (ADR-170) is configured, the store routes
894    /// by append class: the ADR-133 idempotent audit-batch lane — the
895    /// measured bulk of event write volume — persists to the events database
896    /// (forwarded over the events daemon socket in daemon deployments, or
897    /// opened directly in embedded/one-shot contexts), while plain appends
898    /// stay on this runtime's backend, keeping every raw-SQL consumer of the
899    /// legacy `events` table (schedule provenance, kg projection guards,
900    /// GraphQuery's substrate union) correct by construction. Reads merge
901    /// both stores. Unconfigured runtimes (tests, in-memory) keep the legacy
902    /// main-store behavior. Every returned store is decorated at this typed
903    /// accessor boundary so append callers cannot override the namespace or
904    /// actor resolved into the sealed authorization token.
905    pub fn events(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn EventStore>> {
906        Ok(crate::event_store_guard::AttributedEventStore::wrap(
907            self.raw_events_for_namespace(token.namespace().as_str())?,
908            token,
909        ))
910    }
911
912    /// Build the undecorated event store used only by the registry's trusted
913    /// audit composer, which stamps from each resolved `GateRequest` before
914    /// enqueueing. Pack/runtime call sites must use [`Self::events`] instead.
915    pub(crate) fn raw_events_for_namespace(
916        &self,
917        namespace: &str,
918    ) -> RuntimeResult<Arc<dyn EventStore>> {
919        let legacy = self.backend.events_for_namespace(namespace)?;
920        match &self.config.events_split {
921            None => Ok(legacy),
922            Some(split) => {
923                // Read-only is decided before the transport question: a
924                // read-only runtime must neither create nor schema-initialize
925                // an events database, and it must not forward writes to the
926                // events daemon either — a socket in the config describes the
927                // deployment, not this process's authority. Serve merged
928                // reads from a read-only open of the sidecar when it exists
929                // (the storage-layer read-only binding refuses any write that
930                // slips through), and the legacy store alone otherwise: no
931                // sidecar on disk means no lane rows exist, so minting the
932                // file just to read nothing from it would be a write in
933                // disguise.
934                if self.backend.is_read_only() {
935                    if !split.db_path.exists() {
936                        return Ok(legacy);
937                    }
938                    let lane = crate::events_split::direct_backend_with_max_readers(
939                        &split.db_path,
940                        true,
941                        Some(self.backend.pool().config().max_readers),
942                    )?
943                    .events_for_namespace(namespace)?;
944                    return Ok(Arc::new(crate::events_split::SplitEventStore::new(
945                        legacy, lane,
946                    )));
947                }
948                let lane: Arc<dyn EventStore> = match &split.socket_path {
949                    #[cfg(unix)]
950                    Some(socket) => {
951                        let client = crate::events_split::client_for(socket)?;
952                        Arc::new(crate::events_split::ForwardingEventStore::new(
953                            namespace, client,
954                        ))
955                    }
956                    #[cfg(not(unix))]
957                    Some(_socket) => {
958                        return Err(RuntimeError::InvalidInput(
959                            "events-daemon socket forwarding requires a Unix platform; \
960                             configure the events split in direct mode here"
961                                .to_string(),
962                        ));
963                    }
964                    None => crate::events_split::direct_backend_with_max_readers(
965                        &split.db_path,
966                        false,
967                        Some(self.backend.pool().config().max_readers),
968                    )?
969                    .events_for_namespace(namespace)?,
970                };
971                Ok(Arc::new(crate::events_split::SplitEventStore::new(
972                    legacy, lane,
973                )))
974            }
975        }
976    }
977
978    /// Get the raw SQL access capability (for ad-hoc queries).
979    pub fn sql(&self) -> Arc<dyn SqlAccess> {
980        self.backend.sql()
981    }
982
983    /// SQL access to the events-split sidecar database for read purposes,
984    /// when the split (ADR-170) is configured and the sidecar exists on
985    /// disk. `None` means every event row lives in the legacy `events`
986    /// table, so a raw-SQL consumer needs no second lookup. Consumers that
987    /// resolve an event by id or hex prefix against `self.sql()` must also
988    /// consult this store on a miss: the audit-batch lane's rows live only
989    /// in the sidecar.
990    ///
991    /// A writable runtime opens the sidecar through the ordinary writable
992    /// binding: WAL supports one-writer-many-readers, and the read-only
993    /// binding's frozen-snapshot guard refuses any sidecar with a live
994    /// writer's `-shm` beside it — exactly the live-deployment case these
995    /// reads exist for. A read-only runtime keeps the read-only open (it
996    /// must neither create nor schema-initialize a sidecar), which serves
997    /// genuinely frozen snapshots and refuses live ones, matching the
998    /// read-only arm of `events()`. Never creates a sidecar as a side
999    /// effect of a read.
1000    pub fn events_sidecar_sql_read_only(&self) -> RuntimeResult<Option<Arc<dyn SqlAccess>>> {
1001        match &self.config.events_split {
1002            None => Ok(None),
1003            Some(split) => {
1004                if !split.db_path.exists() {
1005                    return Ok(None);
1006                }
1007                let backend = if self.backend.is_read_only() {
1008                    crate::events_split::direct_backend_with_max_readers(
1009                        &split.db_path,
1010                        true,
1011                        Some(self.backend.pool().config().max_readers),
1012                    )?
1013                } else {
1014                    crate::events_split::direct_backend_with_max_readers(
1015                        &split.db_path,
1016                        false,
1017                        Some(self.backend.pool().config().max_readers),
1018                    )?
1019                };
1020                Ok(Some(backend.sql()))
1021            }
1022        }
1023    }
1024
1025    /// Get a VectorStore for the configured embedding model, scoped to the token's namespace.
1026    ///
1027    /// Returns `Unconfigured("embedding_model")` if no model is set.
1028    pub fn vectors(
1029        &self,
1030        token: &NamespaceToken,
1031    ) -> RuntimeResult<Arc<dyn khive_storage::VectorStore>> {
1032        let model = self.resolve_embedding_model(None)?;
1033        self.vectors_for_embedding_model(token, model)
1034    }
1035
1036    /// Get a VectorStore for a specific named embedding model, scoped to the token's namespace.
1037    ///
1038    /// Accepts both built-in lattice model names/aliases and custom provider names
1039    /// registered via [`register_embedder`](Self::register_embedder). Lattice names
1040    /// are routed through the enum-backed path; custom provider names use the
1041    /// provider's declared `dimensions()` directly so that the vector store key
1042    /// is consistent with how vectors were written during `remember`/`recall`.
1043    pub fn vectors_for_model(
1044        &self,
1045        token: &NamespaceToken,
1046        model_name: &str,
1047    ) -> RuntimeResult<Arc<dyn khive_storage::VectorStore>> {
1048        let (model_name, dims) = self.vector_model_metadata(model_name)?;
1049        Ok(self.backend.vectors_for_namespace(
1050            &sanitize_key(&model_name),
1051            &model_name,
1052            dims,
1053            token.namespace().as_str(),
1054        )?)
1055    }
1056
1057    /// Resolve the storage identity and declared dimensions together so guarded
1058    /// SQL publication agrees with VectorStore, including built-in aliases.
1059    pub(crate) fn vector_model_metadata(&self, model_name: &str) -> RuntimeResult<(String, usize)> {
1060        if request_excludes_embedder(model_name) {
1061            return Err(crate::RuntimeError::UnknownModel(model_name.to_string()));
1062        }
1063        let registry = self
1064            .embedder_registry
1065            .read()
1066            .map_err(|_| crate::RuntimeError::Internal("embedder registry lock poisoned".into()))?;
1067        if let Some(model) = parse_embedding_model_alias(model_name) {
1068            // Only proceed via the lattice path if this model is actually in the
1069            // registry; otherwise fall through to the custom-provider path.
1070            let key = model.to_string();
1071            if registry.contains(&key) {
1072                return Ok((key, model.dimensions()));
1073            }
1074        }
1075        registry
1076            .get_provider(model_name)
1077            .map(|provider| (model_name.to_owned(), provider.dimensions()))
1078            .ok_or_else(|| crate::RuntimeError::UnknownModel(model_name.to_string()))
1079    }
1080
1081    /// Get a namespace-scoped vector store for a pack-owned immutable identity.
1082    ///
1083    /// The table key is syntactically validated by [`NamedVectorIdentity`]. This
1084    /// accessor additionally verifies the table's actual sqlite-vec dimension
1085    /// declaration and every persisted `embedding_model` value before returning
1086    /// the store, so reusing one key for incompatible descriptor geometry or
1087    /// semantics fails before a caller can replace rows.
1088    pub async fn vectors_for_named_identity(
1089        &self,
1090        token: &NamespaceToken,
1091        identity: &NamedVectorIdentity,
1092    ) -> RuntimeResult<Arc<dyn VectorStore>> {
1093        let namespace = token.namespace().as_str();
1094        {
1095            let cached = self.named_vector_stores.read().map_err(|_| {
1096                RuntimeError::Internal("named vector store cache lock poisoned".into())
1097            })?;
1098            if let Some(entry) = cached.get(identity.model_key()) {
1099                check_cached_named_vector_identity(&entry.identity, identity)?;
1100                if let Some(store) = entry.by_namespace.get(namespace) {
1101                    return Ok(Arc::clone(store));
1102                }
1103            }
1104        }
1105        let store = self.backend.vectors_for_namespace(
1106            identity.model_key(),
1107            identity.model_name(),
1108            identity.dimensions(),
1109            namespace,
1110        )?;
1111
1112        let table = format!("vec_{}", identity.model_key());
1113        let mut reader = self.sql().reader().await?;
1114        let dimension_row = reader
1115            .query_row(SqlStatement {
1116                sql: "SELECT sql FROM sqlite_schema WHERE type = 'table' AND name = ?1".to_string(),
1117                params: vec![SqlValue::Text(table.clone())],
1118                label: Some("runtime_named_vector_dimension".to_string()),
1119            })
1120            .await?
1121            .ok_or_else(|| {
1122                RuntimeError::Internal(format!(
1123                    "named vector table {table} has no sqlite_schema declaration"
1124                ))
1125            })?;
1126        let table_ddl = match dimension_row.get("sql") {
1127            Some(SqlValue::Text(value)) => value,
1128            other => {
1129                return Err(RuntimeError::Internal(format!(
1130                    "named vector table {table} returned invalid schema metadata: {other:?}"
1131                )))
1132            }
1133        };
1134        let declared_dimensions = vector_dimensions_from_ddl(table_ddl).ok_or_else(|| {
1135            RuntimeError::Internal(format!(
1136                "named vector table {table} has no parseable embedding dimension"
1137            ))
1138        })?;
1139        if declared_dimensions != identity.dimensions() {
1140            return Err(RuntimeError::InvalidInput(format!(
1141                "named vector model_key {:?} is already bound to {declared_dimensions} dimensions, expected {}",
1142                identity.model_key(),
1143                identity.dimensions()
1144            )));
1145        }
1146
1147        let stored_models = reader
1148            .query_all(SqlStatement {
1149                sql: format!(
1150                    "SELECT DISTINCT embedding_model FROM {table} ORDER BY embedding_model LIMIT 2"
1151                ),
1152                params: vec![],
1153                label: Some("runtime_named_vector_model_identity".to_string()),
1154            })
1155            .await?;
1156        for row in stored_models {
1157            let stored = match row.get("embedding_model") {
1158                Some(SqlValue::Text(value)) => value,
1159                other => {
1160                    return Err(RuntimeError::Internal(format!(
1161                        "named vector table {table} returned invalid model identity metadata: {other:?}"
1162                    )))
1163                }
1164            };
1165            if stored != identity.model_name() {
1166                return Err(RuntimeError::InvalidInput(format!(
1167                    "named vector model_key {:?} already contains model {stored:?}, cannot bind it to {:?}",
1168                    identity.model_key(),
1169                    identity.model_name()
1170                )));
1171            }
1172        }
1173
1174        self.backend
1175            .register_embedding_model(
1176                identity.model_key(),
1177                identity.model_name(),
1178                identity.model_key(),
1179                identity.dimensions() as u32,
1180            )
1181            .map_err(|error| {
1182                if matches!(
1183                    &error,
1184                    khive_db::SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
1185                        if code.code == rusqlite::ErrorCode::ConstraintViolation
1186                ) {
1187                    RuntimeError::InvalidInput(format!(
1188                        "named vector model_key {:?} is already bound to a different active model identity",
1189                        identity.model_key()
1190                    ))
1191                } else {
1192                    RuntimeError::Sqlite(error)
1193                }
1194            })?;
1195
1196        let mut cached = self
1197            .named_vector_stores
1198            .write()
1199            .map_err(|_| RuntimeError::Internal("named vector store cache lock poisoned".into()))?;
1200        let entry = cached
1201            .entry(identity.model_key().to_owned())
1202            .or_insert_with(|| CachedNamedVectorStores {
1203                identity: identity.clone(),
1204                by_namespace: HashMap::new(),
1205            });
1206        check_cached_named_vector_identity(&entry.identity, identity)?;
1207        entry
1208            .by_namespace
1209            .insert(namespace.to_owned(), Arc::clone(&store));
1210        Ok(store)
1211    }
1212
1213    /// Output dimensions for a named embedding model, resolved from the
1214    /// embedder registry alone — no storage access. Mirrors
1215    /// [`vectors_for_model`](Self::vectors_for_model)'s resolution order:
1216    /// lattice aliases route through the enum when registered, otherwise the
1217    /// custom provider's declared `dimensions()`. `None` when no such model
1218    /// is registered.
1219    pub fn embedder_dimensions(&self, model_name: &str) -> Option<usize> {
1220        if request_excludes_embedder(model_name) {
1221            return None;
1222        }
1223        if let Some(model) = parse_embedding_model_alias(model_name) {
1224            let key = model.to_string();
1225            let in_registry = self
1226                .embedder_registry
1227                .read()
1228                .map(|reg| reg.contains(&key))
1229                .unwrap_or(false);
1230            if in_registry {
1231                return Some(model.dimensions());
1232            }
1233        }
1234        self.embedder_registry
1235            .read()
1236            .ok()?
1237            .get_provider(model_name)
1238            .map(|p| p.dimensions())
1239    }
1240
1241    fn vectors_for_embedding_model(
1242        &self,
1243        token: &NamespaceToken,
1244        model: EmbeddingModel,
1245    ) -> RuntimeResult<Arc<dyn khive_storage::VectorStore>> {
1246        Ok(self.backend.vectors_for_namespace(
1247            &vec_model_key(model),
1248            &model.to_string(),
1249            model.dimensions(),
1250            token.namespace().as_str(),
1251        )?)
1252    }
1253
1254    /// Get a TextSearch index for the entity corpus (single shared table).
1255    pub fn text(
1256        &self,
1257        token: &NamespaceToken,
1258    ) -> RuntimeResult<Arc<dyn khive_storage::TextSearch>> {
1259        let _ = token;
1260        Ok(self.backend.text("entities")?)
1261    }
1262
1263    /// Get a TextSearch index for the notes corpus (single shared table).
1264    pub fn text_for_notes(
1265        &self,
1266        token: &NamespaceToken,
1267    ) -> RuntimeResult<Arc<dyn khive_storage::TextSearch>> {
1268        let _ = token;
1269        Ok(self.backend.text("notes")?)
1270    }
1271
1272    /// Mint an authorization token for the given namespace.
1273    ///
1274    /// Consults the configured [`crate::Gate`] before minting. With the default
1275    /// `AllowAllGate` this always succeeds. When a real policy-backed gate is
1276    /// installed, this method enforces it and returns `PermissionDenied` on
1277    /// denial.
1278    ///
1279    /// The returned token's read visibility set defaults to `[ns]` — identical
1280    /// to the pre-visibility-set behaviour. Use [`Self::authorize_with_visibility`]
1281    /// to mint a token that can read additional namespaces.
1282    ///
1283    /// When `actor_id` is configured in `RuntimeConfig`, the token carries that
1284    /// actor label so that `comm.inbox` filters by `to_actor`. When
1285    /// unconfigured, the token carries `ActorRef::anonymous()` and inbox falls
1286    /// back to party-line behavior.
1287    pub fn authorize(&self, ns: Namespace) -> RuntimeResult<NamespaceToken> {
1288        let actor = crate::actor_identity::resolve_actor(self.config.actor_id.as_deref());
1289        let req = GateRequest::new(
1290            actor.clone(),
1291            ns.clone(),
1292            "authorize",
1293            serde_json::Value::Null,
1294        );
1295        match self.config.gate.check(&req) {
1296            Ok(ref decision) if decision.is_allow() => {
1297                if let khive_gate::GateDecision::Allow { ref obligations } = decision {
1298                    if !obligations.is_empty() {
1299                        tracing::debug!(
1300                            namespace = %ns.as_str(),
1301                            "authorize: obligations={:?}",
1302                            obligations
1303                        );
1304                    }
1305                }
1306                Ok(NamespaceToken::mint_authorized(ns, actor))
1307            }
1308            Ok(khive_gate::GateDecision::Deny { reason }) => {
1309                Err(crate::RuntimeError::permission_denied("authorize", reason))
1310            }
1311            Ok(_) => Err(crate::RuntimeError::permission_denied(
1312                "authorize",
1313                "gate denied",
1314            )),
1315            Err(e) => {
1316                tracing::warn!(
1317                    namespace = %ns.as_str(),
1318                    error = %crate::secret_gate::bounded_masked_log_text(&e.to_string()),
1319                    "authorize: gate check failed (fail-closed)"
1320                );
1321                Err(crate::RuntimeError::Internal(format!(
1322                    "gate error: {}",
1323                    e.wire_reason()
1324                )))
1325            }
1326        }
1327    }
1328
1329    /// Mint an authorization token with an explicit read-visibility set.
1330    ///
1331    /// `primary` is the **write namespace** — all records created via the
1332    /// returned token land there. `extra_visible` lists additional namespaces
1333    /// the token may read. The primary is always included in the visible set
1334    /// regardless of `extra_visible`.
1335    ///
1336    /// Usage (lambda:leo reading both leo and khive namespaces):
1337    /// ```rust,ignore
1338    /// let tok = rt.authorize_with_visibility(
1339    ///     Namespace::parse("lambda:leo").unwrap(),
1340    ///     vec![Namespace::parse("lambda:khive").unwrap()],
1341    /// )?;
1342    /// ```
1343    pub fn authorize_with_visibility(
1344        &self,
1345        primary: Namespace,
1346        extra_visible: Vec<Namespace>,
1347    ) -> RuntimeResult<NamespaceToken> {
1348        let actor = crate::actor_identity::resolve_actor(self.config.actor_id.as_deref());
1349        let req = GateRequest::new(
1350            actor.clone(),
1351            primary.clone(),
1352            "authorize",
1353            serde_json::Value::Null,
1354        );
1355        match self.config.gate.check(&req) {
1356            Ok(ref decision) if decision.is_allow() => {
1357                if let khive_gate::GateDecision::Allow { ref obligations } = decision {
1358                    if !obligations.is_empty() {
1359                        tracing::debug!(
1360                            namespace = %primary.as_str(),
1361                            "authorize_with_visibility: obligations={:?}",
1362                            obligations
1363                        );
1364                    }
1365                }
1366                // The primary check authorizes writes to `primary` only. Each
1367                // extra namespace grants read visibility, so each one takes
1368                // its own Read-classified gate check before it may enter the
1369                // minted set — a token must never carry visibility the gate
1370                // was not asked about. Any deny or gate error refuses the
1371                // whole mint, naming the offending namespace (fail-closed).
1372                for extra in &extra_visible {
1373                    let extra_req = GateRequest::new(
1374                        actor.clone(),
1375                        extra.clone(),
1376                        "authorize.visible",
1377                        serde_json::Value::Null,
1378                    );
1379                    match self.config.gate.check(&extra_req) {
1380                        Ok(ref extra_decision) if extra_decision.is_allow() => {}
1381                        Ok(khive_gate::GateDecision::Deny { reason }) => {
1382                            return Err(crate::RuntimeError::permission_denied(
1383                                "authorize",
1384                                format!(
1385                                    "visibility namespace {:?} denied: {reason}",
1386                                    extra.as_str()
1387                                ),
1388                            ));
1389                        }
1390                        Ok(_) => {
1391                            return Err(crate::RuntimeError::permission_denied(
1392                                "authorize",
1393                                format!("visibility namespace {:?} denied by gate", extra.as_str()),
1394                            ));
1395                        }
1396                        Err(e) => {
1397                            tracing::warn!(
1398                                namespace = %extra.as_str(),
1399                                error = %crate::secret_gate::bounded_masked_log_text(&e.to_string()),
1400                                "authorize_with_visibility: extra-namespace gate check failed (fail-closed)"
1401                            );
1402                            return Err(crate::RuntimeError::Internal(format!(
1403                                "gate error: {}",
1404                                e.wire_reason()
1405                            )));
1406                        }
1407                    }
1408                }
1409                Ok(NamespaceToken::mint_with_visibility(
1410                    primary,
1411                    extra_visible,
1412                    actor,
1413                ))
1414            }
1415            Ok(khive_gate::GateDecision::Deny { reason }) => {
1416                Err(crate::RuntimeError::permission_denied("authorize", reason))
1417            }
1418            Ok(_) => Err(crate::RuntimeError::permission_denied(
1419                "authorize",
1420                "gate denied",
1421            )),
1422            Err(e) => {
1423                tracing::warn!(
1424                    namespace = %primary.as_str(),
1425                    error = %crate::secret_gate::bounded_masked_log_text(&e.to_string()),
1426                    "authorize_with_visibility: gate check failed (fail-closed)"
1427                );
1428                Err(crate::RuntimeError::Internal(format!(
1429                    "gate error: {}",
1430                    e.wire_reason()
1431                )))
1432            }
1433        }
1434    }
1435
1436    /// Install the pack-aggregated edge endpoint rules.
1437    ///
1438    /// Called by the transport layer after the `VerbRegistry` is built so
1439    /// that runtime-layer edge validation can consult pack rules. Idempotent:
1440    /// later calls overwrite the previous rule set.
1441    pub fn install_edge_rules(&self, rules: Vec<EdgeEndpointRule>) {
1442        if let Ok(mut guard) = self.edge_rules.write() {
1443            *guard = rules;
1444        }
1445    }
1446
1447    /// Install an already-paired blob hydrator into this runtime.
1448    ///
1449    /// Reinstalling the exact same `Arc` is idempotent. A different pair is
1450    /// rejected: replacing it would split or reset the aggregate admission
1451    /// budget while requests may still hold leases.
1452    pub fn install_blob_hydrator(
1453        &self,
1454        hydrator: Arc<crate::blob::BlobHydrator>,
1455    ) -> RuntimeResult<()> {
1456        if hydrator.budget_bytes() != self.config.blob_hydration_bytes {
1457            return Err(RuntimeError::InvalidInput(format!(
1458                "blob hydrator budget {} does not match this runtime's resolved budget {}",
1459                hydrator.budget_bytes(),
1460                self.config.blob_hydration_bytes
1461            )));
1462        }
1463        // The mode gate must sit on THIS seam, not only on `install_blob_store`:
1464        // `BlobHydrator::new` is public, so without it a caller pairs a writable
1465        // store, installs it here, and `blob_store()` hands mutating pack paths
1466        // a writable store on a runtime whose declared mode is read-only.
1467        // Boot paths installing one hydrator across handles of MIXED modes —
1468        // where the blob pack's own backend mode, not each receiving
1469        // handle's, governs mutability — go through
1470        // [`Self::install_shared_blob_hydrator`] instead.
1471        if self.is_read_only() && !hydrator.enforces_read_only() {
1472            return Err(RuntimeError::InvalidInput(
1473                "this runtime is read-only: install the raw store with install_blob_store, \
1474                 which wraps it so every physical mutator refuses"
1475                    .to_string(),
1476            ));
1477        }
1478        self.install_blob_hydrator_slot(hydrator)
1479    }
1480
1481    /// Install a boot-shared hydrator whose mutability is governed by the
1482    /// blob runtime's own mode, not this handle's domain-store mode.
1483    ///
1484    /// ADR-160 D3 installs one hydrator `Arc` on every runtime handle a boot
1485    /// produces, and the documented multi-backend matrix includes a writable
1486    /// blob secondary beside a read-only main: there the shared hydrator is
1487    /// legitimately writable on a read-only domain handle. The mode decision
1488    /// must therefore already be encoded in the hydrator, and it must have
1489    /// been DERIVED, not declared: only hydrators built through
1490    /// [`crate::BlobHydrator::resolve_for_governing_backend`] — whose mode
1491    /// comes from the governing backend's own access mode — are accepted
1492    /// here. A hand-paired hydrator (`BlobHydrator::new` / `for_mode`) is
1493    /// refused so a safe downstream caller cannot use this seam to put a
1494    /// writable store on a read-only runtime; such callers use
1495    /// [`Self::install_blob_hydrator`], which holds hydrator mode against
1496    /// this runtime's own.
1497    ///
1498    /// This gate is a wrong-wiring guard, not an in-process sandbox: which
1499    /// backend governs is the boot host's topology assertion, and a caller
1500    /// who deliberately selects an unrelated writable backend as governing
1501    /// is outside what any runtime seam can enforce (see the trust-model
1502    /// note on [`crate::BlobHydrator::resolve_for_governing_backend`]).
1503    pub fn install_shared_blob_hydrator(
1504        &self,
1505        hydrator: Arc<crate::blob::BlobHydrator>,
1506    ) -> RuntimeResult<()> {
1507        if !hydrator.is_governed() {
1508            return Err(RuntimeError::InvalidInput(
1509                "the shared install seam accepts only hydrators whose mode was derived from a \
1510                 governing backend (BlobHydrator::resolve_for_governing_backend); use \
1511                 install_blob_hydrator for a hand-paired hydrator"
1512                    .to_string(),
1513            ));
1514        }
1515        if hydrator.budget_bytes() != self.config.blob_hydration_bytes {
1516            return Err(RuntimeError::InvalidInput(format!(
1517                "blob hydrator budget {} does not match this runtime's resolved budget {}",
1518                hydrator.budget_bytes(),
1519                self.config.blob_hydration_bytes
1520            )));
1521        }
1522        self.install_blob_hydrator_slot(hydrator)
1523    }
1524
1525    /// Whether `candidate` is the same pairing as the installed `current`:
1526    /// the exact `Arc`, or a distinct hydrator allocation over the same raw
1527    /// store with the same budget and mode. The latter arises when two
1528    /// concurrent first installs each construct a hydrator from one raw
1529    /// store — the `OnceLock` loser must read as idempotent, not as a
1530    /// conflicting install.
1531    fn is_same_blob_pairing(
1532        current: &Arc<crate::blob::BlobHydrator>,
1533        candidate: &Arc<crate::blob::BlobHydrator>,
1534    ) -> bool {
1535        Arc::ptr_eq(current, candidate)
1536            || (Arc::ptr_eq(&current.raw_store(), &candidate.raw_store())
1537                && current.budget_bytes() == candidate.budget_bytes()
1538                && current.enforces_read_only() == candidate.enforces_read_only())
1539    }
1540
1541    /// One-shot slot semantics shared by both install seams: first install
1542    /// wins, an equivalent pairing is idempotent, a different pairing is
1543    /// refused (replacing it would split or reset the aggregate admission
1544    /// budget while requests may still hold leases).
1545    fn install_blob_hydrator_slot(
1546        &self,
1547        hydrator: Arc<crate::blob::BlobHydrator>,
1548    ) -> RuntimeResult<()> {
1549        if let Some(current) = self.blob_hydrator.get() {
1550            return if Self::is_same_blob_pairing(current, &hydrator) {
1551                Ok(())
1552            } else {
1553                Err(RuntimeError::InvalidInput(
1554                    "a different blob hydrator is already installed".to_string(),
1555                ))
1556            };
1557        }
1558
1559        match self.blob_hydrator.set(hydrator) {
1560            Ok(()) => Ok(()),
1561            Err(candidate) => {
1562                let current = self.blob_hydrator.get().ok_or_else(|| {
1563                    RuntimeError::Internal(
1564                        "blob hydrator install raced without a visible winner".to_string(),
1565                    )
1566                })?;
1567                if Self::is_same_blob_pairing(current, &candidate) {
1568                    Ok(())
1569                } else {
1570                    Err(RuntimeError::InvalidInput(
1571                        "a different blob hydrator is already installed".to_string(),
1572                    ))
1573                }
1574            }
1575        }
1576    }
1577
1578    /// Pair and install a store using this runtime's resolved hydration budget.
1579    ///
1580    /// Boot paths that own multiple runtimes should instead construct one
1581    /// [`crate::BlobHydrator`] and call [`Self::install_blob_hydrator`] with
1582    /// the same `Arc` on every handle.
1583    pub fn install_blob_store(
1584        &self,
1585        store: Arc<dyn khive_storage::BlobStore>,
1586    ) -> RuntimeResult<()> {
1587        if let Some(current) = self.blob_hydrator.get() {
1588            // Reinstalling the same raw store is idempotent in BOTH modes:
1589            // on a read-only runtime the installed hydrator wraps the raw
1590            // store, so identity is checked against the raw handle the
1591            // hydrator remembers, not only the (possibly wrapped) paired one.
1592            let current_store = current.store();
1593            if Arc::ptr_eq(&current_store, &store) || Arc::ptr_eq(&current.raw_store(), &store) {
1594                return Ok(());
1595            }
1596        }
1597        // A read-only runtime holds its mode at this seam, not only during
1598        // boot resolution: an arbitrary store installed after launch is
1599        // wrapped so every physical mutator refuses while the bounded read
1600        // surface stays available. Without this, post-boot installation is a
1601        // writable bypass of the runtime's declared mode.
1602        let hydrator = if self.is_read_only() {
1603            crate::blob::BlobHydrator::new_read_only(store, self.config.blob_hydration_bytes)?
1604        } else {
1605            crate::blob::BlobHydrator::new(store, self.config.blob_hydration_bytes)?
1606        };
1607        self.install_blob_hydrator(Arc::new(hydrator))
1608    }
1609
1610    /// Return the installed shared blob hydrator, if boot configured one.
1611    pub fn blob_hydrator(&self) -> Option<Arc<crate::blob::BlobHydrator>> {
1612        self.blob_hydrator.get().cloned()
1613    }
1614
1615    /// Return the installed `BlobStore`, if the boot path resolved and
1616    /// installed one. `None` when no `[storage.blob]` selection was ever
1617    /// installed — e.g. a bare/test runtime constructed without going
1618    /// through the `khive-mcp` boot path.
1619    pub fn blob_store(&self) -> Option<Arc<dyn khive_storage::BlobStore>> {
1620        self.blob_hydrator.get().map(|hydrator| hydrator.store())
1621    }
1622
1623    /// Install the pack-aggregated valid entity and note kinds.
1624    ///
1625    /// Called by the transport layer after the `VerbRegistry` is built so that
1626    /// runtime-layer entity/note creation and import validate kind strings against
1627    /// the merged pack vocabulary. Idempotent: later calls overwrite previous sets.
1628    ///
1629    /// When no kinds are installed (empty lists), kind validation is skipped at
1630    /// the runtime layer. The pack handler layer remains the primary enforcement
1631    /// point; this provides defense-in-depth for direct Rust callers and import.
1632    pub fn install_kind_registry(&self, entity_kinds: Vec<String>, note_kinds: Vec<String>) {
1633        if let Ok(mut guard) = self.valid_entity_kinds.write() {
1634            *guard = entity_kinds;
1635        }
1636        if let Ok(mut guard) = self.valid_note_kinds.write() {
1637            *guard = note_kinds;
1638        }
1639    }
1640
1641    /// Install the pack-owned note kinds aggregated from the pack registry.
1642    ///
1643    /// Called by the transport after the `VerbRegistry` is built, same timing
1644    /// as [`install_kind_registry`](Self::install_kind_registry).
1645    pub fn install_pack_owned_note_kinds(&self, kinds: Vec<String>) {
1646        if let Ok(mut guard) = self.pack_owned_note_kinds.write() {
1647            *guard = kinds;
1648        }
1649    }
1650
1651    /// Whether `kind` is a note kind owned by a pack (see
1652    /// [`install_pack_owned_note_kinds`](Self::install_pack_owned_note_kinds)).
1653    ///
1654    /// Always `false` before the transport installs the list — a bare runtime
1655    /// has no packs, so no kind is pack-owned there.
1656    pub fn is_pack_owned_note_kind(&self, kind: &str) -> bool {
1657        self.pack_owned_note_kinds
1658            .read()
1659            .map(|g| g.iter().any(|k| k == kind))
1660            .unwrap_or(false)
1661    }
1662
1663    /// Validate that `kind` is a pack-registered entity kind.
1664    ///
1665    /// Returns `Ok(())` when no kinds are installed (bare runtime without packs).
1666    /// Returns `InvalidInput` when kinds are installed and `kind` is not among them.
1667    pub(crate) fn validate_entity_kind(&self, kind: &str) -> crate::RuntimeResult<()> {
1668        let guard = self.valid_entity_kinds.read().map_err(|_| {
1669            crate::RuntimeError::Internal("entity kind registry lock poisoned".into())
1670        })?;
1671        if guard.is_empty() {
1672            return Ok(());
1673        }
1674        if guard.iter().any(|k| k == kind) {
1675            Ok(())
1676        } else {
1677            Err(crate::RuntimeError::InvalidInput(format!(
1678                "unknown entity kind {kind:?}; valid: {}",
1679                guard.join(", ")
1680            )))
1681        }
1682    }
1683
1684    /// Validate that `kind` is a pack-registered note kind.
1685    ///
1686    /// Returns `Ok(())` when no kinds are installed (bare runtime without packs).
1687    /// Returns `InvalidInput` when kinds are installed and `kind` is not among them.
1688    pub(crate) fn validate_note_kind(&self, kind: &str) -> crate::RuntimeResult<()> {
1689        let guard = self.valid_note_kinds.read().map_err(|_| {
1690            crate::RuntimeError::Internal("note kind registry lock poisoned".into())
1691        })?;
1692        if guard.is_empty() {
1693            return Ok(());
1694        }
1695        if guard.iter().any(|k| k == kind) {
1696            Ok(())
1697        } else {
1698            Err(crate::RuntimeError::InvalidInput(format!(
1699                "unknown note kind {kind:?}; valid: {}",
1700                guard.join(", ")
1701            )))
1702        }
1703    }
1704
1705    /// Install a pack-supplied entity-type validator.
1706    ///
1707    /// Called by the `KgPack` during registration so that `create_many` can validate
1708    /// `entity_type` values at the runtime layer, closing the hole where direct Rust
1709    /// callers bypass the handler-layer `validate_entity_type` check.
1710    ///
1711    /// The callback receives `(kind, entity_type)` and returns the normalised type
1712    /// string, or `RuntimeError::InvalidInput` if the type is not registered for that
1713    /// kind. Passing `entity_type = None` must return `Ok(None)`.
1714    pub fn install_entity_type_validator(&self, f: EntityTypeValidatorFn) {
1715        if let Ok(mut guard) = self.entity_type_validator.write() {
1716            *guard = Some(f);
1717        }
1718    }
1719
1720    /// Validate and normalise `entity_type` through the pack-installed validator.
1721    ///
1722    /// Returns `Ok(entity_type)` when no validator is installed (bare runtime).
1723    /// Returns `InvalidInput` when a validator is installed and rejects the type.
1724    pub(crate) fn validate_entity_type_for_kind(
1725        &self,
1726        kind: &str,
1727        entity_type: Option<&str>,
1728    ) -> crate::RuntimeResult<Option<String>> {
1729        let guard = self.entity_type_validator.read().map_err(|_| {
1730            crate::RuntimeError::Internal("entity type validator lock poisoned".into())
1731        })?;
1732        match guard.as_ref() {
1733            None => Ok(entity_type.map(str::to_string)),
1734            Some(validate) => validate(kind, entity_type),
1735        }
1736    }
1737
1738    /// Install a pack-owned note-mutation hook.
1739    ///
1740    /// Overwrites any previously-installed hook, same single-slot semantics
1741    /// as [`install_entity_type_validator`](Self::install_entity_type_validator).
1742    /// In practice only one pack (`khive-pack-memory`) installs one today;
1743    /// if a second pack ever needs this, the slot should be widened to a
1744    /// `Vec` at that point rather than silently overwritten.
1745    pub fn install_note_mutation_hook(&self, f: NoteMutationHookFn) {
1746        if let Ok(mut guard) = self.note_mutation_hook.write() {
1747            *guard = Some(f);
1748        }
1749    }
1750
1751    /// Install the pack-aggregated entity-kind update hooks (issue #2943).
1752    ///
1753    /// Called by the transport after the `VerbRegistry` is built, same
1754    /// timing as [`install_kind_registry`](Self::install_kind_registry) —
1755    /// pass `registry.entity_kind_hooks()`. Idempotent: a later call
1756    /// replaces the set.
1757    pub fn install_entity_kind_hooks(&self, hooks: EntityKindHooks) {
1758        if let Ok(mut guard) = self.entity_kind_hooks.write() {
1759            *guard = hooks;
1760        }
1761    }
1762
1763    /// The installed `KindHook` for entity `kind`, if its owning pack
1764    /// registered one via [`install_entity_kind_hooks`](Self::install_entity_kind_hooks).
1765    ///
1766    /// `None` before the transport installs the aggregate (bare runtime) or
1767    /// when no pack registered a hook for this entity kind — the caller
1768    /// treats this the same as a hook whose `validate_entity_update`
1769    /// inherited the trait's `Ok(())` default.
1770    pub(crate) fn entity_kind_hook(&self, kind: &str) -> Option<Arc<dyn KindHook>> {
1771        self.entity_kind_hooks.read().ok().and_then(|guard| {
1772            guard
1773                .iter()
1774                .find(|(k, _)| k == kind)
1775                .map(|(_, hook)| hook.clone())
1776        })
1777    }
1778
1779    /// Install a pack-owned note-write validator.
1780    ///
1781    /// Called during pack registration (`PackRuntime::register_note_write_validator`)
1782    /// so that the covered note-write sites carrying caller-supplied
1783    /// `properties` derive the owning pack's identity properties from the
1784    /// authorization token, closing the gap where a direct Rust caller, the
1785    /// generic `create` verb, or the proposal-apply path (which dispatches no
1786    /// pack hooks) writes them unchecked. Single-slot semantics, same as
1787    /// [`install_note_mutation_hook`](Self::install_note_mutation_hook): a
1788    /// second installing pack overwrites the first, so a validator must return
1789    /// kinds it does not own unchanged.
1790    ///
1791    /// Covered sites — each calls `derive_note_write_properties`
1792    /// before the write: `create_note_inner` (`operations.rs`, the generic
1793    /// `create` verb funnel and every other public `create_note*` variant),
1794    /// `atomic_prepare::prepare_add_note` (the proposal-apply add-note path),
1795    /// and `atomic_message::create_notes_atomic_with_report` (the atomic
1796    /// multi-note writer).
1797    ///
1798    /// NOT covered by this validator: `try_create_note` (`operations.rs`).
1799    /// `try_create_note` is deliberately excluded — its only caller path is
1800    /// `comm.ingest`, where `properties.from_actor` is the external transport
1801    /// sender named by the `from` parameter, not the authenticated caller,
1802    /// and where transport-owned quarantine/channel properties are
1803    /// legitimately established. Running the generic validator there would
1804    /// stamp every inbound message as the ingesting daemon and reject the
1805    /// evidence the trusted ingest handler just derived. `try_create_note`
1806    /// instead runs its own narrower reserved-transport-property check
1807    /// inline (`operations.rs`'s `try_create_note_impl`), which allows the
1808    /// three `message`-kind transport properties only when called through
1809    /// [`Self::try_create_note_as_trusted_ingest`] with a
1810    /// [`crate::pack::ChannelIngestCapability`].
1811    ///
1812    /// The `NoteStore` returned by [`notes`](Self::notes) is covered by a
1813    /// different, narrower mechanism: it is wrapped in
1814    /// `note_store_guard::PolicyEnforcingNoteStore`, which refuses
1815    /// `upsert_note` / `upsert_notes` / `try_insert_note` /
1816    /// `replace_note_if_unchanged` calls that would write a `kind = "message"`
1817    /// note carrying `quarantined` / `channel_kind` / `channel_slug`, and
1818    /// refuses `set_note_property` / `try_patch_note_property` /
1819    /// `patch_note_property_atomic` / `update_note_properties` calls that
1820    /// would patch any of those keys onto any note — unconditionally, since
1821    /// that public accessor has no way to see a trust decision. `try_create_note_impl` itself reaches storage through
1822    /// `Self::raw_notes`, the unwrapped accessor, so its own inline check
1823    /// (which can legitimately allow those properties for trusted ingest)
1824    /// is not double-enforced or contradicted by the wrapper.
1825    /// Register a pack-defined custom fusion strategy under `name` (ADR-012).
1826    ///
1827    /// Unlike `install_entity_type_validator`/`install_note_mutation_hook`,
1828    /// this slot is keyed rather than single-occupancy: multiple packs each
1829    /// register their own named strategy, and a second registration under an
1830    /// already-used `name` replaces the first. Looked up by
1831    /// `FusionStrategy::Custom { name, .. }` at the hybrid-search dispatch
1832    /// boundary in `crate::fusion`; an unregistered name fails closed with
1833    /// `RuntimeError::UnknownFusionStrategy` rather than silently falling
1834    /// back to RRF.
1835    pub fn register_fusion_strategy(
1836        &self,
1837        name: impl Into<String>,
1838        executor: Arc<dyn crate::fusion::FusionExecutor>,
1839    ) {
1840        if let Ok(mut guard) = self.fusion_executors.write() {
1841            guard.insert(name.into(), executor);
1842        }
1843    }
1844
1845    /// Resolve a registered custom fusion executor by name.
1846    ///
1847    /// Returns `RuntimeError::UnknownFusionStrategy` when no pack has
1848    /// registered `name` — callers must invoke this before any
1849    /// empty-input/zero-limit short circuit so a misconfigured name errors
1850    /// on every call, including zero-result ones.
1851    pub(crate) fn fusion_executor(
1852        &self,
1853        name: &str,
1854    ) -> RuntimeResult<Arc<dyn crate::fusion::FusionExecutor>> {
1855        let guard = self
1856            .fusion_executors
1857            .read()
1858            .map_err(|_| RuntimeError::Internal("fusion executor registry lock poisoned".into()))?;
1859        guard
1860            .get(name)
1861            .cloned()
1862            .ok_or_else(|| RuntimeError::UnknownFusionStrategy(name.to_string()))
1863    }
1864
1865    pub fn install_note_write_validator(&self, f: NoteWriteValidatorFn) {
1866        if let Ok(mut guard) = self.note_write_validator.write() {
1867            *guard = Some(f);
1868        }
1869    }
1870
1871    /// Whether a note-write validator is installed on this runtime.
1872    ///
1873    /// Exists so a transport's own tests can assert, per boot path, that the
1874    /// documented startup sequence actually filled the slot. A missing install
1875    /// fails open and silently — an empty slot passes caller-supplied
1876    /// properties straight through, which no write site can distinguish from a
1877    /// validator that approved them — so occupancy is asserted, never assumed.
1878    pub fn has_note_write_validator(&self) -> bool {
1879        self.note_write_validator
1880            .read()
1881            .map(|g| g.is_some())
1882            .unwrap_or(false)
1883    }
1884
1885    /// Run caller-supplied note `properties` through the installed note-write
1886    /// validator, returning the properties to store.
1887    ///
1888    /// Returns them unchanged when no validator is installed (bare runtime).
1889    pub(crate) fn derive_note_write_properties(
1890        &self,
1891        kind: &str,
1892        token: &NamespaceToken,
1893        properties: Option<serde_json::Value>,
1894    ) -> RuntimeResult<Option<serde_json::Value>> {
1895        let validator = self
1896            .note_write_validator
1897            .read()
1898            .map_err(|_| RuntimeError::Internal("note write validator lock poisoned".into()))?
1899            .clone();
1900        match validator {
1901            None => Ok(properties),
1902            Some(validate) => validate(kind, &token.actor().id, properties),
1903        }
1904    }
1905
1906    /// Invoke the pack-installed note-mutation hook, if any.
1907    ///
1908    /// `kind` is the note's `kind` string (e.g. `"memory"`); `id` is the
1909    /// note's UUID. No-op when no hook is installed (bare runtime, or no
1910    /// pack cares). Errors inside the hook are the hook's own concern to
1911    /// handle/log — this call site cannot propagate a failure without
1912    /// changing `update_note`/`delete_note`'s already-committed success
1913    /// return value.
1914    pub(crate) async fn fire_note_mutation_hook(&self, kind: &str, id: uuid::Uuid) {
1915        let hook = self
1916            .note_mutation_hook
1917            .read()
1918            .ok()
1919            .and_then(|guard| guard.clone());
1920        if let Some(hook) = hook {
1921            hook(kind.to_string(), id).await;
1922        }
1923    }
1924
1925    /// Snapshot of currently-installed pack edge rules.
1926    ///
1927    /// This is the same composed rule set `validate_edge_relation_endpoints`
1928    /// consults via `pack_rule_allows` when accepting/rejecting an edge. Public
1929    /// so pack-layer error-hint code (e.g. `khive-pack-kg`'s
1930    /// `valid_relations_for_entity_pair`) can derive hints from the exact
1931    /// source the validator uses, rather than maintaining a separate
1932    /// hand-authored table that can drift out of sync.
1933    pub fn pack_edge_rules(&self) -> Vec<EdgeEndpointRule> {
1934        self.edge_rules
1935            .read()
1936            .map(|g| g.clone())
1937            .unwrap_or_default()
1938    }
1939
1940    /// Borrow the installed pack edge rules for a synchronous calculation.
1941    pub(crate) fn with_pack_edge_rules<T>(&self, f: impl FnOnce(&[EdgeEndpointRule]) -> T) -> T {
1942        match self.edge_rules.read() {
1943            Ok(rules) => f(&rules),
1944            Err(_) => f(&[]),
1945        }
1946    }
1947
1948    /// Return the name of the default embedding model (empty string if none configured).
1949    pub fn default_embedder_name(&self) -> &str {
1950        self.default_embedder_name.as_ref()
1951    }
1952
1953    /// Resolve a model name (or `None` for the default) to an `EmbeddingModel`.
1954    ///
1955    /// Returns `UnknownModel` if the name is not in the registry, or
1956    /// `Unconfigured` if `None` is passed and no default model is set.
1957    pub fn resolve_embedding_model(&self, name: Option<&str>) -> RuntimeResult<EmbeddingModel> {
1958        let model = match name {
1959            Some(raw) => parse_embedding_model_alias(raw)
1960                .ok_or_else(|| crate::RuntimeError::UnknownModel(raw.to_string()))?,
1961            None => self
1962                .config
1963                .embedding_model
1964                .ok_or_else(|| crate::RuntimeError::Unconfigured("embedding_model".into()))?,
1965        };
1966        let key = model.to_string();
1967        if request_excludes_embedder(&key) {
1968            return Err(crate::RuntimeError::UnknownModel(
1969                name.unwrap_or_else(|| self.default_embedder_name())
1970                    .to_string(),
1971            ));
1972        }
1973        let contains = self
1974            .embedder_registry
1975            .read()
1976            .map(|reg| reg.contains(&key))
1977            .unwrap_or(false);
1978        if contains {
1979            Ok(model)
1980        } else {
1981            Err(crate::RuntimeError::UnknownModel(
1982                name.unwrap_or_else(|| self.default_embedder_name())
1983                    .to_string(),
1984            ))
1985        }
1986    }
1987
1988    /// Names of all registered embedding models in this runtime.
1989    ///
1990    /// Includes both built-in lattice models and any custom embedders
1991    /// registered by packs via [`register_embedder`](Self::register_embedder).
1992    /// Useful for operations that must touch every model's storage (e.g.,
1993    /// scoped vector deletion on note delete). The default model is included.
1994    pub fn registered_embedding_model_names(&self) -> Vec<String> {
1995        self.embedder_registry
1996            .read()
1997            .map(|reg| {
1998                reg.names()
1999                    .into_iter()
2000                    .filter(|name| !request_excludes_embedder(name))
2001                    .collect()
2002            })
2003            .unwrap_or_default()
2004    }
2005
2006    /// Get the lazily-initialized embedding service for the named model.
2007    ///
2008    /// Accepts both built-in lattice model names (e.g. `"all-minilm-l6-v2"`,
2009    /// `"paraphrase"`) and custom provider names registered via
2010    /// [`register_embedder`](Self::register_embedder).
2011    ///
2012    /// For lattice model names, aliases (e.g. `"paraphrase"`) are resolved to
2013    /// their canonical key before looking up the registry. For custom providers
2014    /// the name must match exactly as supplied during registration.
2015    ///
2016    /// First call for any name loads the underlying service (cold start cost);
2017    /// subsequent calls are cheap (registry caches the `Arc`).
2018    pub async fn embedder(&self, name: &str) -> RuntimeResult<Arc<dyn EmbeddingService>> {
2019        self.embedder_inner(name, None).await
2020    }
2021
2022    pub(crate) async fn embedder_with_token(
2023        &self,
2024        token: &NamespaceToken,
2025        name: &str,
2026    ) -> RuntimeResult<Arc<dyn EmbeddingService>> {
2027        self.embedder_inner(name, Some(token)).await
2028    }
2029
2030    async fn embedder_inner(
2031        &self,
2032        name: &str,
2033        token: Option<&NamespaceToken>,
2034    ) -> RuntimeResult<Arc<dyn EmbeddingService>> {
2035        // Fall back to the literal name (not the alias table) so custom
2036        // providers registered with non-lattice names stay reachable.
2037        let canonical_key = match parse_embedding_model_alias(name) {
2038            Some(model) => model.to_string(),
2039            None => name.to_owned(),
2040        };
2041        if request_excludes_embedder(&canonical_key) {
2042            return Err(crate::RuntimeError::UnknownModel(name.to_string()));
2043        }
2044        // Clone the entry so we don't hold the RwLockGuard across the
2045        // async OnceCell initialisation (Send bound).
2046        let entry = {
2047            let registry = self.embedder_registry.read().map_err(|_| {
2048                crate::RuntimeError::Internal("embedder registry lock poisoned".into())
2049            })?;
2050            registry
2051                .get_entry(&canonical_key)
2052                .ok_or_else(|| crate::RuntimeError::UnknownModel(name.to_string()))?
2053        };
2054        let (service, init_duration_us) = entry.resolve().await?;
2055        if let Some(duration_us) = init_duration_us {
2056            if let Some(token) = token {
2057                self.emit_embedder_initialized(token, &canonical_key, duration_us)
2058                    .await;
2059            } else if let Ok(token) = self.authorize(self.config.default_namespace.clone()) {
2060                self.emit_embedder_initialized(&token, &canonical_key, duration_us)
2061                    .await;
2062            }
2063        }
2064        Ok(service)
2065    }
2066
2067    async fn emit_embedder_initialized(
2068        &self,
2069        token: &NamespaceToken,
2070        model_name: &str,
2071        duration_us: i64,
2072    ) {
2073        // Lazy embedder construction can happen during daemon warm or an
2074        // assertive request. A snapshot has no durable audit sink, so do not
2075        // resolve an EventStore merely to attempt a known-rejected append.
2076        if self.is_read_only() {
2077            return;
2078        }
2079        let Ok(store) = self.events(token) else {
2080            return;
2081        };
2082        let event = Event::new(
2083            token.namespace().as_str(),
2084            "embedder.init",
2085            EventKind::EmbedderInitialized,
2086            SubstrateKind::Event,
2087            format!("{}:{}", token.actor().kind, token.actor().id),
2088        )
2089        .with_payload(serde_json::json!({
2090            "model_name": model_name,
2091            "duration_us": duration_us,
2092        }))
2093        .with_duration_us(duration_us);
2094        if let Err(err) = store.append_event(event).await {
2095            tracing::warn!(error = %err, model_name, "embedder initialization event append failed");
2096        }
2097    }
2098
2099    /// Register a custom embedding provider with this runtime.
2100    ///
2101    /// The provider is added to the shared [`EmbedderRegistry`] so all clones
2102    /// of this runtime see the new provider immediately. If a provider with the
2103    /// same name already exists it is replaced (last-writer wins — see
2104    /// [`crate::EmbedderRegistry::register`] for the rationale).
2105    ///
2106    /// Packs should call this from [`crate::PackRuntime::register_embedders`] (the
2107    /// hook is invoked by the transport during pack initialisation, before the
2108    /// first verb dispatch).
2109    ///
2110    /// [`EmbedderRegistry`]: crate::embedder_registry::EmbedderRegistry
2111    pub fn register_embedder(
2112        &self,
2113        provider: impl crate::embedder_registry::EmbedderProvider + 'static,
2114    ) {
2115        if let Ok(mut registry) = self.embedder_registry.write() {
2116            registry.register(provider);
2117        } else {
2118            tracing::warn!(
2119                "embedder registry lock poisoned — embedder {} not registered",
2120                std::any::type_name::<dyn crate::embedder_registry::EmbedderProvider>()
2121            );
2122        }
2123    }
2124
2125    /// List registered embedding models via `SqlAccess`, routing through the
2126    /// existing connection pool rather than opening a fresh `Connection` per call.
2127    ///
2128    /// Optionally filter by `engine_name`. Returns an empty vec when the
2129    /// `_embedding_models` table does not yet exist (e.g. no migrations have run
2130    /// or no models have been registered). All other SQL errors are propagated.
2131    pub async fn list_embedding_models(
2132        &self,
2133        engine_filter: Option<&str>,
2134    ) -> RuntimeResult<Vec<khive_db::EmbeddingModelRegistryRecord>> {
2135        use khive_storage::{SqlStatement, SqlValue};
2136
2137        let (sql_text, params) = if let Some(engine) = engine_filter {
2138            (
2139                "SELECT engine_name, model_id, key_version, dim, status, \
2140                 activated_at, superseded_at \
2141                 FROM _embedding_models WHERE engine_name = ?1 \
2142                 ORDER BY engine_name, activated_at IS NULL, activated_at"
2143                    .to_string(),
2144                vec![SqlValue::Text(engine.to_string())],
2145            )
2146        } else {
2147            (
2148                "SELECT engine_name, model_id, key_version, dim, status, \
2149                 activated_at, superseded_at \
2150                 FROM _embedding_models \
2151                 ORDER BY engine_name, activated_at IS NULL, activated_at"
2152                    .to_string(),
2153                vec![],
2154            )
2155        };
2156
2157        let stmt = SqlStatement {
2158            sql: sql_text,
2159            params,
2160            label: Some("list_embedding_models".into()),
2161        };
2162
2163        let mut reader = self
2164            .sql()
2165            .reader()
2166            .await
2167            .map_err(crate::RuntimeError::Storage)?;
2168
2169        let rows = match reader.query_all(stmt).await {
2170            Ok(rows) => rows,
2171            Err(e) if e.to_string().contains("no such table: _embedding_models") => {
2172                return Ok(Vec::new())
2173            }
2174            Err(e) => return Err(crate::RuntimeError::Storage(e)),
2175        };
2176
2177        let mut records = Vec::with_capacity(rows.len());
2178        for row in rows {
2179            macro_rules! required_text {
2180                ($col:expr) => {
2181                    match row.get($col) {
2182                        Some(SqlValue::Text(s)) => s.clone(),
2183                        other => {
2184                            tracing::warn!(column = $col, value = ?other, "skipping registry row: unexpected type");
2185                            continue;
2186                        }
2187                    }
2188                };
2189            }
2190            let engine_name = required_text!("engine_name");
2191            let model_id = required_text!("model_id");
2192            let key_version = required_text!("key_version");
2193            let dimensions = match row.get("dim") {
2194                Some(SqlValue::Integer(n)) => match u32::try_from(*n) {
2195                    Ok(d) => d,
2196                    Err(_) => {
2197                        tracing::warn!(dim = n, "skipping registry row: dim out of u32 range");
2198                        continue;
2199                    }
2200                },
2201                other => {
2202                    tracing::warn!(column = "dim", value = ?other, "skipping registry row: unexpected type");
2203                    continue;
2204                }
2205            };
2206            let status = required_text!("status");
2207            let activated_at = match row.get("activated_at") {
2208                Some(SqlValue::Integer(n)) => Some(*n),
2209                _ => None,
2210            };
2211            let superseded_at = match row.get("superseded_at") {
2212                Some(SqlValue::Integer(n)) => Some(*n),
2213                _ => None,
2214            };
2215            records.push(khive_db::EmbeddingModelRegistryRecord {
2216                engine_name,
2217                model_id,
2218                key_version,
2219                dimensions,
2220                status,
2221                activated_at,
2222                superseded_at,
2223            });
2224        }
2225
2226        Ok(records)
2227    }
2228}
2229
2230fn vector_dimensions_from_ddl(ddl: &str) -> Option<usize> {
2231    let lower = ddl.to_ascii_lowercase();
2232    let suffix = lower.split_once("embedding float[")?.1;
2233    let dimension = suffix.split_once(']')?.0;
2234    if dimension.is_empty() || !dimension.bytes().all(|byte| byte.is_ascii_digit()) {
2235        return None;
2236    }
2237    dimension.parse().ok()
2238}
2239
2240// INLINE TEST JUSTIFICATION: tests here cover KhiveRuntime construction helpers
2241// (in-memory backend wiring, NamespaceToken::for_namespace) that are
2242// pub(crate)-only and cannot be called from the integration test crate.
2243#[cfg(test)]
2244mod tests {
2245    use super::*;
2246    use khive_gate::GateRef;
2247    use serial_test::serial;
2248
2249    #[cfg(target_os = "macos")]
2250    #[test]
2251    fn in_process_runtime_tests_have_4096_open_file_slots() {
2252        let _runtime = KhiveRuntime::memory().expect("test runtime");
2253        let mut limits = libc::rlimit {
2254            rlim_cur: 0,
2255            rlim_max: 0,
2256        };
2257        // SAFETY: `limits` is a writable local value.
2258        assert_eq!(
2259            unsafe { libc::getrlimit(libc::RLIMIT_NOFILE, &mut limits) },
2260            0
2261        );
2262        assert!(
2263            limits.rlim_cur >= IN_PROCESS_TEST_NOFILE_LIMIT,
2264            "a parallel runtime suite needs at least 4096 open-file slots"
2265        );
2266    }
2267
2268    fn test_blob_hydrator() -> (tempfile::TempDir, Arc<crate::BlobHydrator>) {
2269        let root = tempfile::tempdir().expect("blob root");
2270        let store = Arc::new(
2271            khive_db::stores::blob::FsBlobStore::new(root.path().to_path_buf(), 0)
2272                .expect("fs blob store"),
2273        );
2274        let hydrator = Arc::new(
2275            crate::BlobHydrator::new(store, crate::DEFAULT_BLOB_HYDRATION_BYTES)
2276                .expect("blob hydrator"),
2277        );
2278        (root, hydrator)
2279    }
2280
2281    #[test]
2282    fn memory_runtime_creates_successfully() {
2283        let rt = KhiveRuntime::memory().expect("memory runtime should create");
2284        assert!(rt.config().db_path.is_none());
2285    }
2286
2287    #[test]
2288    fn installed_blob_hydrator_is_shared_by_clone_and_core_handles() {
2289        let main_backend = Arc::new(StorageBackend::memory().expect("main backend"));
2290        let pack_backend = Arc::new(StorageBackend::memory().expect("pack backend"));
2291        let mut config = RuntimeConfig::no_embeddings();
2292        config.backend_id = BackendId::parse("assets").expect("valid backend id");
2293        let runtime = KhiveRuntime::from_backend(pack_backend, config)
2294            .with_core_backend(Arc::clone(&main_backend));
2295        let (_root, hydrator) = test_blob_hydrator();
2296
2297        runtime
2298            .install_blob_hydrator(Arc::clone(&hydrator))
2299            .expect("first install");
2300
2301        for handle in [runtime.clone(), runtime.core()] {
2302            let installed = handle.blob_hydrator().expect("installed hydrator");
2303            assert!(Arc::ptr_eq(&installed, &hydrator));
2304        }
2305    }
2306
2307    #[test]
2308    fn blob_hydrator_install_is_idempotent_but_rejects_replacement() {
2309        let runtime = KhiveRuntime::memory().expect("runtime");
2310        let (_first_root, first) = test_blob_hydrator();
2311        let (_second_root, second) = test_blob_hydrator();
2312
2313        runtime
2314            .install_blob_hydrator(Arc::clone(&first))
2315            .expect("first install");
2316        runtime
2317            .install_blob_hydrator(Arc::clone(&first))
2318            .expect("same Arc reinstall is idempotent");
2319
2320        let error = runtime
2321            .install_blob_hydrator(second)
2322            .expect_err("a different hydrator must not replace the installed pair");
2323        assert!(error.to_string().contains("already installed"));
2324        assert!(Arc::ptr_eq(
2325            &runtime.blob_hydrator().expect("original remains"),
2326            &first
2327        ));
2328    }
2329
2330    #[test]
2331    fn blob_hydrator_install_rejects_a_budget_that_disagrees_with_runtime_identity() {
2332        let runtime = KhiveRuntime::memory().expect("runtime");
2333        let root = tempfile::tempdir().expect("blob root");
2334        let store = Arc::new(
2335            khive_db::stores::blob::FsBlobStore::new(root.path().to_path_buf(), 0)
2336                .expect("fs blob store"),
2337        );
2338        let mismatched = Arc::new(
2339            crate::BlobHydrator::new(store, khive_storage::MAX_BLOB_WHOLE_BYTES)
2340                .expect("minimum blob budget"),
2341        );
2342
2343        let error = runtime
2344            .install_blob_hydrator(mismatched)
2345            .expect_err("live admission must match the construction-baked config identity");
2346        assert!(matches!(error, RuntimeError::InvalidInput(_)));
2347        assert!(runtime.blob_hydrator().is_none());
2348    }
2349
2350    #[test]
2351    fn fresh_tail_policy_is_instance_scoped_and_clone_stable() {
2352        let enabled = KhiveRuntime::memory()
2353            .expect("enabled memory runtime")
2354            .with_ann_fresh_tail_enabled(true);
2355        let disabled = KhiveRuntime::memory()
2356            .expect("disabled memory runtime")
2357            .with_ann_fresh_tail_enabled(false);
2358
2359        assert!(enabled.ann_fresh_tail_enabled());
2360        assert!(enabled.clone().ann_fresh_tail_enabled());
2361        assert!(!disabled.ann_fresh_tail_enabled());
2362        assert!(!disabled.clone().ann_fresh_tail_enabled());
2363    }
2364
2365    #[tokio::test]
2366    async fn runtime_db_diagnostics_supplies_both_contention_counter_sources() {
2367        let rt = KhiveRuntime::memory().expect("memory runtime should create");
2368
2369        let report = rt.db_diagnostics().await.expect("diagnostics succeed");
2370
2371        assert!(
2372            report.writer_contention.writer_acquisitions >= 1,
2373            "runtime construction runs migrations through the finite-wait pooled writer"
2374        );
2375        assert_eq!(
2376            report.writer_contention.writer_acquisitions,
2377            report
2378                .writer_contention
2379                .pooled_writer_acquisitions
2380                .saturating_add(report.writer_contention.standalone_writer_acquisitions)
2381                .saturating_add(report.writer_contention.writer_task_acquisitions),
2382            "the public aggregate must equal the class-specific snapshot"
2383        );
2384        assert!(
2385            report.writer_contention.audit_append_failures.is_some(),
2386            "the runtime path must supply its process-wide swallowed-audit counter"
2387        );
2388        assert!(report
2389            .writer_contention
2390            .audit_obligation_append_failures
2391            .is_some());
2392        assert!(report
2393            .writer_contention
2394            .audit_obligation_append_failures_unavailable_reason
2395            .is_none());
2396        assert!(report
2397            .writer_contention
2398            .audit_append_failures_unavailable_reason
2399            .is_none());
2400    }
2401
2402    #[test]
2403    fn backend_data_dir_returns_none_for_memory_backend() {
2404        let rt = KhiveRuntime::memory().expect("memory runtime");
2405        assert!(rt.backend_data_dir().is_none());
2406    }
2407
2408    #[test]
2409    fn backend_data_dir_returns_parent_dir_for_file_backend() {
2410        let dir = tempfile::tempdir().unwrap();
2411        let path = dir.path().join("test.db");
2412        let config = RuntimeConfig {
2413            web: Default::default(),
2414            telemetry: Default::default(),
2415            mounts: Vec::new(),
2416            brain: Default::default(),
2417            git_write: Default::default(),
2418            display_timezone: chrono_tz::Tz::UTC,
2419            events_split: None,
2420            db_path: Some(path),
2421            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2422            default_namespace: Namespace::local(),
2423            embedding_model: None,
2424            additional_embedding_models: vec![],
2425            gate: Arc::new(AllowAllGate),
2426            packs: vec!["kg".to_string()],
2427            backend_id: BackendId::main(),
2428            brain_profile: None,
2429            visible_namespaces: vec![],
2430            allowed_outbound_namespaces: vec![],
2431            actor_id: None,
2432            exec: Default::default(),
2433        };
2434        let rt = KhiveRuntime::new_for_test(config).expect("file runtime");
2435        let data_dir = rt
2436            .backend_data_dir()
2437            .expect("file backend must return Some");
2438        assert_eq!(data_dir, dir.path());
2439    }
2440
2441    /// A sidecar-only event must resolve through the public hex-prefix path:
2442    /// the main-store scan cannot see lane rows, so `resolve_prefix_inner`
2443    /// carries a sidecar arm. The pre-insert assert is the control — the
2444    /// prefix misses until the lane row exists, so a pass cannot come from
2445    /// the legacy scan.
2446    #[tokio::test]
2447    async fn resolve_prefix_finds_sidecar_only_event() {
2448        let dir = tempfile::tempdir().unwrap();
2449        let _registry_guard = crate::events_split::TestRegistryGuard::new(dir.path());
2450        let sidecar_path = dir.path().join("main.db.events.db");
2451        let config = RuntimeConfig {
2452            web: Default::default(),
2453            telemetry: Default::default(),
2454            mounts: Vec::new(),
2455            brain: Default::default(),
2456            git_write: Default::default(),
2457            display_timezone: chrono_tz::Tz::UTC,
2458            events_split: Some(crate::events_split::EventsSplitConfig {
2459                db_path: sidecar_path.clone(),
2460                socket_path: None,
2461            }),
2462            db_path: Some(dir.path().join("main.db")),
2463            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2464            default_namespace: Namespace::local(),
2465            embedding_model: None,
2466            additional_embedding_models: vec![],
2467            gate: Arc::new(AllowAllGate),
2468            packs: vec!["kg".to_string()],
2469            backend_id: BackendId::main(),
2470            brain_profile: None,
2471            visible_namespaces: vec![],
2472            allowed_outbound_namespaces: vec![],
2473            actor_id: None,
2474            exec: Default::default(),
2475        };
2476        let rt = KhiveRuntime::new_for_test(config).expect("file runtime");
2477
2478        let event = khive_storage::Event::new(
2479            "local",
2480            "memory.recall",
2481            khive_types::EventKind::RecallExecuted,
2482            khive_types::SubstrateKind::Note,
2483            "agent:test",
2484        );
2485        let event_id = event.id;
2486        let prefix = event_id.to_string()[..8].to_string();
2487
2488        assert_eq!(
2489            rt.resolve_prefix_unfiltered(&prefix)
2490                .await
2491                .expect("pre-insert resolve"),
2492            None,
2493            "control: prefix must miss before the lane row exists"
2494        );
2495
2496        let lane = crate::events_split::direct_backend_for(&sidecar_path)
2497            .expect("lane backend")
2498            .events_for_namespace("local")
2499            .expect("lane store");
2500        lane.append_event(event).await.expect("lane append");
2501
2502        assert_eq!(
2503            rt.resolve_prefix_unfiltered(&prefix)
2504                .await
2505                .expect("post-insert resolve"),
2506            Some(event_id),
2507            "a sidecar-only event id must resolve by hex prefix"
2508        );
2509    }
2510
2511    #[test]
2512    fn backend_data_dir_returns_none_for_from_backend_with_memory() {
2513        let backend = Arc::new(StorageBackend::memory().expect("memory backend"));
2514        let config = RuntimeConfig {
2515            web: Default::default(),
2516            telemetry: Default::default(),
2517            mounts: Vec::new(),
2518            brain: Default::default(),
2519            git_write: Default::default(),
2520            display_timezone: chrono_tz::Tz::UTC,
2521            events_split: None,
2522            db_path: None,
2523            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2524            default_namespace: Namespace::local(),
2525            embedding_model: None,
2526            additional_embedding_models: vec![],
2527            gate: Arc::new(AllowAllGate),
2528            packs: vec!["kg".to_string()],
2529            backend_id: BackendId::main(),
2530            brain_profile: None,
2531            visible_namespaces: vec![],
2532            allowed_outbound_namespaces: vec![],
2533            actor_id: None,
2534            exec: Default::default(),
2535        };
2536        let rt = KhiveRuntime::from_backend(backend, config);
2537        assert!(rt.backend_data_dir().is_none());
2538    }
2539
2540    #[test]
2541    fn file_runtime_creates_successfully() {
2542        let dir = tempfile::tempdir().unwrap();
2543        let path = dir.path().join("test.db");
2544        let config = RuntimeConfig {
2545            web: Default::default(),
2546            telemetry: Default::default(),
2547            mounts: Vec::new(),
2548            brain: Default::default(),
2549            git_write: Default::default(),
2550            display_timezone: chrono_tz::Tz::UTC,
2551            events_split: None,
2552            db_path: Some(path.clone()),
2553            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2554            default_namespace: Namespace::parse("test").unwrap(),
2555            embedding_model: None,
2556            additional_embedding_models: vec![],
2557            gate: Arc::new(AllowAllGate),
2558            packs: vec!["kg".to_string()],
2559            backend_id: BackendId::main(),
2560            brain_profile: None,
2561            visible_namespaces: vec![],
2562            allowed_outbound_namespaces: vec![],
2563            actor_id: None,
2564            exec: Default::default(),
2565        };
2566        let rt = KhiveRuntime::new_for_test(config).expect("file runtime should create");
2567        assert!(path.exists());
2568        assert_eq!(rt.config().default_namespace.as_str(), "test");
2569    }
2570
2571    #[cfg(unix)]
2572    #[tokio::test]
2573    async fn normal_boot_detects_read_only_snapshot_and_skips_model_registration() {
2574        use std::os::unix::fs::PermissionsExt;
2575
2576        let dir = tempfile::tempdir().unwrap();
2577        let path = dir.path().join("read_only_runtime.db");
2578        let base = RuntimeConfig {
2579            web: Default::default(),
2580            telemetry: Default::default(),
2581            mounts: Vec::new(),
2582            brain: Default::default(),
2583            git_write: Default::default(),
2584            display_timezone: chrono_tz::Tz::UTC,
2585            events_split: None,
2586            db_path: Some(path.clone()),
2587            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2588            default_namespace: Namespace::local(),
2589            embedding_model: None,
2590            additional_embedding_models: vec![],
2591            gate: Arc::new(AllowAllGate),
2592            packs: vec!["kg".to_string()],
2593            backend_id: BackendId::main(),
2594            brain_profile: None,
2595            visible_namespaces: vec![],
2596            allowed_outbound_namespaces: vec![],
2597            actor_id: None,
2598            exec: Default::default(),
2599        };
2600        {
2601            let writable =
2602                KhiveRuntime::new_for_test(base.clone()).expect("create migrated snapshot");
2603            assert!(writable
2604                .list_embedding_models(None)
2605                .await
2606                .expect("registry query")
2607                .is_empty());
2608        }
2609
2610        let mut permissions = std::fs::metadata(&path).unwrap().permissions();
2611        permissions.set_mode(0o444);
2612        std::fs::set_permissions(&path, permissions).unwrap();
2613        // A lingering writable `-shm` from the writable fixture's asynchronous
2614        // connection close is rejected by read-only admission as potentially
2615        // live; freeze any sidecars into the documented frozen-snapshot form.
2616        khive_storage::test_support::freeze_snapshot_sidecars(&path);
2617
2618        let read_only_config = RuntimeConfig {
2619            embedding_model: Some(EmbeddingModel::AllMiniLmL6V2),
2620            ..base
2621        };
2622        let runtime = KhiveRuntime::new_for_test(read_only_config)
2623            .expect("read-only boot must validate instead of migrating/registering");
2624        assert!(runtime.is_read_only());
2625        assert_eq!(
2626            runtime.backend().pool().writer_acquisition_snapshot(),
2627            khive_db::pool::WriterAcquisitionSnapshot::default(),
2628            "the construction-inclusive acquisition baseline must stay at zero"
2629        );
2630        assert!(
2631            runtime
2632                .list_embedding_models(None)
2633                .await
2634                .expect("read-only registry query")
2635                .is_empty(),
2636            "configured models must remain in-memory only during read-only boot"
2637        );
2638    }
2639
2640    #[test]
2641    fn explicit_readonly_constructor_uses_read_only_pool_even_on_writable_file_mode() {
2642        let dir = tempfile::tempdir().unwrap();
2643        let path = dir.path().join("explicit_read_only_runtime.db");
2644        let config = RuntimeConfig {
2645            web: Default::default(),
2646            telemetry: Default::default(),
2647            mounts: Vec::new(),
2648            brain: Default::default(),
2649            git_write: Default::default(),
2650            display_timezone: chrono_tz::Tz::UTC,
2651            events_split: None,
2652            db_path: Some(path.clone()),
2653            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2654            default_namespace: Namespace::local(),
2655            embedding_model: None,
2656            additional_embedding_models: vec![],
2657            gate: Arc::new(AllowAllGate),
2658            packs: vec!["kg".to_string()],
2659            backend_id: BackendId::main(),
2660            brain_profile: None,
2661            visible_namespaces: vec![],
2662            allowed_outbound_namespaces: vec![],
2663            actor_id: None,
2664            exec: Default::default(),
2665        };
2666        KhiveRuntime::new_for_test(config.clone()).expect("create migrated database");
2667        #[cfg(unix)]
2668        khive_storage::test_support::freeze_snapshot_sidecars(&path);
2669
2670        let runtime = KhiveRuntime::new_readonly_for_test(config).expect("explicit read-only boot");
2671        assert!(runtime.is_read_only());
2672        assert_eq!(
2673            runtime.backend().pool().writer_acquisition_snapshot(),
2674            khive_db::pool::WriterAcquisitionSnapshot::default(),
2675            "explicit read-only construction must validate through a reader without ever \
2676             acquiring the writer"
2677        );
2678    }
2679
2680    /// Grants Write on the primary namespace and Read, but not Write, on the
2681    /// extra namespace. The pseudo-verb split is what lets a right-aware gate
2682    /// express ADR-129's asymmetric authority contract.
2683    #[derive(Debug)]
2684    struct ReadOnlyExtraGate {
2685        primary: &'static str,
2686        extra: &'static str,
2687    }
2688
2689    impl khive_gate::Gate for ReadOnlyExtraGate {
2690        fn check(
2691            &self,
2692            req: &khive_gate::GateRequest,
2693        ) -> Result<khive_gate::GateDecision, khive_gate::GateError> {
2694            let allowed = match req.verb.as_str() {
2695                "authorize" => req.namespace.as_str() == self.primary,
2696                "authorize.visible" => req.namespace.as_str() == self.extra,
2697                _ => false,
2698            };
2699            if allowed {
2700                Ok(khive_gate::GateDecision::allow())
2701            } else {
2702                Ok(khive_gate::GateDecision::Deny {
2703                    reason: format!(
2704                        "{} denied for namespace {:?}",
2705                        req.verb,
2706                        req.namespace.as_str()
2707                    ),
2708                })
2709            }
2710        }
2711    }
2712
2713    #[test]
2714    fn authorize_with_visibility_allows_read_only_extra_namespace() {
2715        let primary = Namespace::parse("lambda:caller").expect("primary");
2716        let extra = Namespace::parse("lambda:read-only").expect("extra");
2717        let config = RuntimeConfig {
2718            db_path: None,
2719            packs: vec!["kg".to_string()],
2720            brain_profile: None,
2721            actor_id: None,
2722            gate: Arc::new(ReadOnlyExtraGate {
2723                primary: "lambda:caller",
2724                extra: "lambda:read-only",
2725            }),
2726            ..RuntimeConfig::no_embeddings()
2727        };
2728        let rt = KhiveRuntime::new(config).expect("memory runtime");
2729
2730        let token = rt
2731            .authorize_with_visibility(primary, vec![extra.clone()])
2732            .expect("Write on primary and Read on extra must mint");
2733        assert!(token.visible_namespaces().contains(&extra));
2734    }
2735
2736    /// Denies exactly one namespace; every other request is allowed. Lets the
2737    /// test below prove a refusal comes from the per-extra visibility check
2738    /// rather than from the primary authorization.
2739    #[derive(Debug)]
2740    struct DenyNamespaceGate {
2741        deny: &'static str,
2742    }
2743
2744    impl khive_gate::Gate for DenyNamespaceGate {
2745        fn check(
2746            &self,
2747            req: &khive_gate::GateRequest,
2748        ) -> Result<khive_gate::GateDecision, khive_gate::GateError> {
2749            if req.namespace.as_str() == self.deny {
2750                Ok(khive_gate::GateDecision::Deny {
2751                    reason: "namespace denied by policy".to_string(),
2752                })
2753            } else {
2754                Ok(khive_gate::GateDecision::allow())
2755            }
2756        }
2757    }
2758
2759    #[test]
2760    fn authorize_with_visibility_denies_missing_extra_read() {
2761        let config = RuntimeConfig {
2762            db_path: None,
2763            packs: vec!["kg".to_string()],
2764            brain_profile: None,
2765            actor_id: None,
2766            gate: Arc::new(DenyNamespaceGate {
2767                deny: "lambda:secret",
2768            }),
2769            ..RuntimeConfig::no_embeddings()
2770        };
2771        let rt = KhiveRuntime::new(config).expect("memory runtime");
2772        let primary = Namespace::parse("lambda:caller").expect("primary");
2773        let denied = Namespace::parse("lambda:secret").expect("denied");
2774        let allowed = Namespace::parse("lambda:open").expect("allowed");
2775
2776        // Control: the same mint without the denied namespace succeeds, so
2777        // the refusal below can only come from the per-extra check.
2778        rt.authorize_with_visibility(primary.clone(), vec![allowed.clone()])
2779            .expect("mint with only allowed extras");
2780
2781        let err = rt
2782            .authorize_with_visibility(primary, vec![allowed, denied])
2783            .expect_err("a denied extra namespace must refuse the whole mint");
2784        let msg = err.to_string();
2785        assert!(
2786            msg.contains("lambda:secret"),
2787            "refusal must name the offending namespace: {msg}"
2788        );
2789    }
2790
2791    /// Build a migrated database and reopen it read-only, returning the
2792    /// tempdir that keeps it alive alongside the runtime.
2793    fn make_read_only_runtime() -> (tempfile::TempDir, KhiveRuntime) {
2794        let dir = tempfile::tempdir().unwrap();
2795        let path = dir.path().join("read_only_blob_seam.db");
2796        let config = RuntimeConfig {
2797            web: Default::default(),
2798            telemetry: Default::default(),
2799            mounts: Vec::new(),
2800            brain: Default::default(),
2801            git_write: Default::default(),
2802            display_timezone: chrono_tz::Tz::UTC,
2803            db_path: Some(path.clone()),
2804            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2805            default_namespace: Namespace::local(),
2806            embedding_model: None,
2807            additional_embedding_models: vec![],
2808            gate: Arc::new(AllowAllGate),
2809            packs: vec!["kg".to_string()],
2810            backend_id: BackendId::main(),
2811            brain_profile: None,
2812            visible_namespaces: vec![],
2813            allowed_outbound_namespaces: vec![],
2814            actor_id: None,
2815            events_split: None,
2816            exec: Default::default(),
2817        };
2818        KhiveRuntime::new_for_test(config.clone()).expect("create migrated database");
2819        #[cfg(unix)]
2820        khive_storage::test_support::freeze_snapshot_sidecars(&path);
2821        let runtime = KhiveRuntime::new_readonly_for_test(config).expect("read-only boot");
2822        assert!(runtime.is_read_only());
2823        (dir, runtime)
2824    }
2825
2826    #[tokio::test]
2827    async fn install_blob_store_on_read_only_runtime_refuses_mutators() {
2828        let (_dir, runtime) = make_read_only_runtime();
2829
2830        // Control: the raw store IS writable — seed an object through it —
2831        // so the refusal below can only come from the install-seam wrap.
2832        use khive_storage::BlobStore as _;
2833        let blob_root = tempfile::tempdir().unwrap();
2834        let writable = Arc::new(
2835            khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
2836                .expect("fs blob store"),
2837        );
2838        let seeded = writable
2839            .put(b"seed".to_vec())
2840            .await
2841            .expect("seed put through the raw store");
2842        runtime
2843            .install_blob_store(writable.clone())
2844            .expect("read-only install wraps rather than refusing");
2845
2846        let installed = runtime.blob_store().expect("installed store");
2847        assert!(
2848            installed.exists(&seeded).await.expect("exists"),
2849            "bounded read surface must stay available"
2850        );
2851        let err = installed
2852            .put(b"post-boot".to_vec())
2853            .await
2854            .expect_err("put must refuse on a read-only runtime");
2855        assert!(
2856            err.to_string().contains("read-only"),
2857            "refusal must name the mode: {err}"
2858        );
2859
2860        // Reinstalling the SAME raw store stays idempotent even though the
2861        // installed hydrator holds a wrapper around it: identity is checked
2862        // against the raw handle the hydrator remembers.
2863        runtime
2864            .install_blob_store(writable.clone())
2865            .expect("reinstalling the same raw store must be idempotent");
2866
2867        // The hydrator seam holds the mode too: pairing a writable store
2868        // with `BlobHydrator::new` and installing it directly must refuse,
2869        // or it is a public bypass of everything above.
2870        let bypass_root = tempfile::tempdir().unwrap();
2871        let bypass_store = Arc::new(
2872            khive_db::stores::blob::FsBlobStore::new(bypass_root.path().to_path_buf(), 0)
2873                .expect("fs blob store"),
2874        );
2875        let writable_hydrator = Arc::new(
2876            crate::BlobHydrator::new(bypass_store, crate::DEFAULT_BLOB_HYDRATION_BYTES)
2877                .expect("construct writable hydrator"),
2878        );
2879        let err = runtime
2880            .install_blob_hydrator(writable_hydrator)
2881            .expect_err("a writable hydrator must be refused on a read-only runtime");
2882        assert!(
2883            err.to_string().contains("read-only"),
2884            "refusal must name the mode: {err}"
2885        );
2886    }
2887
2888    #[tokio::test]
2889    async fn mode_aware_hydrator_constructor_satisfies_read_only_install() {
2890        // The boot path constructs its hydrator with `for_mode` from the
2891        // blob runtime's configured mode; a read-only construction must pass
2892        // the read-only install gate, serve bounded reads, refuse mutation,
2893        // and treat a second hydrator allocation over the same raw store as
2894        // the same pairing (the concurrent-first-install loser shape).
2895        let (_dir, runtime) = make_read_only_runtime();
2896        let blob_root = tempfile::tempdir().unwrap();
2897        let raw = Arc::new(
2898            khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
2899                .expect("fs blob store"),
2900        ) as Arc<dyn khive_storage::BlobStore>;
2901        let seeded = raw.put(b"seed".to_vec()).await.expect("seed put");
2902        let hydrator = Arc::new(
2903            crate::BlobHydrator::for_mode(
2904                Arc::clone(&raw),
2905                crate::DEFAULT_BLOB_HYDRATION_BYTES,
2906                true,
2907            )
2908            .expect("mode-aware read-only construction"),
2909        );
2910        runtime
2911            .install_blob_hydrator(hydrator)
2912            .expect("a for_mode(read_only) hydrator must pass the read-only gate");
2913        let installed = runtime.blob_store().expect("installed store");
2914        assert!(installed.exists(&seeded).await.expect("exists"));
2915        installed
2916            .put(b"post".to_vec())
2917            .await
2918            .expect_err("mutation must refuse through the wrapped store");
2919
2920        // A DISTINCT hydrator allocation over the same raw store, budget,
2921        // and mode is the same pairing: installing it must be idempotent,
2922        // not a conflicting-install error.
2923        let twin = Arc::new(
2924            crate::BlobHydrator::for_mode(
2925                Arc::clone(&raw),
2926                crate::DEFAULT_BLOB_HYDRATION_BYTES,
2927                true,
2928            )
2929            .expect("twin construction"),
2930        );
2931        runtime
2932            .install_blob_hydrator(twin)
2933            .expect("an equivalent pairing must read as idempotent");
2934    }
2935
2936    #[tokio::test]
2937    async fn shared_install_permits_governed_writable_hydrator_on_read_only_handle() {
2938        // The documented multi-backend matrix includes a writable blob
2939        // secondary beside a read-only main: boot installs ONE writable
2940        // hydrator on every handle, including read-only ones. The plain
2941        // seam must still refuse that pairing, and the shared seam accepts
2942        // it ONLY when the hydrator's mode was derived from a governing
2943        // backend — a hand-paired writable hydrator is refused too, so a
2944        // safe downstream caller cannot use the shared seam to defeat a
2945        // read-only handle's guarantee.
2946        let (_dir, runtime) = make_read_only_runtime();
2947        let blob_root = tempfile::tempdir().unwrap();
2948        let raw = Arc::new(
2949            khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
2950                .expect("fs blob store"),
2951        ) as Arc<dyn khive_storage::BlobStore>;
2952        let hand_paired = Arc::new(
2953            crate::BlobHydrator::for_mode(raw, crate::DEFAULT_BLOB_HYDRATION_BYTES, false)
2954                .expect("writable construction"),
2955        );
2956        runtime
2957            .install_blob_hydrator(Arc::clone(&hand_paired))
2958            .expect_err("the plain seam must refuse a writable hydrator on a read-only handle");
2959        runtime
2960            .install_shared_blob_hydrator(hand_paired)
2961            .expect_err("the shared seam must refuse a hand-paired (ungoverned) hydrator");
2962
2963        // The sanctioned path: a WRITABLE governing backend (the blob pack's
2964        // backend in the mixed-mode topology) derives a governed writable
2965        // hydrator, and the shared seam installs it on the read-only handle.
2966        let governing = khive_db::StorageBackend::memory().expect("memory backend");
2967        let cfg = crate::KhiveConfig {
2968            storage: crate::engine_config::StorageSectionConfig {
2969                blob: Some(crate::engine_config::BlobConfig::Fs {
2970                    root: Some(blob_root.path().to_string_lossy().into_owned()),
2971                    floor_bytes: Some(0),
2972                }),
2973            },
2974            ..crate::KhiveConfig::default()
2975        };
2976        let governed = Arc::new(
2977            crate::BlobHydrator::resolve_for_governing_backend(
2978                &cfg,
2979                &governing,
2980                &governing,
2981                crate::DEFAULT_BLOB_HYDRATION_BYTES,
2982            )
2983            .expect("governed construction"),
2984        );
2985        runtime
2986            .install_shared_blob_hydrator(governed)
2987            .expect("the shared seam accepts a governed hydrator");
2988        let installed = runtime.blob_store().expect("installed store");
2989        let put = installed
2990            .put(b"shared-write".to_vec())
2991            .await
2992            .expect("the governed writable capability must actually mutate");
2993        assert!(installed.exists(&put).await.expect("exists"));
2994    }
2995
2996    /// A `~/`-prefixed `--db`/`KHIVE_DB` override must resolve, boot, and
2997    /// fingerprint identically to the equivalent absolute path. Before this
2998    /// fix, `resolve_db_anchor` left a leading `~` unexpanded in
2999    /// `RuntimeConfig.db_path`, so single-backend boot (`KhiveRuntime::new`)
3000    /// opened a literal `./~/...` file under the process cwd while
3001    /// `compute_config_id` (which canonicalizes/expands separately) still
3002    /// fingerprinted the real `$HOME` path — two processes pointed at the
3003    /// same logical database could open different files yet share a
3004    /// `config_id`, letting daemon dispatch route requests to the wrong one.
3005    #[test]
3006    #[serial]
3007    fn tilde_prefixed_db_override_resolves_and_boots_like_the_absolute_equivalent() {
3008        if crate::test_process::run_in_child() {
3009            return;
3010        }
3011
3012        let original_home = std::env::var_os("HOME");
3013        let original_cwd = std::env::current_dir().expect("read cwd");
3014        let home_dir = tempfile::tempdir().expect("home tempdir");
3015        let work_dir = tempfile::tempdir().expect("work tempdir");
3016        std::env::set_var("HOME", home_dir.path());
3017        std::env::set_current_dir(work_dir.path()).expect("chdir into isolated work dir");
3018
3019        let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3020            let tilde_anchor = crate::config::resolve_db_anchor(Some("~/data.db"))
3021                .expect("an explicit path always anchors");
3022            let expected = home_dir.path().join("data.db");
3023            assert_eq!(
3024                tilde_anchor, expected,
3025                "resolve_db_anchor must expand a leading ~ to $HOME before it ever \
3026                 reaches RuntimeConfig.db_path"
3027            );
3028
3029            let absolute_anchor = crate::config::resolve_db_anchor(Some(
3030                expected.to_str().expect("utf8 tempdir path"),
3031            ))
3032            .expect("an explicit path always anchors");
3033            assert_eq!(
3034                tilde_anchor, absolute_anchor,
3035                "a ~-prefixed override and its equivalent absolute path must resolve to \
3036                 the identical anchor"
3037            );
3038
3039            let make_config = |db_path: std::path::PathBuf| RuntimeConfig {
3040                web: Default::default(),
3041                telemetry: Default::default(),
3042                mounts: Vec::new(),
3043                brain: Default::default(),
3044                git_write: Default::default(),
3045                display_timezone: chrono_tz::Tz::UTC,
3046                events_split: None,
3047                db_path: Some(db_path),
3048                blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3049                default_namespace: Namespace::local(),
3050                embedding_model: None,
3051                additional_embedding_models: vec![],
3052                gate: Arc::new(AllowAllGate),
3053                packs: vec!["kg".to_string()],
3054                backend_id: BackendId::main(),
3055                brain_profile: None,
3056                visible_namespaces: vec![],
3057                allowed_outbound_namespaces: vec![],
3058                actor_id: None,
3059                exec: Default::default(),
3060            };
3061
3062            let tilde_cfg = make_config(tilde_anchor.clone());
3063
3064            let rt =
3065                KhiveRuntime::new_for_test(tilde_cfg).expect("boot must open the expanded path");
3066            assert_eq!(
3067                rt.backend_data_dir().expect("file backend"),
3068                home_dir.path(),
3069                "single-backend boot must open the file under the expanded $HOME \
3070                 directory, not a literal ~ path relative to cwd"
3071            );
3072            assert!(
3073                expected.exists(),
3074                "the database file must be created at the expanded $HOME path"
3075            );
3076            assert!(
3077                !work_dir.path().join("~").exists(),
3078                "boot must never create a literal '~' directory under the process cwd"
3079            );
3080        }));
3081
3082        match &original_home {
3083            Some(h) => std::env::set_var("HOME", h),
3084            None => std::env::remove_var("HOME"),
3085        }
3086        let _ = std::env::set_current_dir(&original_cwd);
3087        outcome.expect("test body panicked");
3088    }
3089
3090    #[test]
3091    fn from_backend_uses_provided_backend() {
3092        let backend = Arc::new(StorageBackend::memory().expect("memory backend"));
3093        let config = RuntimeConfig {
3094            web: Default::default(),
3095            telemetry: Default::default(),
3096            mounts: Vec::new(),
3097            brain: Default::default(),
3098            git_write: Default::default(),
3099            display_timezone: chrono_tz::Tz::UTC,
3100            events_split: None,
3101            db_path: None,
3102            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3103            default_namespace: Namespace::local(),
3104            embedding_model: None,
3105            additional_embedding_models: vec![],
3106            gate: Arc::new(AllowAllGate),
3107            packs: vec!["kg".to_string()],
3108            backend_id: BackendId::parse("lore").expect("valid backend id"),
3109            brain_profile: None,
3110            visible_namespaces: vec![],
3111            allowed_outbound_namespaces: vec![],
3112            actor_id: None,
3113            exec: Default::default(),
3114        };
3115        let rt = KhiveRuntime::from_backend(backend, config);
3116        assert_eq!(rt.backend_id().as_str(), "lore");
3117        assert!(rt.config().db_path.is_none());
3118    }
3119
3120    #[test]
3121    fn backend_id_defaults_to_main() {
3122        let rt = KhiveRuntime::memory().unwrap();
3123        assert_eq!(rt.backend_id().as_str(), BackendId::MAIN);
3124    }
3125
3126    #[test]
3127    fn store_accessors_return_ok() {
3128        let rt = KhiveRuntime::memory().unwrap();
3129        let tok = NamespaceToken::local();
3130        assert!(rt.entities(&tok).is_ok());
3131        assert!(rt.graph(&tok).is_ok());
3132        assert!(rt.notes(&tok).is_ok());
3133        assert!(rt.events(&tok).is_ok());
3134    }
3135
3136    fn attributed_event_runtime() -> (KhiveRuntime, NamespaceToken) {
3137        let runtime = KhiveRuntime::new(RuntimeConfig {
3138            db_path: None,
3139            actor_id: Some("lambda:enrolled".to_string()),
3140            ..RuntimeConfig::no_embeddings()
3141        })
3142        .expect("memory runtime");
3143        let token = runtime
3144            .authorize(Namespace::local())
3145            .expect("configured actor is allowed by the default gate");
3146        (runtime, token)
3147    }
3148
3149    fn forged_event(verb: &str) -> Event {
3150        Event::new(
3151            "caller-selected-namespace",
3152            verb,
3153            EventKind::Audit,
3154            SubstrateKind::Event,
3155            "caller-selected-actor",
3156        )
3157    }
3158
3159    fn assert_token_attribution(event: &Event) {
3160        assert_eq!(event.namespace, "local");
3161        assert_eq!(event.actor, "actor:lambda:enrolled");
3162    }
3163
3164    #[tokio::test]
3165    async fn token_scoped_event_store_stamps_resolved_attribution_on_single_append() {
3166        let (runtime, token) = attributed_event_runtime();
3167        let store = runtime.events(&token).expect("event store");
3168        let event = forged_event("single");
3169        let id = event.id;
3170
3171        store.append_event(event).await.expect("append");
3172
3173        let stored = store
3174            .get_event(id)
3175            .await
3176            .expect("read")
3177            .expect("the token-stamped event remains visible to the token");
3178        assert_token_attribution(&stored);
3179        assert_eq!(stored.verb, "single", "non-attribution fields survive");
3180    }
3181
3182    #[tokio::test]
3183    async fn token_scoped_event_store_stamps_resolved_attribution_on_batch_paths() {
3184        let (runtime, token) = attributed_event_runtime();
3185        let store = runtime.events(&token).expect("event store");
3186        let ordinary = forged_event("batch");
3187        let ordinary_id = ordinary.id;
3188
3189        store
3190            .append_events(vec![ordinary])
3191            .await
3192            .expect("ordinary batch append");
3193        let stored = store
3194            .get_event(ordinary_id)
3195            .await
3196            .expect("read")
3197            .expect("ordinary batch event remains visible to the token");
3198        assert_token_attribution(&stored);
3199
3200        let idempotent = forged_event("idempotent_batch");
3201        let idempotent_id = idempotent.id;
3202        let outcome = store
3203            .append_events_idempotent(vec![idempotent])
3204            .await
3205            .expect("idempotent batch append");
3206        assert_eq!(
3207            outcome.rows,
3208            vec![khive_storage::event::EventAppendDisposition::Inserted]
3209        );
3210        let stored = store
3211            .get_event(idempotent_id)
3212            .await
3213            .expect("read")
3214            .expect("idempotent batch event remains visible to the token");
3215        assert_token_attribution(&stored);
3216    }
3217
3218    #[test]
3219    fn vectors_returns_unconfigured_without_model() {
3220        let rt = KhiveRuntime::memory().unwrap();
3221        let tok = NamespaceToken::local();
3222        match rt.vectors(&tok) {
3223            Err(crate::RuntimeError::Unconfigured(s)) => assert_eq!(s, "embedding_model"),
3224            Err(other) => panic!("expected Unconfigured, got {:?}", other),
3225            Ok(_) => panic!("expected Err, got Ok"),
3226        }
3227    }
3228
3229    #[test]
3230    fn vec_model_key_sanitizes_dots_and_dashes() {
3231        assert_eq!(
3232            vec_model_key(EmbeddingModel::BgeSmallEnV15),
3233            "bge_small_en_v1_5"
3234        );
3235        assert_eq!(
3236            vec_model_key(EmbeddingModel::BgeBaseEnV15),
3237            "bge_base_en_v1_5"
3238        );
3239        assert_eq!(
3240            vec_model_key(EmbeddingModel::AllMiniLmL6V2),
3241            "all_minilm_l6_v2"
3242        );
3243    }
3244
3245    #[test]
3246    fn default_config_uses_allow_all_gate() {
3247        let cfg = RuntimeConfig::default();
3248        assert_eq!(cfg.default_namespace.as_str(), "local");
3249        let _: GateRef = cfg.gate.clone();
3250    }
3251
3252    #[test]
3253    fn parse_pack_list_handles_comma_and_whitespace() {
3254        assert_eq!(parse_pack_list("kg"), vec!["kg".to_string()]);
3255        assert_eq!(
3256            parse_pack_list("kg,gtd"),
3257            vec!["kg".to_string(), "gtd".to_string()]
3258        );
3259        assert_eq!(
3260            parse_pack_list("  kg ,  gtd  "),
3261            vec!["kg".to_string(), "gtd".to_string()]
3262        );
3263        assert_eq!(
3264            parse_pack_list("kg gtd"),
3265            vec!["kg".to_string(), "gtd".to_string()]
3266        );
3267        assert_eq!(parse_pack_list(",,"), Vec::<String>::new());
3268        assert_eq!(parse_pack_list(""), Vec::<String>::new());
3269    }
3270
3271    #[test]
3272    fn default_config_packs_loads_production_set() {
3273        let prior = std::env::var("KHIVE_PACKS").ok();
3274        // SAFETY: test function runs single-threaded; no other threads read or write KHIVE_PACKS.
3275        unsafe {
3276            std::env::remove_var("KHIVE_PACKS");
3277        }
3278        // The default distribution loads all production packs.
3279        let cfg = RuntimeConfig::default();
3280        assert_eq!(cfg.packs, RuntimeConfig::built_in_packs());
3281        assert!(cfg.packs.contains(&"kg".to_string()));
3282        assert!(cfg.packs.contains(&"gtd".to_string()));
3283        assert!(cfg.packs.contains(&"memory".to_string()));
3284        assert!(cfg.packs.contains(&"brain".to_string()));
3285        assert!(cfg.packs.contains(&"comm".to_string()));
3286        assert!(cfg.packs.contains(&"schedule".to_string()));
3287        assert!(cfg.packs.contains(&"knowledge".to_string()));
3288        // session loads by default so its background mirror warm-hook runs in
3289        // production; its handlers are all operator-only subhandlers (0 wire verbs).
3290        assert!(cfg.packs.contains(&"session".to_string()));
3291        assert!(cfg.packs.contains(&"git".to_string()));
3292        assert!(cfg.packs.contains(&"code".to_string()));
3293        assert!(cfg.packs.contains(&"workspace".to_string()));
3294        // blob loads by default; a normal file-backed boot installs a
3295        // default FsBlobStore beside the database file with no config
3296        // needed, so its verbs are live in default deployments too (only an
3297        // in-memory backend leaves them unconfigured).
3298        assert!(cfg.packs.contains(&"blob".to_string()));
3299        // tool loads by default: the registry, discovery and use-policy verbs
3300        // (ADR-180) are live in default deployments.
3301        assert!(cfg.packs.contains(&"tool".to_string()));
3302        assert!(cfg.packs.contains(&"exec".to_string()));
3303        assert_eq!(cfg.packs.len(), 14);
3304        if let Some(v) = prior {
3305            // SAFETY: single-threaded test cleanup; restores KHIVE_PACKS to its prior value.
3306            unsafe {
3307                std::env::set_var("KHIVE_PACKS", v);
3308            }
3309        }
3310    }
3311
3312    #[test]
3313    fn default_config_uses_minilm_when_env_unset() {
3314        let prior = std::env::var("KHIVE_EMBEDDING_MODEL").ok();
3315        // SAFETY: tests are serial by default for env mutation here; if other tests
3316        // mutate this var, mark them with the same scope.
3317        unsafe {
3318            std::env::remove_var("KHIVE_EMBEDDING_MODEL");
3319        }
3320        let cfg = RuntimeConfig::default();
3321        assert_eq!(cfg.embedding_model, Some(EmbeddingModel::AllMiniLmL6V2));
3322        if let Some(v) = prior {
3323            // SAFETY: single-threaded test cleanup; restores KHIVE_EMBEDDING_MODEL to its prior value.
3324            unsafe {
3325                std::env::set_var("KHIVE_EMBEDDING_MODEL", v);
3326            }
3327        }
3328    }
3329
3330    // ---- Actor config tests ----
3331
3332    use crate::engine_config::{ActorConfig, KhiveConfig};
3333
3334    fn khive_cfg_with_actor(id: &str) -> KhiveConfig {
3335        KhiveConfig {
3336            engines: vec![],
3337            actor: ActorConfig {
3338                id: Some(id.to_string()),
3339                display_name: None,
3340                ..Default::default()
3341            },
3342            ..KhiveConfig::default()
3343        }
3344    }
3345
3346    #[test]
3347    fn runtime_config_from_khive_config_actor_id_does_not_override_default_namespace() {
3348        // `[actor] id` must not set `default_namespace`: writes stay pinned to
3349        // `local`. A non-`'local'` actor.id is folded into the default read
3350        // visible-set, but that does not change default_namespace. This test
3351        // asserts the write-routing invariant only.
3352        let base = RuntimeConfig {
3353            web: Default::default(),
3354            telemetry: Default::default(),
3355            mounts: Vec::new(),
3356            brain: Default::default(),
3357            git_write: Default::default(),
3358            display_timezone: chrono_tz::Tz::UTC,
3359            events_split: None,
3360            db_path: None,
3361            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3362            default_namespace: Namespace::local(),
3363            embedding_model: None,
3364            additional_embedding_models: vec![],
3365            gate: Arc::new(AllowAllGate),
3366            packs: vec!["kg".to_string()],
3367            backend_id: BackendId::main(),
3368            brain_profile: None,
3369            visible_namespaces: vec![],
3370            allowed_outbound_namespaces: vec![],
3371            actor_id: None,
3372            exec: Default::default(),
3373        };
3374        let cfg = khive_cfg_with_actor("lambda:khive");
3375        let result = runtime_config_from_khive_config(&cfg, base);
3376        assert_eq!(
3377            result.default_namespace.as_str(),
3378            "local",
3379            "actor.id must not become default_namespace (ADR-007 Rev 4 Rule 0); writes pin to local"
3380        );
3381    }
3382
3383    #[test]
3384    fn runtime_config_from_khive_config_empty_actor_id_keeps_base_namespace() {
3385        let base = RuntimeConfig {
3386            web: Default::default(),
3387            telemetry: Default::default(),
3388            mounts: Vec::new(),
3389            brain: Default::default(),
3390            git_write: Default::default(),
3391            display_timezone: chrono_tz::Tz::UTC,
3392            events_split: None,
3393            db_path: None,
3394            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3395            default_namespace: Namespace::parse("lambda:base").unwrap(),
3396            embedding_model: None,
3397            additional_embedding_models: vec![],
3398            gate: Arc::new(AllowAllGate),
3399            packs: vec!["kg".to_string()],
3400            backend_id: BackendId::main(),
3401            brain_profile: None,
3402            visible_namespaces: vec![],
3403            allowed_outbound_namespaces: vec![],
3404            actor_id: None,
3405            exec: Default::default(),
3406        };
3407        let cfg = KhiveConfig {
3408            engines: vec![],
3409            actor: ActorConfig {
3410                id: Some(String::new()),
3411                display_name: None,
3412                ..Default::default()
3413            },
3414            ..KhiveConfig::default()
3415        };
3416        let result = runtime_config_from_khive_config(&cfg, base);
3417        assert_eq!(
3418            result.default_namespace.as_str(),
3419            "lambda:base",
3420            "empty actor.id must not override base namespace"
3421        );
3422    }
3423
3424    #[test]
3425    fn runtime_config_from_khive_config_absent_actor_id_keeps_base_namespace() {
3426        let base = RuntimeConfig {
3427            web: Default::default(),
3428            telemetry: Default::default(),
3429            mounts: Vec::new(),
3430            brain: Default::default(),
3431            git_write: Default::default(),
3432            display_timezone: chrono_tz::Tz::UTC,
3433            events_split: None,
3434            db_path: None,
3435            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3436            default_namespace: Namespace::parse("lambda:base").unwrap(),
3437            embedding_model: None,
3438            additional_embedding_models: vec![],
3439            gate: Arc::new(AllowAllGate),
3440            packs: vec!["kg".to_string()],
3441            backend_id: BackendId::main(),
3442            brain_profile: None,
3443            visible_namespaces: vec![],
3444            allowed_outbound_namespaces: vec![],
3445            actor_id: None,
3446            exec: Default::default(),
3447        };
3448        let cfg = KhiveConfig::default(); // no actor.id
3449        let result = runtime_config_from_khive_config(&cfg, base);
3450        assert_eq!(
3451            result.default_namespace.as_str(),
3452            "lambda:base",
3453            "absent actor.id must not override base namespace"
3454        );
3455    }
3456
3457    #[test]
3458    fn runtime_config_from_khive_config_actor_id_with_engines() {
3459        let base = RuntimeConfig {
3460            web: Default::default(),
3461            telemetry: Default::default(),
3462            mounts: Vec::new(),
3463            brain: Default::default(),
3464            git_write: Default::default(),
3465            display_timezone: chrono_tz::Tz::UTC,
3466            events_split: None,
3467            db_path: None,
3468            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3469            default_namespace: Namespace::local(),
3470            embedding_model: None,
3471            additional_embedding_models: vec![],
3472            gate: Arc::new(AllowAllGate),
3473            packs: vec!["kg".to_string()],
3474            backend_id: BackendId::main(),
3475            brain_profile: None,
3476            visible_namespaces: vec![],
3477            allowed_outbound_namespaces: vec![],
3478            actor_id: None,
3479            exec: Default::default(),
3480        };
3481        let cfg = KhiveConfig {
3482            engines: vec![crate::engine_config::EngineConfig {
3483                name: "default".to_string(),
3484                model: "all-minilm-l6-v2".to_string(),
3485                default: true,
3486                fusion_weight: None,
3487                dims: None,
3488            }],
3489            actor: ActorConfig {
3490                id: Some("lambda:test".to_string()),
3491                display_name: None,
3492                ..Default::default()
3493            },
3494            ..KhiveConfig::default()
3495        };
3496        let result = runtime_config_from_khive_config(&cfg, base);
3497        assert_eq!(
3498            result.default_namespace.as_str(),
3499            "local",
3500            "actor.id must not override default_namespace (ADR-007 Rev 4 Rule 0); \
3501             writes pin to local; engine config is still applied"
3502        );
3503        assert!(result.embedding_model.is_some());
3504    }
3505
3506    // ---- [display] timezone (ADR-169) wiring tests ----
3507
3508    #[test]
3509    fn runtime_config_from_khive_config_display_timezone_overrides_base() {
3510        let base = RuntimeConfig {
3511            web: Default::default(),
3512            telemetry: Default::default(),
3513            mounts: Vec::new(),
3514            brain: Default::default(),
3515            git_write: Default::default(),
3516            display_timezone: chrono_tz::Tz::UTC,
3517            events_split: None,
3518            db_path: None,
3519            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3520            default_namespace: Namespace::local(),
3521            embedding_model: None,
3522            additional_embedding_models: vec![],
3523            gate: Arc::new(AllowAllGate),
3524            packs: vec!["kg".to_string()],
3525            backend_id: BackendId::main(),
3526            brain_profile: None,
3527            visible_namespaces: vec![],
3528            allowed_outbound_namespaces: vec![],
3529            actor_id: None,
3530            exec: Default::default(),
3531        };
3532        let cfg = KhiveConfig {
3533            display: crate::engine_config::DisplaySectionConfig {
3534                timezone: Some("America/New_York".to_string()),
3535            },
3536            ..KhiveConfig::default()
3537        };
3538        let result = runtime_config_from_khive_config(&cfg, base);
3539        assert_eq!(
3540            result.display_timezone,
3541            "America/New_York".parse::<chrono_tz::Tz>().unwrap(),
3542            "[display] timezone in khive.toml must override base.display_timezone"
3543        );
3544    }
3545
3546    #[test]
3547    fn runtime_config_from_khive_config_absent_display_timezone_keeps_base() {
3548        let base = RuntimeConfig {
3549            web: Default::default(),
3550            telemetry: Default::default(),
3551            mounts: Vec::new(),
3552            brain: Default::default(),
3553            git_write: Default::default(),
3554            display_timezone: "Asia/Tokyo".parse().unwrap(),
3555            events_split: None,
3556            db_path: None,
3557            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3558            default_namespace: Namespace::local(),
3559            embedding_model: None,
3560            additional_embedding_models: vec![],
3561            gate: Arc::new(AllowAllGate),
3562            packs: vec!["kg".to_string()],
3563            backend_id: BackendId::main(),
3564            brain_profile: None,
3565            visible_namespaces: vec![],
3566            allowed_outbound_namespaces: vec![],
3567            actor_id: None,
3568            exec: Default::default(),
3569        };
3570        let cfg = KhiveConfig::default(); // no [display] section
3571        let result = runtime_config_from_khive_config(&cfg, base);
3572        assert_eq!(
3573            result.display_timezone,
3574            "Asia/Tokyo".parse::<chrono_tz::Tz>().unwrap(),
3575            "absent [display] timezone must preserve base.display_timezone unchanged"
3576        );
3577    }
3578
3579    // ---- base.actor_id (env-resolved actor) preservation tests ----
3580    //
3581    // Regression coverage: a project config found without an `[actor] id` used
3582    // to silently drop `base.actor_id` (e.g. the value `RuntimeConfig::default()`
3583    // read from `KHIVE_ACTOR`) because both return arms spread an unconditional
3584    // `actor_id: None` over `..base`. The fix falls back to `base.actor_id`
3585    // when the TOML supplies no `[actor] id`, in both arms.
3586
3587    #[test]
3588    #[serial]
3589    fn runtime_config_from_khive_config_engines_present_preserves_env_actor_when_toml_has_none() {
3590        let prior = std::env::var("KHIVE_ACTOR").ok();
3591        // SAFETY: test is #[serial]; no other test in this crate reads/writes KHIVE_ACTOR.
3592        unsafe {
3593            std::env::set_var("KHIVE_ACTOR", "lambda:test-env-actor");
3594        }
3595        let base = RuntimeConfig::default();
3596        assert_eq!(base.actor_id.as_deref(), Some("lambda:test-env-actor"));
3597
3598        let cfg = KhiveConfig {
3599            engines: vec![crate::engine_config::EngineConfig {
3600                name: "default".to_string(),
3601                model: "all-minilm-l6-v2".to_string(),
3602                default: true,
3603                fusion_weight: None,
3604                dims: None,
3605            }],
3606            actor: ActorConfig::default(), // no [actor] id
3607            ..KhiveConfig::default()
3608        };
3609        let result = runtime_config_from_khive_config(&cfg, base);
3610        assert_eq!(
3611            result.actor_id.as_deref(),
3612            Some("lambda:test-env-actor"),
3613            "engines-present arm must preserve base.actor_id (env actor) when TOML has no [actor] id"
3614        );
3615
3616        // SAFETY: restores prior KHIVE_ACTOR value (test cleanup).
3617        unsafe {
3618            match prior {
3619                Some(v) => std::env::set_var("KHIVE_ACTOR", v),
3620                None => std::env::remove_var("KHIVE_ACTOR"),
3621            }
3622        }
3623    }
3624
3625    #[test]
3626    #[serial]
3627    fn runtime_config_from_khive_config_engines_empty_preserves_env_actor_when_toml_has_none() {
3628        let prior = std::env::var("KHIVE_ACTOR").ok();
3629        // SAFETY: test is #[serial]; no other test in this crate reads/writes KHIVE_ACTOR.
3630        unsafe {
3631            std::env::set_var("KHIVE_ACTOR", "lambda:test-env-actor");
3632        }
3633        let base = RuntimeConfig::default();
3634        assert_eq!(base.actor_id.as_deref(), Some("lambda:test-env-actor"));
3635
3636        let cfg = KhiveConfig {
3637            engines: vec![],
3638            actor: ActorConfig::default(), // no [actor] id
3639            ..KhiveConfig::default()
3640        };
3641        let result = runtime_config_from_khive_config(&cfg, base);
3642        assert_eq!(
3643            result.actor_id.as_deref(),
3644            Some("lambda:test-env-actor"),
3645            "engines-empty early-return arm must preserve base.actor_id (env actor) when TOML has no [actor] id"
3646        );
3647
3648        // SAFETY: restores prior KHIVE_ACTOR value (test cleanup).
3649        unsafe {
3650            match prior {
3651                Some(v) => std::env::set_var("KHIVE_ACTOR", v),
3652                None => std::env::remove_var("KHIVE_ACTOR"),
3653            }
3654        }
3655    }
3656
3657    #[test]
3658    #[serial]
3659    fn runtime_config_from_khive_config_toml_actor_wins_over_env_actor() {
3660        let prior = std::env::var("KHIVE_ACTOR").ok();
3661        // SAFETY: test is #[serial]; no other test in this crate reads/writes KHIVE_ACTOR.
3662        unsafe {
3663            std::env::set_var("KHIVE_ACTOR", "lambda:test-env-actor");
3664        }
3665        let base = RuntimeConfig::default();
3666        assert_eq!(base.actor_id.as_deref(), Some("lambda:test-env-actor"));
3667
3668        let cfg = khive_cfg_with_actor("lambda:toml-actor");
3669        let result = runtime_config_from_khive_config(&cfg, base);
3670        assert_eq!(
3671            result.actor_id.as_deref(),
3672            Some("lambda:toml-actor"),
3673            "TOML [actor] id must win over the env-resolved base.actor_id"
3674        );
3675
3676        // SAFETY: restores prior KHIVE_ACTOR value (test cleanup).
3677        unsafe {
3678            match prior {
3679                Some(v) => std::env::set_var("KHIVE_ACTOR", v),
3680                None => std::env::remove_var("KHIVE_ACTOR"),
3681            }
3682        }
3683    }
3684
3685    // ---- list_embedding_models tests ----
3686
3687    // ---- core_backend accessor tests ----
3688
3689    /// Create a migrated in-memory backend (for tests that need raw Arc<StorageBackend>).
3690    fn migrated_memory_backend() -> Arc<StorageBackend> {
3691        let backend = StorageBackend::memory().expect("memory backend");
3692        {
3693            let mut writer = backend.pool().try_writer().expect("writer");
3694            khive_db::run_migrations(writer.conn_mut()).expect("migrations");
3695        }
3696        Arc::new(backend)
3697    }
3698
3699    fn secondary_config() -> RuntimeConfig {
3700        RuntimeConfig {
3701            web: Default::default(),
3702            telemetry: Default::default(),
3703            mounts: Vec::new(),
3704            brain: Default::default(),
3705            git_write: Default::default(),
3706            exec: Default::default(),
3707            display_timezone: chrono_tz::Tz::UTC,
3708            events_split: None,
3709            db_path: None,
3710            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3711            default_namespace: Namespace::local(),
3712            embedding_model: None,
3713            additional_embedding_models: vec![],
3714            gate: Arc::new(AllowAllGate),
3715            packs: vec!["kg".to_string()],
3716            backend_id: BackendId::parse("lore").expect("valid backend id"),
3717            brain_profile: None,
3718            visible_namespaces: vec![],
3719            allowed_outbound_namespaces: vec![],
3720            actor_id: None,
3721        }
3722    }
3723
3724    #[test]
3725    fn core_on_main_runtime_returns_same_backend_id() {
3726        // For a main-bound runtime, core() must return a clone with backend_id == "main".
3727        let rt = KhiveRuntime::memory().unwrap();
3728        assert_eq!(rt.backend_id().as_str(), BackendId::MAIN);
3729        let core_rt = rt.core();
3730        assert_eq!(core_rt.backend_id().as_str(), BackendId::MAIN);
3731    }
3732
3733    #[tokio::test]
3734    async fn core_on_main_runtime_round_trips_note() {
3735        // core() on a main-bound runtime (core_backend = None) returns self.clone(),
3736        // so a note written through core() is readable through the original runtime.
3737        let rt = KhiveRuntime::memory().unwrap();
3738        let tok = NamespaceToken::local();
3739
3740        let note = rt
3741            .core()
3742            .create_note(
3743                &tok,
3744                "observation",
3745                None,
3746                "adr073-main-round-trip",
3747                None,
3748                None,
3749                vec![],
3750            )
3751            .await
3752            .expect("create_note via core()");
3753
3754        let found = rt
3755            .notes(&tok)
3756            .expect("notes store")
3757            .get_note(note.id)
3758            .await
3759            .expect("get_note");
3760
3761        assert!(
3762            found.is_some(),
3763            "note written via core() must be visible through original rt"
3764        );
3765    }
3766
3767    /// Proves note→main and aux→secondary writes are each isolated.
3768    ///
3769    /// Backend A = main; backend B = secondary.
3770    /// rt_secondary is bound to B with core_backend = Some(A).
3771    ///
3772    /// Direction 1 (note → main):
3773    ///   rt_secondary.core().create_note(...) must land in A (visible from rt_main)
3774    ///   and NOT in B (not visible from rt_secondary).
3775    ///
3776    /// Direction 2 (aux → secondary):
3777    ///   A raw SQL write via rt_secondary.sql() lands in B only; A is untouched.
3778    #[tokio::test]
3779    async fn cross_backend_split_note_to_main_aux_to_secondary() {
3780        use khive_storage::{SqlStatement, SqlValue};
3781
3782        let main_arc = migrated_memory_backend();
3783        let secondary_arc = migrated_memory_backend();
3784
3785        let main_config = RuntimeConfig {
3786            web: Default::default(),
3787            telemetry: Default::default(),
3788            mounts: Vec::new(),
3789            brain: Default::default(),
3790            git_write: Default::default(),
3791            display_timezone: chrono_tz::Tz::UTC,
3792            events_split: None,
3793            db_path: None,
3794            blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3795            default_namespace: Namespace::local(),
3796            embedding_model: None,
3797            additional_embedding_models: vec![],
3798            gate: Arc::new(AllowAllGate),
3799            packs: vec!["kg".to_string()],
3800            backend_id: BackendId::main(),
3801            brain_profile: None,
3802            visible_namespaces: vec![],
3803            allowed_outbound_namespaces: vec![],
3804            actor_id: None,
3805            exec: Default::default(),
3806        };
3807
3808        let rt_main = KhiveRuntime::from_backend(main_arc.clone(), main_config);
3809        let rt_secondary = KhiveRuntime::from_backend(secondary_arc, secondary_config())
3810            .with_core_backend(main_arc.clone());
3811
3812        let tok = NamespaceToken::local();
3813
3814        // ── Direction 1: note must land in A (main), not in B (secondary) ──
3815
3816        let note = rt_secondary
3817            .core()
3818            .create_note(
3819                &tok,
3820                "observation",
3821                None,
3822                "adr073-split-test",
3823                None,
3824                None,
3825                vec![],
3826            )
3827            .await
3828            .expect("create_note via core()");
3829        let note_id = note.id;
3830
3831        // Visible from main (A).
3832        let in_main = rt_main
3833            .notes(&tok)
3834            .expect("main notes store")
3835            .get_note(note_id)
3836            .await
3837            .expect("get_note from main");
3838        assert!(
3839            in_main.is_some(),
3840            "note written via core() must appear in main backend A"
3841        );
3842
3843        // Not visible from secondary (B).
3844        let in_secondary = rt_secondary
3845            .notes(&tok)
3846            .expect("secondary notes store")
3847            .get_note(note_id)
3848            .await
3849            .expect("get_note from secondary");
3850        assert!(
3851            in_secondary.is_none(),
3852            "note written to main via core() must NOT appear in secondary backend B"
3853        );
3854
3855        // ── Direction 2: aux write via rt_secondary.sql() lands in B, not A ──
3856
3857        {
3858            let mut writer = rt_secondary.sql().writer().await.expect("secondary writer");
3859            writer
3860                .execute(SqlStatement {
3861                    sql: "CREATE TABLE IF NOT EXISTS _test_adr073_aux \
3862                          (marker TEXT PRIMARY KEY)"
3863                        .into(),
3864                    params: vec![],
3865                    label: None,
3866                })
3867                .await
3868                .expect("create aux table in B");
3869            writer
3870                .execute(SqlStatement {
3871                    sql: "INSERT INTO _test_adr073_aux VALUES (?1)".into(),
3872                    params: vec![SqlValue::Text("b-side-sentinel".into())],
3873                    label: None,
3874                })
3875                .await
3876                .expect("insert into aux table in B");
3877        }
3878
3879        // Row is present in B.
3880        let mut reader_b = rt_secondary.sql().reader().await.expect("secondary reader");
3881        let rows_b = reader_b
3882            .query_all(SqlStatement {
3883                sql: "SELECT marker FROM _test_adr073_aux".into(),
3884                params: vec![],
3885                label: None,
3886            })
3887            .await
3888            .expect("select from B");
3889        assert_eq!(rows_b.len(), 1, "aux row must exist in B");
3890        match rows_b[0].get("marker") {
3891            Some(SqlValue::Text(s)) => {
3892                assert_eq!(s, "b-side-sentinel", "sentinel value must match")
3893            }
3894            other => panic!("expected Text('b-side-sentinel'), got {other:?}"),
3895        }
3896
3897        // Row is absent from A (table does not exist there).
3898        let mut reader_a = rt_main.sql().reader().await.expect("main reader");
3899        let result_a = reader_a
3900            .query_all(SqlStatement {
3901                sql: "SELECT marker FROM _test_adr073_aux".into(),
3902                params: vec![],
3903                label: None,
3904            })
3905            .await;
3906        // A does not have this table → must error or return no rows.
3907        match result_a {
3908            Err(e) => assert!(
3909                e.to_string().contains("no such table"),
3910                "expected 'no such table' error from A, got: {e}"
3911            ),
3912            Ok(rows) => assert!(
3913                rows.is_empty(),
3914                "aux table must not have rows in A, got {} rows",
3915                rows.len()
3916            ),
3917        }
3918    }
3919
3920    #[test]
3921    fn constructors_leave_core_backend_none_by_behavior() {
3922        // core() on any standard constructor returns a clone with same backend_id —
3923        // proof that core_backend = None (returns self.clone(), not a different backend).
3924        let rt_mem = KhiveRuntime::memory().unwrap();
3925        assert_eq!(rt_mem.core().backend_id().as_str(), BackendId::MAIN);
3926
3927        let backend = migrated_memory_backend();
3928        let rt_from = KhiveRuntime::from_backend(
3929            backend,
3930            RuntimeConfig {
3931                web: Default::default(),
3932                telemetry: Default::default(),
3933                mounts: Vec::new(),
3934                brain: Default::default(),
3935                git_write: Default::default(),
3936                display_timezone: chrono_tz::Tz::UTC,
3937                events_split: None,
3938                db_path: None,
3939                blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3940                default_namespace: Namespace::local(),
3941                embedding_model: None,
3942                additional_embedding_models: vec![],
3943                gate: Arc::new(AllowAllGate),
3944                packs: vec!["kg".to_string()],
3945                backend_id: BackendId::parse("lore").expect("valid backend id"),
3946                brain_profile: None,
3947                visible_namespaces: vec![],
3948                allowed_outbound_namespaces: vec![],
3949                actor_id: None,
3950                exec: Default::default(),
3951            },
3952        );
3953        // from_backend with backend_id="lore" and no core_backend: core() returns
3954        // self.clone() which has backend_id="lore" (not "main").
3955        assert_eq!(rt_from.core().backend_id().as_str(), "lore");
3956    }
3957
3958    #[test]
3959    fn with_core_backend_sets_core_then_core_returns_main_id() {
3960        // After wiring, core() must return a runtime with backend_id == "main".
3961        let main_arc = migrated_memory_backend();
3962        let secondary_arc = migrated_memory_backend();
3963
3964        let rt_secondary = KhiveRuntime::from_backend(secondary_arc, secondary_config())
3965            .with_core_backend(main_arc);
3966
3967        assert_eq!(rt_secondary.backend_id().as_str(), "lore");
3968        assert_eq!(
3969            rt_secondary.core().backend_id().as_str(),
3970            BackendId::MAIN,
3971            "core() on a secondary runtime must return a main-bound handle"
3972        );
3973    }
3974
3975    #[test]
3976    fn attachment_store_rejects_secondary_handle_and_accepts_its_core_projection() {
3977        let main_arc = migrated_memory_backend();
3978        let secondary_arc = migrated_memory_backend();
3979        let rt_secondary = KhiveRuntime::from_backend(secondary_arc, secondary_config())
3980            .with_core_backend(main_arc);
3981
3982        let error = match rt_secondary.attachments() {
3983            Ok(_) => panic!("a secondary runtime must not expose attachment mutation"),
3984            Err(error) => error,
3985        };
3986        assert!(matches!(error, RuntimeError::InvalidInput(_)));
3987        assert!(
3988            error.to_string().contains("canonical main backend"),
3989            "secondary refusal must explain the liveness authority: {error}"
3990        );
3991        rt_secondary
3992            .core()
3993            .attachments()
3994            .expect("core projection must expose the main attachment store");
3995    }
3996
3997    #[tokio::test]
3998    async fn record_plus_attachment_publication_rejects_a_secondary_runtime() {
3999        use khive_storage::{BlobStore as _, NewAttachment};
4000
4001        let main_arc = migrated_memory_backend();
4002        let secondary_arc = migrated_memory_backend();
4003        let rt_secondary =
4004            KhiveRuntime::from_backend(Arc::clone(&secondary_arc), secondary_config())
4005                .with_core_backend(Arc::clone(&main_arc));
4006        let blob_root = tempfile::tempdir().expect("blob root");
4007        let blob_store = Arc::new(
4008            khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
4009                .expect("blob store"),
4010        );
4011        let content_ref = blob_store.put(b"secondary-ref".to_vec()).await.unwrap();
4012        rt_secondary
4013            .install_blob_store(blob_store.clone())
4014            .expect("shared blob store");
4015        let token = rt_secondary.authorize(Namespace::local()).unwrap();
4016
4017        let error = rt_secondary
4018            .create_entity_with_attachments(
4019                &token,
4020                "artifact",
4021                Some("visual_asset"),
4022                "must route through core",
4023                None,
4024                None,
4025                vec![],
4026                vec![NewAttachment {
4027                    role: "content".to_string(),
4028                    content_ref: content_ref.clone(),
4029                    media_type: None,
4030                    size_bytes: Some(13),
4031                }],
4032            )
4033            .await
4034            .expect_err("secondary attachment publication must fail closed");
4035        assert!(error.to_string().contains("canonical main backend"));
4036        assert!(rt_secondary
4037            .list_entities(&token, None, None, 10, 0)
4038            .await
4039            .unwrap()
4040            .is_empty());
4041        assert!(rt_secondary
4042            .core()
4043            .list_entities(&token, None, None, 10, 0)
4044            .await
4045            .unwrap()
4046            .is_empty());
4047        assert!(
4048            blob_store.exists(&content_ref).await.unwrap(),
4049            "refusal must not mutate the already-published object"
4050        );
4051    }
4052
4053    #[tokio::test]
4054    async fn list_embedding_models_returns_empty_when_table_absent() {
4055        // A brand-new in-memory runtime has migrations applied, so _embedding_models
4056        // IS created. But with no rows inserted, the result must be empty.
4057        let rt = KhiveRuntime::memory().expect("memory runtime");
4058        let records = rt
4059            .list_embedding_models(None)
4060            .await
4061            .expect("list ok on empty table");
4062        assert!(records.is_empty());
4063    }
4064
4065    #[tokio::test]
4066    async fn list_embedding_models_returns_row_after_insert() {
4067        use khive_storage::{SqlStatement, SqlValue};
4068
4069        let rt = KhiveRuntime::memory().expect("memory runtime");
4070        let sql = rt.sql();
4071
4072        let now = 1_000_000i64;
4073        let id = uuid::Uuid::new_v4();
4074        let canonical_key = b"test_engine:test-model-v1:v1:384".to_vec();
4075
4076        let mut writer = sql.writer().await.expect("writer");
4077        writer
4078            .execute(SqlStatement {
4079                sql: "INSERT INTO _embedding_models \
4080                      (id, engine_name, model_id, key_version, dim, output_dim, status, \
4081                       activated_at, superseded_at, superseded_by, canonical_key, created_at) \
4082                      VALUES (?1, ?2, ?3, ?4, ?5, NULL, ?6, ?7, NULL, NULL, ?8, ?9)"
4083                    .into(),
4084                params: vec![
4085                    SqlValue::Blob(id.as_bytes().to_vec()),
4086                    SqlValue::Text("test_engine".into()),
4087                    SqlValue::Text("test-model-v1".into()),
4088                    SqlValue::Text("v1".into()),
4089                    SqlValue::Integer(384),
4090                    SqlValue::Text("active".into()),
4091                    SqlValue::Integer(now),
4092                    SqlValue::Blob(canonical_key),
4093                    SqlValue::Integer(now),
4094                ],
4095                label: None,
4096            })
4097            .await
4098            .expect("insert row");
4099        drop(writer);
4100
4101        let records = rt.list_embedding_models(None).await.expect("list ok");
4102        assert_eq!(records.len(), 1);
4103        assert_eq!(records[0].engine_name, "test_engine");
4104        assert_eq!(records[0].model_id, "test-model-v1");
4105        assert_eq!(records[0].key_version, "v1");
4106        assert_eq!(records[0].dimensions, 384);
4107        assert_eq!(records[0].status, "active");
4108
4109        // engine filter — match
4110        let filtered = rt
4111            .list_embedding_models(Some("test_engine"))
4112            .await
4113            .expect("filter ok");
4114        assert_eq!(filtered.len(), 1);
4115
4116        // engine filter — no match
4117        let no_match = rt
4118            .list_embedding_models(Some("other_engine"))
4119            .await
4120            .expect("no-match ok");
4121        assert!(no_match.is_empty());
4122    }
4123
4124    #[test]
4125    fn named_vector_identity_rejects_ambiguous_or_unsafe_values() {
4126        assert!(NamedVectorIdentity::new("", "model", 4).is_err());
4127        assert!(NamedVectorIdentity::new("bad-key", "model", 4).is_err());
4128        assert!(NamedVectorIdentity::new("valid_key", " model", 4).is_err());
4129        assert!(NamedVectorIdentity::new("valid_key", "model", 0).is_err());
4130        assert!(NamedVectorIdentity::new("valid_key", "model", 8193).is_err());
4131        assert!(NamedVectorIdentity::new("k".repeat(128), "m".repeat(512), 4).is_ok());
4132        assert!(NamedVectorIdentity::new("k".repeat(129), "model", 4).is_err());
4133        assert!(NamedVectorIdentity::new("valid_key", "m".repeat(513), 4).is_err());
4134        assert_eq!(
4135            NamedVectorIdentity::new("valid_key", "model", 4)
4136                .expect("valid identity")
4137                .dimensions(),
4138            4
4139        );
4140    }
4141
4142    #[tokio::test]
4143    async fn named_vector_store_rejects_dimension_or_model_key_rebinding() {
4144        let rt = KhiveRuntime::memory().expect("memory runtime");
4145        let token = rt.authorize(Namespace::local()).expect("authorize");
4146        let original = NamedVectorIdentity::new("visual_contract", "model-a", 4).unwrap();
4147        rt.vectors_for_named_identity(&token, &original)
4148            .await
4149            .expect("create named vector store");
4150        let registered = rt
4151            .list_embedding_models(Some("visual_contract"))
4152            .await
4153            .expect("list model registry");
4154        assert!(registered.iter().any(|record| {
4155            record.model_id == "model-a"
4156                && record.key_version == "visual_contract"
4157                && record.dimensions == 4
4158        }));
4159        let wrong_dimensions = NamedVectorIdentity::new("visual_contract", "model-a", 5).unwrap();
4160        let Err(dimension_error) = rt
4161            .vectors_for_named_identity(&token, &wrong_dimensions)
4162            .await
4163        else {
4164            panic!("same key cannot change dimensions");
4165        };
4166        assert!(dimension_error.to_string().contains("dimensions"));
4167
4168        let wrong_model = NamedVectorIdentity::new("visual_contract", "model-b", 4).unwrap();
4169        let Err(model_error) = rt.vectors_for_named_identity(&token, &wrong_model).await else {
4170            panic!("same key cannot change model identity");
4171        };
4172        assert!(model_error.to_string().contains("already bound"));
4173    }
4174
4175    #[tokio::test]
4176    async fn repeated_named_vector_lookup_avoids_writer_acquisition() {
4177        let runtime = KhiveRuntime::memory().expect("memory runtime");
4178        let token = runtime.authorize(Namespace::local()).expect("authorize");
4179        let identity = NamedVectorIdentity::new("visual_cached", "model-a", 4).unwrap();
4180        runtime
4181            .vectors_for_named_identity(&token, &identity)
4182            .await
4183            .expect("first lookup validates and registers");
4184        let writer_before = runtime.backend().pool().writer_acquisition_snapshot();
4185
4186        runtime
4187            .clone()
4188            .vectors_for_named_identity(&token, &identity)
4189            .await
4190            .expect("clone reuses verified store");
4191        assert_eq!(
4192            runtime.backend().pool().writer_acquisition_snapshot(),
4193            writer_before,
4194            "repeated reads must not reach vector-table setup or model registration"
4195        );
4196    }
4197
4198    #[tokio::test]
4199    async fn core_projection_reuses_main_named_vector_cache() {
4200        let main_backend = migrated_memory_backend();
4201        let main =
4202            KhiveRuntime::from_backend(Arc::clone(&main_backend), RuntimeConfig::no_embeddings());
4203        let secondary = KhiveRuntime::from_backend(migrated_memory_backend(), secondary_config())
4204            .with_core_embedders_from(&main)
4205            .with_core_backend(Arc::clone(&main_backend));
4206        let core = secondary.core();
4207        let token = core.authorize(Namespace::local()).expect("authorize");
4208        let identity = NamedVectorIdentity::new("core_visual_cached", "model-a", 4).unwrap();
4209        core.vectors_for_named_identity(&token, &identity)
4210            .await
4211            .expect("first lookup validates on main");
4212        let writer_before = main_backend.pool().writer_acquisition_snapshot();
4213
4214        secondary
4215            .core()
4216            .vectors_for_named_identity(&token, &identity)
4217            .await
4218            .expect("new core projection reuses main store");
4219        assert_eq!(
4220            main_backend.pool().writer_acquisition_snapshot(),
4221            writer_before
4222        );
4223    }
4224
4225    #[tokio::test]
4226    async fn concurrent_named_vector_first_bind_has_one_immutable_winner() {
4227        let rt = KhiveRuntime::memory().expect("memory runtime");
4228        let token = rt.authorize(Namespace::local()).expect("authorize");
4229        let first = NamedVectorIdentity::new("visual_race", "model-a", 4).unwrap();
4230        let second = NamedVectorIdentity::new("visual_race", "model-b", 4).unwrap();
4231
4232        let (first_result, second_result) = tokio::join!(
4233            rt.vectors_for_named_identity(&token, &first),
4234            rt.vectors_for_named_identity(&token, &second),
4235        );
4236        assert_ne!(
4237            first_result.is_ok(),
4238            second_result.is_ok(),
4239            "the active engine_name uniqueness rule must select exactly one first binding"
4240        );
4241
4242        let (winner, loser) = if first_result.is_ok() {
4243            (&first, &second)
4244        } else {
4245            (&second, &first)
4246        };
4247        rt.vectors_for_named_identity(&token, winner)
4248            .await
4249            .expect("winning identity remains idempotent");
4250        let error = match rt.vectors_for_named_identity(&token, loser).await {
4251            Ok(_) => panic!("losing identity cannot rebind the empty table"),
4252            Err(error) => error,
4253        };
4254        assert!(error.to_string().contains("already bound"));
4255
4256        let registered = rt
4257            .list_embedding_models(Some("visual_race"))
4258            .await
4259            .expect("list race registry");
4260        assert_eq!(registered.len(), 1);
4261        assert_eq!(registered[0].model_id, winner.model_name());
4262    }
4263
4264    #[tokio::test]
4265    async fn named_vector_registry_keeps_immutable_revisions_active_together() {
4266        let rt = KhiveRuntime::memory().expect("memory runtime");
4267        let token = rt.authorize(Namespace::local()).expect("authorize");
4268        let first = NamedVectorIdentity::new("visual_revision_a", "visual-model", 4).unwrap();
4269        let second = NamedVectorIdentity::new("visual_revision_b", "visual-model", 4).unwrap();
4270
4271        rt.vectors_for_named_identity(&token, &first)
4272            .await
4273            .expect("open first immutable space");
4274        rt.vectors_for_named_identity(&token, &second)
4275            .await
4276            .expect("open second immutable space");
4277
4278        let registered = rt.list_embedding_models(None).await.expect("list registry");
4279        assert!(registered.iter().any(|record| {
4280            record.engine_name == "visual_revision_a"
4281                && record.model_id == "visual-model"
4282                && record.key_version == "visual_revision_a"
4283                && record.status == "active"
4284        }));
4285        assert!(registered.iter().any(|record| {
4286            record.engine_name == "visual_revision_b"
4287                && record.model_id == "visual-model"
4288                && record.key_version == "visual_revision_b"
4289                && record.status == "active"
4290        }));
4291    }
4292}