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::path::PathBuf;
9use std::sync::{Arc, Mutex, OnceLock, RwLock, Weak};
10
11use khive_db::{ConnectionPool, StorageBackend};
12#[cfg(test)]
13use khive_gate::AllowAllGate;
14use khive_gate::GateRequest;
15use khive_storage::types::{SqlStatement, SqlValue};
16use khive_storage::{
17    AttachmentStore, EntityStore, Event, EventStore, GraphStore, NoteStore, SqlAccess, VectorStore,
18};
19use khive_types::{EdgeEndpointRule, EventKind, Namespace, SubstrateKind};
20use lattice_embed::{EmbeddingModel, EmbeddingService};
21
22use crate::config::{
23    build_embedder_registry, parse_embedding_model_alias, register_configured_embedding_models,
24    sanitize_key, vec_model_key,
25};
26use crate::error::{RuntimeError, RuntimeResult};
27use crate::note_search_ann::NoteSearchAnnProvider;
28use crate::pack::KindHook;
29
30#[path = "runtime/config_access.rs"]
31mod config_access;
32mod embedder_init;
33mod events_disk_policy;
34mod serving_policy;
35
36#[cfg(all(test, target_os = "macos"))]
37const IN_PROCESS_TEST_NOFILE_LIMIT: libc::rlim_t = 4096;
38#[cfg(all(test, target_os = "macos"))]
39static IN_PROCESS_TEST_NOFILE_INIT: std::sync::Once = std::sync::Once::new();
40
41#[cfg(all(test, target_os = "macos"))]
42fn ensure_in_process_test_nofile_limit() {
43    IN_PROCESS_TEST_NOFILE_INIT.call_once(|| {
44        let mut limits = libc::rlimit {
45            rlim_cur: 0,
46            rlim_max: 0,
47        };
48        // SAFETY: `limits` is writable, and only this test binary's soft
49        // limit may change; the inherited hard limit is preserved.
50        assert_eq!(unsafe { libc::getrlimit(libc::RLIMIT_NOFILE, &mut limits) }, 0);
51        assert!(
52            limits.rlim_max >= IN_PROCESS_TEST_NOFILE_LIMIT,
53            "in-process SQLite tests require a hard open-file limit of at least {IN_PROCESS_TEST_NOFILE_LIMIT}"
54        );
55        if limits.rlim_cur < IN_PROCESS_TEST_NOFILE_LIMIT {
56            limits.rlim_cur = IN_PROCESS_TEST_NOFILE_LIMIT;
57            // SAFETY: the new soft limit does not exceed the observed hard
58            // limit, which is left unchanged.
59            assert_eq!(unsafe { libc::setrlimit(libc::RLIMIT_NOFILE, &limits) }, 0);
60        }
61    });
62}
63
64tokio::task_local! {
65    static REQUEST_EMBEDDER_EXCLUSIONS: Arc<HashSet<String>>;
66}
67
68/// Run one request with daemon-only embedding models excluded from registry access.
69pub fn scope_request_embedder_exclusions<F>(
70    excluded: Vec<String>,
71    future: F,
72) -> impl Future<Output = F::Output>
73where
74    F: Future,
75{
76    REQUEST_EMBEDDER_EXCLUSIONS.scope(Arc::new(excluded.into_iter().collect()), future)
77}
78
79/// Carry the current request's embedder exclusions into a spawned task.
80pub fn inherit_request_embedder_scope<F>(future: F) -> impl Future<Output = F::Output>
81where
82    F: Future,
83{
84    let exclusions = REQUEST_EMBEDDER_EXCLUSIONS.try_with(Arc::clone).ok();
85    async move {
86        match exclusions {
87            Some(exclusions) => REQUEST_EMBEDDER_EXCLUSIONS.scope(exclusions, future).await,
88            None => future.await,
89        }
90    }
91}
92
93fn request_excludes_embedder(name: &str) -> bool {
94    let canonical = parse_embedding_model_alias(name)
95        .map(|model| model.to_string())
96        .unwrap_or_else(|| name.to_string());
97    REQUEST_EMBEDDER_EXCLUSIONS
98        .try_with(|excluded| excluded.contains(&canonical))
99        .unwrap_or(false)
100}
101
102/// Callback type for pack-installed entity-type validators.
103///
104/// Receives `(kind, entity_type)` and returns the normalised type string,
105/// or `RuntimeError::InvalidInput` if the type is not registered for that kind.
106/// When `entity_type` is `None`, the implementation must return `Ok(None)`.
107pub type EntityTypeValidatorFn =
108    Arc<dyn Fn(&str, Option<&str>) -> Result<Option<String>, RuntimeError> + Send + Sync>;
109
110/// Pack-aggregated entity-kind update hooks: `(entity kind, hook)` for every
111/// kind whose owning pack both declares it and registers a `KindHook`.
112///
113/// Named rather than written inline because the runtime stores it behind an
114/// `Arc<RwLock<..>>` and passes it across the transport boundary, so the bare
115/// form appears three times and reads as noise at each one.
116pub type EntityKindHooks = Vec<(String, Arc<dyn KindHook>)>;
117
118/// Callback type for a pack-installed note-mutation hook.
119///
120/// Invoked by `update_note` (when the note's text/embedding actually
121/// changed) and `delete_note` (soft or hard) with `(note_kind, note_id)`,
122/// after the mutation has been durably applied. Returns a boxed future so
123/// the hook can await async cache-invalidation work (e.g.
124/// `khive-pack-memory`'s ANN warm-cache generation bump) without
125/// `khive-runtime` depending on any pack crate: dependencies point the
126/// other way, so the runtime exposes an extension point and the pack
127/// installs into it, same shape as `EntityTypeValidatorFn`, just async.
128pub type NoteMutationHookFn = Arc<
129    dyn Fn(String, uuid::Uuid) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>
130        + Send
131        + Sync,
132>;
133
134/// Callback type for a pack-installed note-write validator.
135///
136/// The pack that owns a note kind carrying derivable identity installs one so
137/// that the identity is a function of the authorization token rather than of
138/// caller input, on every write path including direct callers that bypass the
139/// handler layer — same rationale as [`EntityTypeValidatorFn`], which exists
140/// for exactly that reason on the entity side.
141///
142/// Kinds the installing pack does not own must be returned unchanged: the slot
143/// is single-occupancy (like `note_mutation_hook`), so a validator that
144/// rewrote foreign kinds would silently govern every other pack's notes.
145pub type NoteWriteValidatorFn = Arc<
146    dyn Fn(&str, &str, Option<serde_json::Value>) -> Result<Option<serde_json::Value>, RuntimeError>
147        + Send
148        + Sync,
149>;
150
151/// Immutable identity for a non-text vector store owned by a pack consumer.
152///
153/// This does not register an [`crate::EmbedderProvider`]. It gives a pack that
154/// performs its own governed inference a narrow path to a namespace-scoped
155/// Khive vector table while keeping model-key and dimension validation at the
156/// runtime boundary.
157#[derive(Clone, Debug, Eq, PartialEq)]
158pub struct NamedVectorIdentity {
159    model_key: String,
160    model_name: String,
161    dimensions: usize,
162}
163
164struct CachedNamedVectorStores {
165    identity: NamedVectorIdentity,
166    by_namespace: HashMap<String, Arc<dyn VectorStore>>,
167}
168
169fn check_cached_named_vector_identity(
170    cached: &NamedVectorIdentity,
171    requested: &NamedVectorIdentity,
172) -> RuntimeResult<()> {
173    if cached.dimensions() != requested.dimensions() {
174        return Err(RuntimeError::InvalidInput(format!(
175            "named vector model_key {:?} is already bound to {} dimensions, expected {}",
176            requested.model_key(),
177            cached.dimensions(),
178            requested.dimensions()
179        )));
180    }
181    if cached.model_name() != requested.model_name() {
182        return Err(RuntimeError::InvalidInput(format!(
183            "named vector model_key {:?} is already bound to a different active model identity",
184            requested.model_key()
185        )));
186    }
187    Ok(())
188}
189
190impl NamedVectorIdentity {
191    const MAX_MODEL_KEY_BYTES: usize = 128;
192    const MAX_MODEL_NAME_BYTES: usize = 512;
193
194    /// Validate and construct a named vector identity.
195    pub fn new(
196        model_key: impl Into<String>,
197        model_name: impl Into<String>,
198        dimensions: usize,
199    ) -> RuntimeResult<Self> {
200        let model_key = model_key.into();
201        let model_name = model_name.into();
202        if model_key.is_empty()
203            || model_key.len() > Self::MAX_MODEL_KEY_BYTES
204            || !model_key
205                .chars()
206                .all(|c| c.is_ascii_alphanumeric() || c == '_')
207        {
208            return Err(RuntimeError::InvalidInput(format!(
209                "named vector model_key must be 1..={} bytes of ASCII alphanumeric/underscore",
210                Self::MAX_MODEL_KEY_BYTES
211            )));
212        }
213        if model_name.trim().is_empty()
214            || model_name.trim() != model_name
215            || model_name.len() > Self::MAX_MODEL_NAME_BYTES
216        {
217            return Err(RuntimeError::InvalidInput(format!(
218                "named vector model_name must be 1..={} bytes with no surrounding whitespace",
219                Self::MAX_MODEL_NAME_BYTES
220            )));
221        }
222        if !(1..=8192).contains(&dimensions) {
223            return Err(RuntimeError::InvalidInput(format!(
224                "named vector dimensions must be in 1..=8192, got {dimensions}"
225            )));
226        }
227        Ok(Self {
228            model_key,
229            model_name,
230            dimensions,
231        })
232    }
233
234    pub fn model_key(&self) -> &str {
235        &self.model_key
236    }
237
238    pub fn model_name(&self) -> &str {
239        &self.model_name
240    }
241
242    pub fn dimensions(&self) -> usize {
243        self.dimensions
244    }
245}
246
247pub use crate::config::{
248    assert_captured_db_anchor_consistent, assert_db_anchor_consistent, expand_tilde,
249    parse_pack_list, resolve_db_anchor, resolve_project_actor_id, runtime_config_from_khive_config,
250    BackendId, BackendIdError, NamespaceToken, RuntimeConfig,
251};
252
253// ---- KhiveRuntime ----
254
255/// Composable runtime handle used by the MCP server.
256///
257/// Wraps a `StorageBackend` and provides namespace-scoped accessor methods
258/// Snapshot of the main runtime's embedder wiring (registry handle, default
259/// name, and the config model fields the default-resolution path reads).
260/// Carried by secondary-pack runtimes; consumed by [`KhiveRuntime::core`].
261#[derive(Clone)]
262struct CoreEmbedderState {
263    registry: Arc<std::sync::RwLock<crate::embedder_registry::EmbedderRegistry>>,
264    default_embedder_name: Arc<str>,
265    embedding_model: Option<EmbeddingModel>,
266    additional_embedding_models: Vec<EmbeddingModel>,
267}
268
269#[derive(Clone)]
270struct NoteKindEntry {
271    name: String,
272    embedding_policy: crate::NoteEmbeddingPolicy,
273    registered: bool,
274}
275
276/// An already-open serving backend eligible for operator diagnostics.
277/// Aliases of the same canonical database file share one entry.
278#[derive(Clone)]
279pub struct OpenedDiagnosticBackend {
280    pub backend_names: Vec<String>,
281    pub canonical_path: Option<PathBuf>,
282    pub pool: Arc<ConnectionPool>,
283}
284
285struct LateOpenedDiagnosticBackend {
286    backend_names: Vec<&'static str>,
287    pool: Weak<ConnectionPool>,
288}
289
290fn same_diagnostic_database(a: &OpenedDiagnosticBackend, b: &OpenedDiagnosticBackend) -> bool {
291    match (a.pool.canonical_path(), b.pool.canonical_path()) {
292        (Some(a), Some(b)) => a == b,
293        (None, None) => Arc::ptr_eq(&a.pool, &b.pool),
294        _ => false,
295    }
296}
297
298/// for each storage capability, plus a lazily-loaded embedder.
299#[derive(Clone)]
300pub struct KhiveRuntime {
301    pub(crate) visibility_receipts: Arc<crate::visibility_receipts::ReceiptCapability>,
302    pub(crate) visibility_cutover: Arc<crate::visibility_receipts::ReceiptCutover>,
303    core_visibility_cutover: Arc<crate::visibility_receipts::ReceiptCutover>,
304    backend: Arc<StorageBackend>,
305    /// Successful named-vector bindings and their namespace-scoped stores.
306    /// Shared by runtime clones so repeated reads do not enter the writer or
307    /// rescan the vector table after the first verified binding.
308    named_vector_stores: Arc<RwLock<HashMap<String, CachedNamedVectorStores>>>,
309    /// The main backend's cache, used when a secondary runtime creates a
310    /// `core()` handle. It must never reuse a secondary backend's store.
311    core_named_vector_stores: Option<Arc<RwLock<HashMap<String, CachedNamedVectorStores>>>>,
312    /// When `Some`, holds the main backend so that `core()` can return a
313    /// main-bound runtime handle without constructing a new connection.
314    /// `None` when this runtime is already bound to the main backend.
315    core_backend: Option<Arc<StorageBackend>>,
316    config: RuntimeConfig,
317    outbound_email_policy: crate::OutboundEmailPolicy,
318    /// All SQLite backends declared by the host process, including those
319    /// assigned to other packs. The code pack fences these from ingest.
320    declared_backend_db_paths: Arc<[PathBuf]>,
321    /// Pools opened during serving host composition, grouped by canonical
322    /// database file.
323    diagnostic_backends: Arc<[OpenedDiagnosticBackend]>,
324    /// Pools opened later by this serving runtime and its pack handles.
325    /// Weak references keep diagnostics from extending their lifetime.
326    late_diagnostic_backends: Arc<Mutex<Vec<LateOpenedDiagnosticBackend>>>,
327    /// ADR-118 exact-leg policy, sampled once at runtime construction.
328    /// Request-time memory/knowledge serving must never re-read the process
329    /// environment because tests and embedded runtimes share one process.
330    ann_fresh_tail_enabled: bool,
331    /// Pack-extensible embedder registry.
332    ///
333    /// Shared across clones via `Arc<RwLock<_>>` so that
334    /// [`register_embedder`](Self::register_embedder) after clone is visible
335    /// to all handles. Built-in lattice models are pre-registered during
336    /// construction; packs may add more via [`PackRuntime::register_embedders`].
337    embedder_registry: Arc<std::sync::RwLock<crate::embedder_registry::EmbedderRegistry>>,
338    default_embedder_name: Arc<str>,
339    /// The MAIN runtime's embedder wiring, carried by secondary-backend
340    /// runtimes so that `core()`-routed writes embed with the main runtime's
341    /// models even when this pack's own registry is empty (`no_embed`).
342    /// `None` on the main runtime, and on secondaries wired before the boot
343    /// path calls [`with_core_embedders_from`](Self::with_core_embedders_from)
344    /// — `core()` then falls back to this runtime's own embedder state.
345    core_embedders: Option<CoreEmbedderState>,
346    /// Pack-extensible edge endpoint rules. Shared across clones
347    /// via `Arc<RwLock<_>>`; installed once by the transport after the
348    /// `VerbRegistry` is built. Empty until installed
349    edge_rules: Arc<RwLock<Vec<EdgeEndpointRule>>>,
350    /// Pack-aggregated valid entity kinds and note-kind policy entries.
351    ///
352    /// Installed by the transport layer after building the `VerbRegistry`.
353    /// When non-empty, `create_entity`, `create_note_inner`, and `import_kg`
354    /// reject kinds not in these sets. When empty (no packs loaded, e.g.
355    /// bare runtime in unit tests), kind validation is skipped — the pack
356    /// handler layer is the primary enforcement point.
357    valid_entity_kinds: Arc<RwLock<Vec<String>>>,
358    valid_note_kinds: Arc<RwLock<Vec<NoteKindEntry>>>,
359    /// Pack-installed entity-type validator.
360    ///
361    /// When `Some`, `create_many` calls this function to validate and normalise
362    /// each `(kind, entity_type)` pair before writing. When `None` (bare runtime
363    /// without packs), entity-type validation is skipped — the pack handler layer
364    /// is the primary enforcement point, same as for `valid_entity_kinds`.
365    entity_type_validator: Arc<RwLock<Option<EntityTypeValidatorFn>>>,
366    /// Pack-installed note-mutation hook.
367    ///
368    /// When `Some`, `update_note` (on text change) and `delete_note` (soft
369    /// or hard) call this after the mutation is durably applied, so a pack
370    /// that caches derived state keyed by note content (e.g. `khive-pack-memory`'s
371    /// warm ANN index) can invalidate/advance its own generation counter even
372    /// when the mutation arrived through a different pack's verb (e.g. KG's
373    /// `update`/`delete` on a `kind="memory"` note) that has no dependency on
374    /// the reacting pack. `None` when no pack installs one (bare runtime, or
375    /// no pack cares about note-mutation notifications) — the call becomes a
376    /// no-op check of an `Option`.
377    note_mutation_hook: Arc<RwLock<Option<NoteMutationHookFn>>>,
378    /// Backend-matched, pack-installed ANN source for note search. Absent on a
379    /// bare runtime or when no memory graph provider serves this backend.
380    note_search_ann_provider: Arc<RwLock<Option<Arc<dyn NoteSearchAnnProvider>>>>,
381    /// Pack-installed note-write validator.
382    ///
383    /// When `Some`, every runtime note-materialisation site that accepts
384    /// caller-supplied `properties` routes them through this function before
385    /// the `Note` is built, so a pack-owned identity property is derived from
386    /// the authorization token instead of trusted from caller input. `None`
387    /// on a bare runtime (no packs) — the properties pass through unchanged.
388    note_write_validator: Arc<RwLock<Option<NoteWriteValidatorFn>>>,
389    /// Pack-installed entity-kind update-validation hooks (issue #2943).
390    ///
391    /// Every `(entity kind, hook)` pair for which an owning pack declares
392    /// the entity kind and registers a `KindHook`, aggregated once by
393    /// `VerbRegistry::entity_kind_hooks` and installed by the transport
394    /// after the registry is built — same timing and rationale as
395    /// `entity_type_validator`: `khive-runtime` does not hold a
396    /// `VerbRegistry`, so this is the extension point that lets
397    /// `prepare_guarded_entity_update` reach a pack's `KindHook` on the
398    /// generic entity `update` path, the counterpart to
399    /// `prepare_note_update_hook` on the note side. Empty until installed
400    /// (bare runtime, or no pack registers an entity-kind hook), which
401    /// leaves the dispatch a no-op.
402    entity_kind_hooks: Arc<RwLock<EntityKindHooks>>,
403    /// Pack-owned note kinds — every note kind declared by a pack other than
404    /// the generic-CRUD pack, installed by the transport from the registry
405    /// (see `VerbRegistry::pack_owned_note_kinds`). Records of these kinds are
406    /// maintained by their owning pack's own verbs, so `update`'s `properties`
407    /// patch is refused on them at the runtime layer and their owned identity
408    /// properties survive a `merge` unchanged. Empty until installed (bare
409    /// runtime), which leaves both rules inert.
410    pack_owned_note_kinds: Arc<RwLock<Vec<String>>>,
411    /// The immutable runtime-owned store/hydration pair (ADR-160 D3).
412    ///
413    /// Boot resolves one store and constructs exactly one shared hydrator,
414    /// then installs that same `Arc` on every runtime handle. The one-shot
415    /// slot rejects replacement so pack runtimes cannot silently split the
416    /// aggregate byte budget after startup. Bare runtimes leave it unset.
417    blob_hydrator: Arc<OnceLock<Arc<crate::blob::BlobHydrator>>>,
418    /// Pack-registered custom fusion executors (ADR-012), keyed by the name
419    /// carried in `FusionStrategy::Custom { name, .. }`.
420    ///
421    /// Unlike `entity_type_validator`/`note_mutation_hook` (single-occupancy —
422    /// one pack owns the slot), multiple packs each register their own named
423    /// strategy under this shared map, so it is keyed rather than a bare
424    /// `Option`. Empty until a pack calls
425    /// [`register_fusion_strategy`](Self::register_fusion_strategy); an
426    /// unregistered `Custom` name at dispatch time is
427    /// `RuntimeError::UnknownFusionStrategy`, never a silent fallback.
428    fusion_executors: Arc<RwLock<HashMap<String, Arc<dyn crate::fusion::FusionExecutor>>>>,
429}
430
431impl KhiveRuntime {
432    /// Create a new runtime with the given config.
433    ///
434    /// The config's `db_path` is used to open or create the SQLite backend.
435    /// This direct constructor is intended for fresh/current single-backend
436    /// databases and tests. Production and multi-backend hosts must use the
437    /// async khive-mcp/kkernel builders so secondary inventory and any
438    /// application-assisted V21 cutover complete before serving. The
439    /// [`from_backend`](Self::from_backend) seam is likewise only for an
440    /// already-prepared backend.
441    pub fn new(mut config: RuntimeConfig) -> RuntimeResult<Self> {
442        let wal_ceiling = config.resolve_wal_ceiling_policy(false)?;
443        let disk_guard = config.resolve_disk_guard_policy(false)?;
444        // Refuse a missing lock directory before the constructor below creates
445        // the database's parent directory.
446        let volume_lock_dir = if config.db_path.is_some() {
447            let configured = config.volume_lock_dir.clone();
448            Some(khive_db::require_volume_lock_dir(configured)?)
449        } else {
450            None
451        };
452        Self::new_with_file_backend(config, true, |path| {
453            StorageBackend::sqlite_with_max_readers_and_policies(
454                path,
455                None,
456                wal_ceiling,
457                disk_guard.expect("file-backed disk policy"),
458                volume_lock_dir.expect("file-backed volume-lock directory"),
459            )
460        })
461    }
462
463    /// Construct a fixture runtime with a small concurrent reader pool.
464    #[cfg(any(test, feature = "test-internals"))]
465    pub fn new_for_test(mut config: RuntimeConfig) -> RuntimeResult<Self> {
466        let wal_ceiling = config.resolve_wal_ceiling_policy(false)?;
467        let disk_guard = config.resolve_disk_guard_policy(false)?;
468        Self::new_with_file_backend(config, true, |path| {
469            StorageBackend::sqlite_for_test_with_policies(
470                path,
471                wal_ceiling,
472                disk_guard.expect("file-backed disk policy"),
473            )
474        })
475    }
476
477    pub(crate) fn new_with_file_backend(
478        config: RuntimeConfig,
479        create_parent: bool,
480        open_file: impl FnOnce(&std::path::Path) -> Result<StorageBackend, khive_db::SqliteError>,
481    ) -> RuntimeResult<Self> {
482        #[cfg(unix)]
483        crate::events_split::socket_path::validate_configured_events_socket(&config)?;
484        #[cfg(all(test, target_os = "macos"))]
485        ensure_in_process_test_nofile_limit();
486        let backend = match &config.db_path {
487            Some(path) => {
488                if let Some(parent) = path.parent().filter(|_| create_parent) {
489                    std::fs::create_dir_all(parent).ok();
490                }
491                open_file(path)?
492            }
493            None => {
494                if config.wal_ceiling_configured_bytes != 0 || config.wal_ceiling_bytes != 0 {
495                    return Err(khive_db::SqliteError::InvalidConfig(
496                        "nonzero wal_ceiling_bytes requires a file-backed SQLite backend".into(),
497                    )
498                    .into());
499                }
500                StorageBackend::memory()?
501            }
502        };
503        // Writable backends migrate before handlers touch the DB. A detected
504        // read-only snapshot is validated at the current schema version without
505        // attempting migration DDL.
506        let schema_version = backend.prepare_core_schema()?;
507        if schema_version < khive_db::migrations::ATTACHMENT_CUTOVER_VERSION {
508            return Err(khive_db::SqliteError::InvalidData(
509                "database requires the host application-assisted V21 attachment cutover; \
510                 start through khive-mcp/kkernel boot instead of constructing KhiveRuntime \
511                 directly"
512                    .into(),
513            )
514            .into());
515        }
516        if !backend.is_read_only() {
517            register_configured_embedding_models(&backend, &config)?;
518        }
519        Ok(Self::assemble_from_backend(
520            Arc::new(backend),
521            config,
522            false,
523        ))
524    }
525
526    /// Open a runtime for read-only inspection (no model registration, no DB creation).
527    ///
528    /// File-backed databases are opened with SQLite read-only/query-only flags
529    /// and must already be at this build's current schema version. No migrations
530    /// or configured-model registration writes are attempted. A `None` path
531    /// retains the historical ephemeral in-memory behavior for tests.
532    pub fn new_readonly(mut config: RuntimeConfig) -> RuntimeResult<Self> {
533        let wal_ceiling = config.resolve_wal_ceiling_policy(true)?;
534        Self::new_readonly_with_file_backend(config, |path| {
535            StorageBackend::sqlite_read_only_with_max_readers_and_wal_ceiling(
536                path,
537                None,
538                wal_ceiling,
539            )
540        })
541    }
542
543    /// Construct a read-only fixture runtime with a small reader pool.
544    #[cfg(any(test, feature = "test-internals"))]
545    pub fn new_readonly_for_test(mut config: RuntimeConfig) -> RuntimeResult<Self> {
546        let wal_ceiling = config.resolve_wal_ceiling_policy(true)?;
547        Self::new_readonly_with_file_backend(config, |path| {
548            StorageBackend::sqlite_read_only_with_max_readers_and_wal_ceiling(
549                path,
550                Some(2),
551                wal_ceiling,
552            )
553        })
554    }
555
556    fn new_readonly_with_file_backend(
557        config: RuntimeConfig,
558        open_file: impl FnOnce(&std::path::Path) -> Result<StorageBackend, khive_db::SqliteError>,
559    ) -> RuntimeResult<Self> {
560        #[cfg(unix)]
561        crate::events_split::socket_path::validate_configured_events_socket(&config)?;
562        #[cfg(all(test, target_os = "macos"))]
563        ensure_in_process_test_nofile_limit();
564        let backend = match &config.db_path {
565            Some(path) => open_file(path)?,
566            None => {
567                if config.wal_ceiling_configured_bytes != 0 || config.wal_ceiling_bytes != 0 {
568                    return Err(khive_db::SqliteError::InvalidConfig(
569                        "nonzero wal_ceiling_bytes requires a file-backed SQLite backend".into(),
570                    )
571                    .into());
572                }
573                StorageBackend::memory()?
574            }
575        };
576        backend.prepare_core_schema()?;
577        Ok(Self::assemble_from_backend(
578            Arc::new(backend),
579            config,
580            false,
581        ))
582    }
583
584    /// Construct a runtime from an already-opened backend.
585    ///
586    /// This is a low-level, infallible assembly seam for already-prepared
587    /// multi-backend deployments. It does not migrate or require completion of
588    /// the V21 attachment cutover, and it does not read the backend: receipt
589    /// admission is checked by the first receipt operation and cached once it
590    /// succeeds. Production hosts must first run the async kkernel/khive-mcp
591    /// coordinator and must not expose a server over a pending or incomplete
592    /// backend. Prefer [`Self::from_prepared_backend`] when constructing one
593    /// fallible host runtime.
594    ///
595    /// The returned runtime has `db_path = None` and `embedding_model = None`; all
596    /// storage access is through the provided `backend`. Set `backend_id` and
597    /// `default_namespace` via the config builder pattern if non-defaults are needed.
598    pub fn from_backend(backend: Arc<StorageBackend>, config: RuntimeConfig) -> Self {
599        if !backend.is_read_only() {
600            if let Err(err) = register_configured_embedding_models(&backend, &config) {
601                tracing::warn!(error = %err, "failed to register configured embedding models");
602            }
603        }
604        Self::assemble_from_backend(backend, config, false)
605    }
606
607    /// Construct a single-backend runtime after a host boot coordinator has
608    /// completed schema preparation and any application-assisted cutover.
609    ///
610    /// Unlike [`Self::from_backend`], configured embedding-model registration
611    /// is fallible here, preserving [`Self::new`]'s single-backend startup
612    /// semantics. This method never runs migrations itself.
613    pub fn from_prepared_backend(
614        backend: Arc<StorageBackend>,
615        config: RuntimeConfig,
616    ) -> RuntimeResult<Self> {
617        if backend.attachment_cutover_status()?
618            != khive_db::migrations::AttachmentCutoverStatus::Complete
619        {
620            return Err(khive_db::SqliteError::InvalidData(
621                "from_prepared_backend requires a complete V21 attachment cutover".into(),
622            )
623            .into());
624        }
625        backend.validate_memory_visibility_cutover()?;
626        if !backend.is_read_only() {
627            register_configured_embedding_models(&backend, &config)?;
628        }
629        Ok(Self::assemble_from_backend(backend, config, true))
630    }
631
632    fn assemble_from_backend(
633        backend: Arc<StorageBackend>,
634        config: RuntimeConfig,
635        cutover_validated: bool,
636    ) -> Self {
637        if config.backend_id.as_str() == BackendId::MAIN {
638            backend.pool().main_pool_generation();
639        }
640        let ann_fresh_tail_enabled = crate::config::ann_fresh_tail_enabled_from_env();
641        let (registry, default_embedder_name) = build_embedder_registry(&config);
642        let visibility_receipts = Arc::new(
643            crate::visibility_receipts::ReceiptCapability::from_config(&config),
644        );
645        let visibility_cutover = Arc::new(crate::visibility_receipts::ReceiptCutover::new(
646            backend.clone(),
647            cutover_validated,
648        ));
649        Self {
650            visibility_receipts,
651            core_visibility_cutover: visibility_cutover.clone(),
652            visibility_cutover,
653            backend,
654            named_vector_stores: Arc::new(RwLock::new(HashMap::new())),
655            core_named_vector_stores: None,
656            core_backend: None,
657            config,
658            outbound_email_policy: Default::default(),
659            declared_backend_db_paths: Vec::new().into(),
660            diagnostic_backends: Vec::new().into(),
661            late_diagnostic_backends: Arc::new(Mutex::new(Vec::new())),
662            ann_fresh_tail_enabled,
663            embedder_registry: Arc::new(std::sync::RwLock::new(registry)),
664            default_embedder_name,
665            core_embedders: None,
666            edge_rules: Arc::new(RwLock::new(Vec::new())),
667            valid_entity_kinds: Arc::new(RwLock::new(Vec::new())),
668            valid_note_kinds: Arc::new(RwLock::new(Vec::new())),
669            entity_type_validator: Arc::new(RwLock::new(None)),
670            note_mutation_hook: Arc::new(RwLock::new(None)),
671            note_search_ann_provider: Arc::new(RwLock::new(None)),
672            note_write_validator: Arc::new(RwLock::new(None)),
673            entity_kind_hooks: Arc::new(RwLock::new(Vec::new())),
674            pack_owned_note_kinds: Arc::new(RwLock::new(Vec::new())),
675            blob_hydrator: Arc::new(OnceLock::new()),
676            fusion_executors: Arc::new(RwLock::new(HashMap::new())),
677        }
678    }
679
680    /// Wire this runtime as a secondary-backend runtime pointing at `core`.
681    ///
682    /// After this call, `self.core()` returns a handle to `core` rather than
683    /// cloning `self`. The caller (the boot path, not pack code) is responsible
684    /// for passing the correct main backend.
685    /// Binding a different core clears its named-vector cache and prior main
686    /// embedder wiring. Call [`Self::with_core_embedders_from`] with the new main
687    /// runtime after rebinding when core-routed writes require its embedders.
688    ///
689    /// Panics in debug builds if `self.config.backend_id == BackendId::MAIN`,
690    /// because the main runtime does not need a core pointer.
691    pub fn with_core_backend(mut self, core: Arc<StorageBackend>) -> Self {
692        debug_assert_ne!(
693            self.config.backend_id.as_str(),
694            BackendId::MAIN,
695            "with_core_backend must not be called on the main runtime"
696        );
697        core.pool().main_pool_generation();
698        if self.visibility_cutover.is_bound_to(&core) {
699            self.core_visibility_cutover = self.visibility_cutover.clone();
700        } else if !self.core_visibility_cutover.is_bound_to(&core) {
701            self.core_visibility_cutover = Arc::new(
702                crate::visibility_receipts::ReceiptCutover::new(core.clone(), false),
703            );
704        }
705        if self
706            .core_backend
707            .as_ref()
708            .is_some_and(|previous| !Arc::ptr_eq(previous, &core))
709        {
710            self.core_named_vector_stores = None;
711            self.core_embedders = None;
712        }
713        if self.core_named_vector_stores.is_none() {
714            self.core_named_vector_stores = Some(Arc::new(RwLock::new(HashMap::new())));
715        }
716        self.core_backend = Some(core);
717        self
718    }
719
720    /// Carry the main runtime's embedder wiring for `core()`-routed writes.
721    ///
722    /// Boot-path companion to [`with_core_backend`](Self::with_core_backend).
723    /// Without it, `core()` shares this pack runtime's own embedder registry —
724    /// which under `[packs.<name>] no_embed = true` is empty, so core-routed
725    /// concept writes would silently skip embedding on the shared graph.
726    pub fn with_core_embedders_from(mut self, main: &KhiveRuntime) -> Self {
727        debug_assert!(
728            main.core_backend.is_none(),
729            "with_core_embedders_from takes the MAIN runtime"
730        );
731        self.core_embedders = Some(CoreEmbedderState {
732            registry: main.embedder_registry.clone(),
733            default_embedder_name: main.default_embedder_name.clone(),
734            embedding_model: main.config.embedding_model,
735            additional_embedding_models: main.config.additional_embedding_models.clone(),
736        });
737        self.core_named_vector_stores = Some(main.named_vector_stores.clone());
738        if Arc::ptr_eq(&self.backend, &main.backend) {
739            self.named_vector_stores = main.named_vector_stores.clone();
740        }
741        self
742    }
743
744    /// Return a runtime handle bound to the main (shared-graph) backend.
745    ///
746    /// When `self` is already the main runtime (`core_backend` is `None`),
747    /// this returns a clone of `self` — no new backend reference is acquired.
748    ///
749    /// When `self` is a secondary-backend runtime (`core_backend` is `Some`),
750    /// this returns a new `KhiveRuntime` backed by the main
751    /// `Arc<StorageBackend>` and sharing all registry state (`embedder_registry`,
752    /// `edge_rules`, `valid_entity_kinds`, `valid_note_kinds`,
753    /// `entity_type_validator`, `note_mutation_hook`, `entity_kind_hooks`) with `self`.
754    /// No database I/O occurs; no embedding models are reloaded.
755    ///
756    /// Use `core()` for notes and entities that must reside in the shared graph
757    /// so that `memory.recall`, cross-pack search, and `annotates` edges work.
758    /// Use `self` (or `self.sql()`) for pack-auxiliary bulk tables.
759    ///
760    /// Handlers that call `core()` more than once per request or loop should bind
761    /// `let core = self.core();` once and reuse it, since each call clones
762    /// `RuntimeConfig` (a heap-allocated struct containing `Vec<String>` fields).
763    pub fn core(&self) -> KhiveRuntime {
764        match &self.core_backend {
765            // A main-assigned pack runtime has no core pointer, but may still
766            // carry main's embedder wiring: with `no_embed` its OWN registry
767            // is empty, and core-routed concept writes must embed regardless
768            // of which backend the pack was assigned to.
769            None => match &self.core_embedders {
770                None => self.clone(),
771                Some(core_embedders) => {
772                    let mut core = self.clone();
773                    core.config.embedding_model = core_embedders.embedding_model;
774                    core.config.additional_embedding_models =
775                        core_embedders.additional_embedding_models.clone();
776                    core.embedder_registry = core_embedders.registry.clone();
777                    core.default_embedder_name = core_embedders.default_embedder_name.clone();
778                    core.core_embedders = None;
779                    core
780                }
781            },
782            Some(main_arc) => {
783                let mut core_config = self.config.clone();
784                core_config.backend_id = BackendId::main();
785                // Core-routed writes embed with the MAIN runtime's wiring when
786                // the boot path supplied it (see `with_core_embedders_from`);
787                // both the registry handle and the config model fields must
788                // come from main, since default-model resolution reads
789                // `config.embedding_model` (`resolve_embedding_model`).
790                let (embedder_registry, default_embedder_name) = match &self.core_embedders {
791                    Some(core_embedders) => {
792                        core_config.embedding_model = core_embedders.embedding_model;
793                        core_config.additional_embedding_models =
794                            core_embedders.additional_embedding_models.clone();
795                        (
796                            core_embedders.registry.clone(),
797                            core_embedders.default_embedder_name.clone(),
798                        )
799                    }
800                    None => (
801                        self.embedder_registry.clone(),
802                        self.default_embedder_name.clone(),
803                    ),
804                };
805                KhiveRuntime {
806                    visibility_receipts: self.visibility_receipts.clone(),
807                    visibility_cutover: self.core_visibility_cutover.clone(),
808                    core_visibility_cutover: self.core_visibility_cutover.clone(),
809                    backend: main_arc.clone(),
810                    named_vector_stores: self
811                        .core_named_vector_stores
812                        .clone()
813                        .unwrap_or_else(|| Arc::new(RwLock::new(HashMap::new()))),
814                    core_named_vector_stores: None,
815                    core_backend: None,
816                    config: core_config,
817                    outbound_email_policy: self.outbound_email_policy.clone(),
818                    declared_backend_db_paths: self.declared_backend_db_paths.clone(),
819                    diagnostic_backends: self.diagnostic_backends.clone(),
820                    late_diagnostic_backends: self.late_diagnostic_backends.clone(),
821                    ann_fresh_tail_enabled: self.ann_fresh_tail_enabled,
822                    embedder_registry,
823                    default_embedder_name,
824                    core_embedders: None,
825                    edge_rules: self.edge_rules.clone(),
826                    valid_entity_kinds: self.valid_entity_kinds.clone(),
827                    valid_note_kinds: self.valid_note_kinds.clone(),
828                    entity_type_validator: self.entity_type_validator.clone(),
829                    note_mutation_hook: self.note_mutation_hook.clone(),
830                    note_search_ann_provider: self.note_search_ann_provider.clone(),
831                    note_write_validator: self.note_write_validator.clone(),
832                    entity_kind_hooks: self.entity_kind_hooks.clone(),
833                    pack_owned_note_kinds: self.pack_owned_note_kinds.clone(),
834                    blob_hydrator: self.blob_hydrator.clone(),
835                    fusion_executors: self.fusion_executors.clone(),
836                }
837            }
838        }
839    }
840
841    /// Create an in-memory runtime (for tests and ephemeral use).
842    pub fn memory() -> RuntimeResult<Self> {
843        Self::new(RuntimeConfig {
844            db_path: None,
845            packs: vec!["kg".to_string()],
846            brain_profile: None,
847            actor_id: None,
848            ..RuntimeConfig::no_embeddings()
849        })
850    }
851
852    /// Return the [`BackendId`] for this runtime's backend.
853    ///
854    /// Used by `SubstrateCoordinator` in `kkernel`
855    /// to identify which backend owns a given node, and to detect cross-backend merges.
856    pub fn backend_id(&self) -> &BackendId {
857        &self.config.backend_id
858    }
859
860    /// Whether two runtime handles share one opened physical store. The host
861    /// deduplicates same-path aliases; separately opened hard-link aliases
862    /// compare the identity pinned when SQLite opened each file.
863    pub fn shares_backend_storage_with(&self, other: &Self) -> bool {
864        if Arc::ptr_eq(&self.backend, &other.backend) {
865            return true;
866        }
867        #[cfg(any(unix, windows))]
868        {
869            matches!(
870                (
871                    self.backend.pool().opened_file_identity_record(),
872                    other.backend.pool().opened_file_identity_record()
873                ),
874                (Some(left), Some(right)) if left == right
875            )
876        }
877        #[cfg(not(any(unix, windows)))]
878        {
879            false
880        }
881    }
882
883    /// Install only pools the host actually opened. A bare runtime defaults
884    /// to its already-open main pool without opening or creating another file.
885    pub fn with_diagnostic_backends(mut self, backends: Arc<[OpenedDiagnosticBackend]>) -> Self {
886        self.diagnostic_backends = backends;
887        self
888    }
889
890    /// Share late-opened diagnostics with pack runtimes from the same serving
891    /// host. Independent runtimes retain independent observers.
892    pub fn with_diagnostic_observer_from(mut self, main: &KhiveRuntime) -> Self {
893        self.late_diagnostic_backends = Arc::clone(&main.late_diagnostic_backends);
894        self
895    }
896
897    fn register_late_diagnostic_pool(
898        &self,
899        backend_name: &'static str,
900        pool: &Arc<ConnectionPool>,
901    ) {
902        let mut backends = self
903            .late_diagnostic_backends
904            .lock()
905            .unwrap_or_else(std::sync::PoisonError::into_inner);
906        backends.retain(|entry| entry.pool.strong_count() > 0);
907        if let Some(existing) = backends.iter_mut().find(|entry| {
908            entry
909                .pool
910                .upgrade()
911                .is_some_and(|live| Arc::ptr_eq(&live, pool))
912        }) {
913            if !existing.backend_names.contains(&backend_name) {
914                existing.backend_names.push(backend_name);
915            }
916            return;
917        }
918        backends.push(LateOpenedDiagnosticBackend {
919            backend_names: vec![backend_name],
920            pool: Arc::downgrade(pool),
921        });
922    }
923
924    fn live_late_diagnostic_backends(&self) -> Vec<OpenedDiagnosticBackend> {
925        let mut backends = self
926            .late_diagnostic_backends
927            .lock()
928            .unwrap_or_else(std::sync::PoisonError::into_inner);
929        let mut live = Vec::new();
930        backends.retain(|entry| {
931            let Some(pool) = entry.pool.upgrade() else {
932                return false;
933            };
934            live.push(OpenedDiagnosticBackend {
935                backend_names: entry
936                    .backend_names
937                    .iter()
938                    .map(|name| name.to_string())
939                    .collect(),
940                canonical_path: pool.canonical_path().map(PathBuf::from),
941                pool,
942            });
943            true
944        });
945        live
946    }
947
948    pub fn diagnostic_backends(&self) -> Arc<[OpenedDiagnosticBackend]> {
949        let mut opened: Vec<OpenedDiagnosticBackend> = if self.diagnostic_backends.is_empty() {
950            let main = self.core().backend.pool_arc();
951            vec![OpenedDiagnosticBackend {
952                backend_names: vec![BackendId::MAIN.to_string()],
953                canonical_path: main.canonical_path().map(PathBuf::from),
954                pool: main,
955            }]
956        } else {
957            self.diagnostic_backends.iter().cloned().collect()
958        };
959        for late in self.live_late_diagnostic_backends() {
960            if let Some(existing) = opened
961                .iter_mut()
962                .find(|existing| same_diagnostic_database(existing, &late))
963            {
964                for name in late.backend_names {
965                    if !existing.backend_names.contains(&name) {
966                        existing.backend_names.push(name);
967                    }
968                }
969            } else {
970                opened.push(late);
971            }
972        }
973        opened.into()
974    }
975
976    /// Whether this runtime selects the vector arm for a hybrid search —
977    /// true exactly when a default embedding model is configured. Single
978    /// source of truth for the policy every fan-out and single-backend
979    /// dispatch path uses to report `arm_participation`/`vector_selected`.
980    pub fn vector_arm_selected(&self) -> bool {
981        self.config.embedding_model.is_some()
982    }
983
984    /// Return a reference to the underlying storage backend.
985    ///
986    /// This is an embedder/infrastructure surface (connection pools, schema
987    /// plans, diagnostics). Stores obtained from it are NOT wrapped by the
988    /// message-evidence policy that [`Self::notes`] enforces: an embedder
989    /// holding the backend already holds root-equivalent access to the
990    /// database file, so the policy boundary sits at the typed accessors
991    /// pack code uses, not here. Pack code must not take note stores from
992    /// this surface.
993    pub fn backend(&self) -> &StorageBackend {
994        &self.backend
995    }
996
997    /// Whether this runtime's bound backend is explicitly or filesystem-mode
998    /// detected read-only.
999    pub fn is_read_only(&self) -> bool {
1000        self.backend.is_read_only()
1001    }
1002
1003    /// Return the directory containing the backend's database file, or `None`
1004    /// for an in-memory backend.
1005    pub fn backend_data_dir(&self) -> Option<std::path::PathBuf> {
1006        self.backend.data_dir()
1007    }
1008
1009    /// Root directory for this database's ANN segment tree (`<db-file>.ann/`
1010    /// beside the file), or `None` for an in-memory backend. Scoped to the
1011    /// database file itself so two databases sharing a parent directory can
1012    /// never adopt each other's segments.
1013    pub fn backend_ann_root(&self) -> Option<std::path::PathBuf> {
1014        self.backend.ann_root()
1015    }
1016
1017    /// Writer-contention, graph-edge integrity, and WAL/checkpoint diagnostics
1018    /// (ADR-091/ADR-135 operator surface): pooled writer and audit-failure
1019    /// counters, build identity, duplicate edge-ID and list-ledger counts,
1020    /// checkpoint counters, a PASSIVE checkpoint probe, WAL file size, and
1021    /// explicitly qualified WAL-pin census. Not write-free: the
1022    /// PASSIVE probe may backfill WAL frames into the database (normal
1023    /// checkpoint I/O). It never changes logical state, escalates to TRUNCATE,
1024    /// creates a missing database file, or deletes sidecar evidence — see
1025    /// `khive_db::diagnostics` for the narrowings that make those claims hold.
1026    ///
1027    /// Always targets the *main* backend via [`Self::core`], regardless of
1028    /// which backend this runtime handle is bound to, so a report never
1029    /// describes a database this handle is not the canonical owner of.
1030    pub async fn db_diagnostics(&self) -> RuntimeResult<khive_db::diagnostics::DbDiagnostics> {
1031        // No `VerbRegistry` handle is reachable from a bare `KhiveRuntime`
1032        // (the audit-batch seam is owned by whichever registry was built
1033        // over this runtime's `EventStore`, not by the runtime itself), so
1034        // the batch-health fields report unavailable with a reason here.
1035        // Callers that hold the registry — e.g. the `db_diagnostics` verb
1036        // handler — use `Self::db_diagnostics_with_audit_metrics` with
1037        // `VerbRegistry::audit_batch_metrics()` instead.
1038        self.db_diagnostics_with_audit_metrics(None).await
1039    }
1040
1041    /// As [`Self::db_diagnostics`], but with the caller supplying the
1042    /// ADR-133 audit-batch health counters from whichever `VerbRegistry`
1043    /// owns the seam over this runtime's `EventStore` (typically
1044    /// `VerbRegistry::audit_batch_metrics()`). `None` behaves identically to
1045    /// [`Self::db_diagnostics`].
1046    pub async fn db_diagnostics_with_audit_metrics(
1047        &self,
1048        runtime_audit_batch_metrics: Option<khive_db::diagnostics::RuntimeAuditBatchMetrics>,
1049    ) -> RuntimeResult<khive_db::diagnostics::DbDiagnostics> {
1050        let pool = self.core().backend.pool_arc();
1051        // Match housekeeping's compiled legacy-record fallback (ADR-091
1052        // Amendment 6), independent of checkpoint or local sweep overrides.
1053        let legacy_sweep_interval = khive_db::SessionSweepConfig::default().interval;
1054        let build_hash = crate::build_info::BUILD_INFO
1055            .is_stamped()
1056            .then_some(crate::build_info::BUILD_INFO.source_revision);
1057        let build = khive_db::diagnostics::BuildIdentity::from_env(
1058            crate::build_info::PACKAGE_VERSION,
1059            build_hash,
1060        );
1061
1062        let mut report = khive_db::diagnostics::collect_with_runtime_audit_metrics_interruptibly(
1063            pool,
1064            build,
1065            legacy_sweep_interval,
1066            crate::pack::audit_append_failure_count(),
1067            runtime_audit_batch_metrics,
1068        )
1069        .await
1070        .map_err(RuntimeError::from)?;
1071        report.writer_contention.audit_obligation_append_failures =
1072            Some(crate::pack::audit_obligation_append_failure_count());
1073        report
1074            .writer_contention
1075            .audit_obligation_append_failures_unavailable_reason = None;
1076        let (ann_routes, fallback_routes) = crate::note_search_ann::route_totals();
1077        report.note_search_ann_route_total = ann_routes;
1078        report.note_search_fallback_route_total = fallback_routes;
1079        Ok(report)
1080    }
1081
1082    /// Collect the same per-file report as the primary diagnostic surface for
1083    /// one opened backend. The process identity comes from main; invoking
1084    /// `ProcessIdentity::current` on a secondary would miscount main generations.
1085    pub async fn db_diagnostics_for_opened_backend_with_audit_metrics(
1086        &self,
1087        backend: &OpenedDiagnosticBackend,
1088        runtime_audit_batch_metrics: Option<khive_db::diagnostics::RuntimeAuditBatchMetrics>,
1089    ) -> RuntimeResult<khive_db::diagnostics::DbDiagnostics> {
1090        let main_pool = self.core().backend.pool_arc();
1091        let process = khive_db::diagnostics::ProcessIdentity::current(&main_pool);
1092        let build_hash = crate::build_info::BUILD_INFO
1093            .is_stamped()
1094            .then_some(crate::build_info::BUILD_INFO.source_revision);
1095        let build = khive_db::diagnostics::BuildIdentity::from_env(
1096            crate::build_info::PACKAGE_VERSION,
1097            build_hash,
1098        );
1099        let mut report =
1100            khive_db::diagnostics::collect_with_runtime_audit_metrics_for_process_interruptibly(
1101                Arc::clone(&backend.pool),
1102                build,
1103                process,
1104                khive_db::SessionSweepConfig::default().interval,
1105                crate::pack::audit_append_failure_count(),
1106                runtime_audit_batch_metrics,
1107            )
1108            .await
1109            .map_err(RuntimeError::from)?;
1110        report.writer_contention.audit_obligation_append_failures =
1111            Some(crate::pack::audit_obligation_append_failure_count());
1112        report
1113            .writer_contention
1114            .audit_obligation_append_failures_unavailable_reason = None;
1115        let (ann_routes, fallback_routes) = crate::note_search_ann::route_totals();
1116        report.note_search_ann_route_total = ann_routes;
1117        report.note_search_fallback_route_total = fallback_routes;
1118        Ok(report)
1119    }
1120
1121    // ---- Store accessors (token-scoped) ----
1122
1123    /// Get an EntityStore scoped to the token's namespace.
1124    pub fn entities(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn EntityStore>> {
1125        Ok(self
1126            .backend
1127            .entities_for_namespace(token.namespace().as_str())?)
1128    }
1129
1130    /// Get a GraphStore scoped to the token's namespace.
1131    pub fn graph(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn GraphStore>> {
1132        Ok(self
1133            .backend
1134            .graph_for_namespace(token.namespace().as_str())?)
1135    }
1136
1137    /// Get a NoteStore scoped to the token's namespace.
1138    ///
1139    /// Wrapped in `note_store_guard::PolicyEnforcingNoteStore`, which
1140    /// refuses any insert/upsert of a `kind = "message"` note carrying
1141    /// `quarantined` / `channel_kind` / `channel_slug` — the transport-owned
1142    /// evidence `comm.health` trusts at face value — and refuses patching
1143    /// those keys through the property-mutation seams on any note kind, so
1144    /// the guard cannot be sidestepped by inserting a clean message note and
1145    /// patching the evidence onto it afterward. Full-row writes also preserve
1146    /// existing channel-health coordinates while allowing heartbeat metadata to
1147    /// change. The trusted channel-ingest path does not go through this accessor; see
1148    /// `Self::raw_notes` and [`Self::try_create_note_as_trusted_ingest`].
1149    pub fn notes(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn NoteStore>> {
1150        Ok(crate::note_store_guard::PolicyEnforcingNoteStore::wrap(
1151            self.raw_notes(token)?,
1152        ))
1153    }
1154
1155    /// Get the unwrapped, policy-free NoteStore scoped to the token's namespace.
1156    ///
1157    /// Bypasses `note_store_guard::PolicyEnforcingNoteStore`. Callers
1158    /// within this crate that have already enforced the reserved-transport-
1159    /// property policy themselves (namely `try_create_note_impl`, which
1160    /// applies it conditionally based on whether the caller presented a
1161    /// [`crate::pack::ChannelIngestCapability`]) use this to reach storage
1162    /// directly rather than run a redundant, less-informed check. Not exposed
1163    /// outside this crate — every other caller must use [`Self::notes`].
1164    pub(crate) fn raw_notes(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn NoteStore>> {
1165        Ok(self
1166            .backend
1167            .notes_for_namespace(token.namespace().as_str())?)
1168    }
1169
1170    /// Return the role-keyed attachment substrate on the canonical main backend.
1171    ///
1172    /// Attachment rows are the process-shared BlobStore's sole SQL liveness
1173    /// authority. A runtime bound directly to a secondary pack backend must call
1174    /// [`Self::core`] first; accepting a secondary mutation here would create a
1175    /// reference that the main-database GC sweep cannot see or fence.
1176    pub fn attachments(&self) -> RuntimeResult<Arc<dyn AttachmentStore>> {
1177        if self.config.backend_id.as_str() != BackendId::MAIN {
1178            return Err(RuntimeError::InvalidInput(format!(
1179                "attachments are owned by the canonical main backend; runtime backend {:?} must route through KhiveRuntime::core()",
1180                self.config.backend_id.as_str()
1181            )));
1182        }
1183        Ok(self.backend.attachments()?)
1184    }
1185
1186    /// Get an EventStore scoped to the token's namespace.
1187    ///
1188    /// When the events-daemon split (ADR-170) is configured, the store routes
1189    /// by append class: the ADR-133 idempotent audit-batch lane — the
1190    /// measured bulk of event write volume — persists to the events database
1191    /// (forwarded over the events daemon socket in daemon deployments, or
1192    /// opened directly in embedded/one-shot contexts), while plain appends
1193    /// stay on this runtime's backend, keeping every raw-SQL consumer of the
1194    /// legacy `events` table (schedule provenance, kg projection guards,
1195    /// GraphQuery's substrate union) correct by construction. Reads merge
1196    /// both stores. Unconfigured runtimes (tests, in-memory) keep the legacy
1197    /// main-store behavior. Every returned store is decorated at this typed
1198    /// accessor boundary so append callers cannot override the namespace or
1199    /// actor resolved into the sealed authorization token.
1200    pub fn events(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn EventStore>> {
1201        Ok(crate::event_store_guard::AttributedEventStore::wrap(
1202            self.raw_events_for_namespace(token.namespace().as_str())?,
1203            token,
1204        ))
1205    }
1206
1207    /// The event sidecar inherits the already-open MAIN pool's WAL policy.
1208    /// A secondary pack's config or pool does not govern this shared file.
1209    fn events_wal_ceiling_policy(&self) -> khive_db::WalCeilingPolicy {
1210        self.core_backend
1211            .as_ref()
1212            .unwrap_or(&self.backend)
1213            .pool()
1214            .config()
1215            .wal_ceiling
1216    }
1217
1218    /// Build the undecorated event store used only by the registry's trusted
1219    /// audit composer, which stamps from each resolved `GateRequest` before
1220    /// enqueueing. Pack/runtime call sites must use [`Self::events`] instead.
1221    pub(crate) fn raw_events_for_namespace(
1222        &self,
1223        namespace: &str,
1224    ) -> RuntimeResult<Arc<dyn EventStore>> {
1225        let legacy = self.backend.events_for_namespace(namespace)?;
1226        match &self.config.events_split {
1227            None => Ok(legacy),
1228            Some(split) => {
1229                // Read-only is decided before the transport question: a
1230                // read-only runtime must neither create nor schema-initialize
1231                // an events database, and it must not forward writes to the
1232                // events daemon either — a socket in the config describes the
1233                // deployment, not this process's authority. Serve merged
1234                // reads from a read-only open of the sidecar when it exists
1235                // (the storage-layer read-only binding refuses any write that
1236                // slips through), and the legacy store alone otherwise: no
1237                // sidecar on disk means no lane rows exist, so minting the
1238                // file just to read nothing from it would be a write in
1239                // disguise.
1240                if self.backend.is_read_only() {
1241                    if !split.db_path.exists() {
1242                        return Ok(legacy);
1243                    }
1244                    let lane_backend = crate::events_split::direct_backend_with_policies(
1245                        &split.db_path,
1246                        true,
1247                        Some(self.backend.pool().config().max_readers),
1248                        self.events_wal_ceiling_policy(),
1249                        None,
1250                        self.events_volume_lock_dir(),
1251                    )?;
1252                    self.register_late_diagnostic_pool("events", &lane_backend.pool_arc());
1253                    let lane = lane_backend.events_for_namespace(namespace)?;
1254                    return Ok(Arc::new(crate::events_split::SplitEventStore::new(
1255                        legacy, lane,
1256                    )));
1257                }
1258                let lane: Arc<dyn EventStore> = match &split.socket_path {
1259                    #[cfg(unix)]
1260                    Some(socket) => {
1261                        let client = crate::events_split::client_for(socket)?;
1262                        Arc::new(crate::events_split::ForwardingEventStore::new(
1263                            namespace, client,
1264                        ))
1265                    }
1266                    #[cfg(not(unix))]
1267                    Some(_socket) => {
1268                        return Err(RuntimeError::InvalidInput(
1269                            "events-daemon socket forwarding requires a Unix platform; \
1270                             configure the events split in direct mode here"
1271                                .to_string(),
1272                        ));
1273                    }
1274                    None => {
1275                        let lane_backend = crate::events_split::direct_backend_with_policies(
1276                            &split.db_path,
1277                            false,
1278                            Some(self.backend.pool().config().max_readers),
1279                            self.events_wal_ceiling_policy(),
1280                            Some(self.events_disk_guard_policy()?),
1281                            self.events_volume_lock_dir(),
1282                        )?;
1283                        self.register_late_diagnostic_pool("events", &lane_backend.pool_arc());
1284                        lane_backend.events_for_namespace(namespace)?
1285                    }
1286                };
1287                Ok(Arc::new(crate::events_split::SplitEventStore::new(
1288                    legacy, lane,
1289                )))
1290            }
1291        }
1292    }
1293
1294    /// Get the raw SQL access capability (for ad-hoc queries).
1295    pub fn sql(&self) -> Arc<dyn SqlAccess> {
1296        self.backend.sql()
1297    }
1298
1299    /// SQL access to the events-split sidecar database for read purposes,
1300    /// when the split (ADR-170) is configured and the sidecar exists on
1301    /// disk. `None` means every event row lives in the legacy `events`
1302    /// table, so a raw-SQL consumer needs no second lookup. Consumers that
1303    /// resolve an event by id or hex prefix against `self.sql()` must also
1304    /// consult this store on a miss: the audit-batch lane's rows live only
1305    /// in the sidecar.
1306    ///
1307    /// A writable runtime opens the sidecar through the ordinary writable
1308    /// binding: WAL supports one-writer-many-readers, and the read-only
1309    /// binding's frozen-snapshot guard refuses any sidecar with a live
1310    /// writer's `-shm` beside it — exactly the live-deployment case these
1311    /// reads exist for. A read-only runtime keeps the read-only open (it
1312    /// must neither create nor schema-initialize a sidecar), which serves
1313    /// genuinely frozen snapshots and refuses live ones, matching the
1314    /// read-only arm of `events()`. Never creates a sidecar as a side
1315    /// effect of a read.
1316    pub fn events_sidecar_sql_read_only(&self) -> RuntimeResult<Option<Arc<dyn SqlAccess>>> {
1317        match &self.config.events_split {
1318            None => Ok(None),
1319            Some(split) => {
1320                if !split.db_path.exists() {
1321                    return Ok(None);
1322                }
1323                let backend = if self.backend.is_read_only() {
1324                    crate::events_split::direct_backend_with_policies(
1325                        &split.db_path,
1326                        true,
1327                        Some(self.backend.pool().config().max_readers),
1328                        self.events_wal_ceiling_policy(),
1329                        None,
1330                        self.events_volume_lock_dir(),
1331                    )?
1332                } else {
1333                    crate::events_split::direct_backend_with_policies(
1334                        &split.db_path,
1335                        false,
1336                        Some(self.backend.pool().config().max_readers),
1337                        self.events_wal_ceiling_policy(),
1338                        Some(self.events_disk_guard_policy()?),
1339                        self.events_volume_lock_dir(),
1340                    )?
1341                };
1342                self.register_late_diagnostic_pool("events", &backend.pool_arc());
1343                Ok(Some(backend.sql()))
1344            }
1345        }
1346    }
1347
1348    /// Get a VectorStore for the configured embedding model, scoped to the token's namespace.
1349    ///
1350    /// Returns `Unconfigured("embedding_model")` if no model is set.
1351    pub fn vectors(
1352        &self,
1353        token: &NamespaceToken,
1354    ) -> RuntimeResult<Arc<dyn khive_storage::VectorStore>> {
1355        let model = self.resolve_embedding_model(None)?;
1356        self.vectors_for_embedding_model(token, model)
1357    }
1358
1359    /// Get a VectorStore for a specific named embedding model, scoped to the token's namespace.
1360    ///
1361    /// Accepts both built-in lattice model names/aliases and custom provider names
1362    /// registered via [`register_embedder`](Self::register_embedder). Lattice names
1363    /// are routed through the enum-backed path; custom provider names use the
1364    /// provider's declared `dimensions()` directly so that the vector store key
1365    /// is consistent with how vectors were written during `remember`/`recall`.
1366    pub fn vectors_for_model(
1367        &self,
1368        token: &NamespaceToken,
1369        model_name: &str,
1370    ) -> RuntimeResult<Arc<dyn khive_storage::VectorStore>> {
1371        let (model_name, dims) = self.vector_model_metadata(model_name)?;
1372        Ok(self.backend.vectors_for_namespace(
1373            &sanitize_key(&model_name),
1374            &model_name,
1375            dims,
1376            token.namespace().as_str(),
1377        )?)
1378    }
1379
1380    /// Resolve the storage identity and declared dimensions together so guarded
1381    /// SQL publication agrees with VectorStore, including built-in aliases.
1382    pub(crate) fn vector_model_metadata(&self, model_name: &str) -> RuntimeResult<(String, usize)> {
1383        if request_excludes_embedder(model_name) {
1384            return Err(crate::RuntimeError::UnknownModel(model_name.to_string()));
1385        }
1386        let registry = self
1387            .embedder_registry
1388            .read()
1389            .map_err(|_| crate::RuntimeError::Internal("embedder registry lock poisoned".into()))?;
1390        if let Some(model) = parse_embedding_model_alias(model_name) {
1391            // Only proceed via the lattice path if this model is actually in the
1392            // registry; otherwise fall through to the custom-provider path.
1393            let key = model.to_string();
1394            if registry.contains(&key) {
1395                return Ok((key, model.dimensions()));
1396            }
1397        }
1398        registry
1399            .get_provider(model_name)
1400            .map(|provider| (model_name.to_owned(), provider.dimensions()))
1401            .ok_or_else(|| crate::RuntimeError::UnknownModel(model_name.to_string()))
1402    }
1403
1404    /// Get a namespace-scoped vector store for a pack-owned immutable identity.
1405    ///
1406    /// The table key is syntactically validated by [`NamedVectorIdentity`]. This
1407    /// accessor additionally verifies the table's actual sqlite-vec dimension
1408    /// declaration and every persisted `embedding_model` value before returning
1409    /// the store, so reusing one key for incompatible descriptor geometry or
1410    /// semantics fails before a caller can replace rows.
1411    pub async fn vectors_for_named_identity(
1412        &self,
1413        token: &NamespaceToken,
1414        identity: &NamedVectorIdentity,
1415    ) -> RuntimeResult<Arc<dyn VectorStore>> {
1416        let namespace = token.namespace().as_str();
1417        {
1418            let cached = self.named_vector_stores.read().map_err(|_| {
1419                RuntimeError::Internal("named vector store cache lock poisoned".into())
1420            })?;
1421            if let Some(entry) = cached.get(identity.model_key()) {
1422                check_cached_named_vector_identity(&entry.identity, identity)?;
1423                if let Some(store) = entry.by_namespace.get(namespace) {
1424                    return Ok(Arc::clone(store));
1425                }
1426            }
1427        }
1428        let store = self.backend.vectors_for_namespace(
1429            identity.model_key(),
1430            identity.model_name(),
1431            identity.dimensions(),
1432            namespace,
1433        )?;
1434
1435        let table = format!("vec_{}", identity.model_key());
1436        let mut reader = self.sql().reader().await?;
1437        let dimension_row = reader
1438            .query_row(SqlStatement {
1439                sql: "SELECT sql FROM sqlite_schema WHERE type = 'table' AND name = ?1".to_string(),
1440                params: vec![SqlValue::Text(table.clone())],
1441                label: Some("runtime_named_vector_dimension".to_string()),
1442            })
1443            .await?
1444            .ok_or_else(|| {
1445                RuntimeError::Internal(format!(
1446                    "named vector table {table} has no sqlite_schema declaration"
1447                ))
1448            })?;
1449        let table_ddl = match dimension_row.get("sql") {
1450            Some(SqlValue::Text(value)) => value,
1451            other => {
1452                return Err(RuntimeError::Internal(format!(
1453                    "named vector table {table} returned invalid schema metadata: {other:?}"
1454                )))
1455            }
1456        };
1457        let declared_dimensions = vector_dimensions_from_ddl(table_ddl).ok_or_else(|| {
1458            RuntimeError::Internal(format!(
1459                "named vector table {table} has no parseable embedding dimension"
1460            ))
1461        })?;
1462        if declared_dimensions != identity.dimensions() {
1463            return Err(RuntimeError::InvalidInput(format!(
1464                "named vector model_key {:?} is already bound to {declared_dimensions} dimensions, expected {}",
1465                identity.model_key(),
1466                identity.dimensions()
1467            )));
1468        }
1469
1470        let stored_models = reader
1471            .query_all(SqlStatement {
1472                sql: format!(
1473                    "SELECT DISTINCT embedding_model FROM {table} ORDER BY embedding_model LIMIT 2"
1474                ),
1475                params: vec![],
1476                label: Some("runtime_named_vector_model_identity".to_string()),
1477            })
1478            .await?;
1479        for row in stored_models {
1480            let stored = match row.get("embedding_model") {
1481                Some(SqlValue::Text(value)) => value,
1482                other => {
1483                    return Err(RuntimeError::Internal(format!(
1484                        "named vector table {table} returned invalid model identity metadata: {other:?}"
1485                    )))
1486                }
1487            };
1488            if stored != identity.model_name() {
1489                return Err(RuntimeError::InvalidInput(format!(
1490                    "named vector model_key {:?} already contains model {stored:?}, cannot bind it to {:?}",
1491                    identity.model_key(),
1492                    identity.model_name()
1493                )));
1494            }
1495        }
1496
1497        self.backend
1498            .register_embedding_model(
1499                identity.model_key(),
1500                identity.model_name(),
1501                identity.model_key(),
1502                identity.dimensions() as u32,
1503            )
1504            .map_err(|error| {
1505                if matches!(
1506                    &error,
1507                    khive_db::SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
1508                        if code.code == rusqlite::ErrorCode::ConstraintViolation
1509                ) {
1510                    RuntimeError::InvalidInput(format!(
1511                        "named vector model_key {:?} is already bound to a different active model identity",
1512                        identity.model_key()
1513                    ))
1514                } else {
1515                    RuntimeError::Sqlite(error)
1516                }
1517            })?;
1518
1519        let mut cached = self
1520            .named_vector_stores
1521            .write()
1522            .map_err(|_| RuntimeError::Internal("named vector store cache lock poisoned".into()))?;
1523        let entry = cached
1524            .entry(identity.model_key().to_owned())
1525            .or_insert_with(|| CachedNamedVectorStores {
1526                identity: identity.clone(),
1527                by_namespace: HashMap::new(),
1528            });
1529        check_cached_named_vector_identity(&entry.identity, identity)?;
1530        entry
1531            .by_namespace
1532            .insert(namespace.to_owned(), Arc::clone(&store));
1533        Ok(store)
1534    }
1535
1536    /// Output dimensions for a named embedding model, resolved from the
1537    /// embedder registry alone — no storage access. Mirrors
1538    /// [`vectors_for_model`](Self::vectors_for_model)'s resolution order:
1539    /// lattice aliases route through the enum when registered, otherwise the
1540    /// custom provider's declared `dimensions()`. `None` when no such model
1541    /// is registered.
1542    pub fn embedder_dimensions(&self, model_name: &str) -> Option<usize> {
1543        if request_excludes_embedder(model_name) {
1544            return None;
1545        }
1546        if let Some(model) = parse_embedding_model_alias(model_name) {
1547            let key = model.to_string();
1548            let in_registry = self
1549                .embedder_registry
1550                .read()
1551                .map(|reg| reg.contains(&key))
1552                .unwrap_or(false);
1553            if in_registry {
1554                return Some(model.dimensions());
1555            }
1556        }
1557        self.embedder_registry
1558            .read()
1559            .ok()?
1560            .get_provider(model_name)
1561            .map(|p| p.dimensions())
1562    }
1563
1564    fn vectors_for_embedding_model(
1565        &self,
1566        token: &NamespaceToken,
1567        model: EmbeddingModel,
1568    ) -> RuntimeResult<Arc<dyn khive_storage::VectorStore>> {
1569        Ok(self.backend.vectors_for_namespace(
1570            &vec_model_key(model),
1571            &model.to_string(),
1572            model.dimensions(),
1573            token.namespace().as_str(),
1574        )?)
1575    }
1576
1577    /// Get a TextSearch index for the entity corpus (single shared table).
1578    pub fn text(
1579        &self,
1580        token: &NamespaceToken,
1581    ) -> RuntimeResult<Arc<dyn khive_storage::TextSearch>> {
1582        let _ = token;
1583        Ok(self.backend.text("entities")?)
1584    }
1585
1586    /// Get a TextSearch index for the notes corpus (single shared table).
1587    pub fn text_for_notes(
1588        &self,
1589        token: &NamespaceToken,
1590    ) -> RuntimeResult<Arc<dyn khive_storage::TextSearch>> {
1591        let _ = token;
1592        Ok(self.backend.text("notes")?)
1593    }
1594
1595    /// Mint an authorization token for the given namespace.
1596    ///
1597    /// Consults the configured [`crate::Gate`] before minting. With the default
1598    /// `AllowAllGate` this always succeeds. When a real policy-backed gate is
1599    /// installed, this method enforces it and returns `PermissionDenied` on
1600    /// denial.
1601    ///
1602    /// The returned token's read visibility set defaults to `[ns]` — identical
1603    /// to the pre-visibility-set behaviour. Use [`Self::authorize_with_visibility`]
1604    /// to mint a token that can read additional namespaces.
1605    ///
1606    /// When `actor_id` is configured in `RuntimeConfig`, the token carries that
1607    /// actor label so that `comm.inbox` filters by `to_actor`. When
1608    /// unconfigured, the token carries `ActorRef::anonymous()` and inbox falls
1609    /// back to party-line behavior.
1610    pub fn authorize(&self, ns: Namespace) -> RuntimeResult<NamespaceToken> {
1611        let actor = crate::actor_identity::resolve_actor(self.config.actor_id.as_deref());
1612        let req = GateRequest::new(
1613            actor.clone(),
1614            ns.clone(),
1615            "authorize",
1616            serde_json::Value::Null,
1617        );
1618        match self.config.gate.check(&req) {
1619            Ok(ref decision) if decision.is_allow() => {
1620                if let khive_gate::GateDecision::Allow { ref obligations } = decision {
1621                    if !obligations.is_empty() {
1622                        tracing::debug!(
1623                            namespace = %ns.as_str(),
1624                            "authorize: obligations={:?}",
1625                            obligations
1626                        );
1627                    }
1628                }
1629                Ok(NamespaceToken::mint_authorized(ns, actor))
1630            }
1631            Ok(khive_gate::GateDecision::Deny { reason }) => {
1632                Err(crate::RuntimeError::permission_denied("authorize", reason))
1633            }
1634            Ok(_) => Err(crate::RuntimeError::permission_denied(
1635                "authorize",
1636                "gate denied",
1637            )),
1638            Err(e) => {
1639                tracing::warn!(
1640                    namespace = %ns.as_str(),
1641                    error = %crate::secret_gate::bounded_masked_log_text(&e.to_string()),
1642                    "authorize: gate check failed (fail-closed)"
1643                );
1644                Err(crate::RuntimeError::Internal(format!(
1645                    "gate error: {}",
1646                    e.wire_reason()
1647                )))
1648            }
1649        }
1650    }
1651
1652    /// Mint an authorization token with an explicit read-visibility set.
1653    ///
1654    /// `primary` is the **write namespace** — all records created via the
1655    /// returned token land there. `extra_visible` lists additional namespaces
1656    /// the token may read. The primary is always included in the visible set
1657    /// regardless of `extra_visible`.
1658    ///
1659    /// Usage (lambda:leo reading both leo and khive namespaces):
1660    /// ```rust,ignore
1661    /// let tok = rt.authorize_with_visibility(
1662    ///     Namespace::parse("lambda:leo").unwrap(),
1663    ///     vec![Namespace::parse("lambda:khive").unwrap()],
1664    /// )?;
1665    /// ```
1666    pub fn authorize_with_visibility(
1667        &self,
1668        primary: Namespace,
1669        extra_visible: Vec<Namespace>,
1670    ) -> RuntimeResult<NamespaceToken> {
1671        let actor = crate::actor_identity::resolve_actor(self.config.actor_id.as_deref());
1672        let req = GateRequest::new(
1673            actor.clone(),
1674            primary.clone(),
1675            "authorize",
1676            serde_json::Value::Null,
1677        );
1678        match self.config.gate.check(&req) {
1679            Ok(ref decision) if decision.is_allow() => {
1680                if let khive_gate::GateDecision::Allow { ref obligations } = decision {
1681                    if !obligations.is_empty() {
1682                        tracing::debug!(
1683                            namespace = %primary.as_str(),
1684                            "authorize_with_visibility: obligations={:?}",
1685                            obligations
1686                        );
1687                    }
1688                }
1689                // The primary check authorizes writes to `primary` only. Each
1690                // extra namespace grants read visibility, so each one takes
1691                // its own Read-classified gate check before it may enter the
1692                // minted set — a token must never carry visibility the gate
1693                // was not asked about. Any deny or gate error refuses the
1694                // whole mint, naming the offending namespace (fail-closed).
1695                for extra in &extra_visible {
1696                    let extra_req = GateRequest::new(
1697                        actor.clone(),
1698                        extra.clone(),
1699                        "authorize.visible",
1700                        serde_json::Value::Null,
1701                    );
1702                    match self.config.gate.check(&extra_req) {
1703                        Ok(ref extra_decision) if extra_decision.is_allow() => {}
1704                        Ok(khive_gate::GateDecision::Deny { reason }) => {
1705                            return Err(crate::RuntimeError::permission_denied(
1706                                "authorize",
1707                                format!(
1708                                    "visibility namespace {:?} denied: {reason}",
1709                                    extra.as_str()
1710                                ),
1711                            ));
1712                        }
1713                        Ok(_) => {
1714                            return Err(crate::RuntimeError::permission_denied(
1715                                "authorize",
1716                                format!("visibility namespace {:?} denied by gate", extra.as_str()),
1717                            ));
1718                        }
1719                        Err(e) => {
1720                            tracing::warn!(
1721                                namespace = %extra.as_str(),
1722                                error = %crate::secret_gate::bounded_masked_log_text(&e.to_string()),
1723                                "authorize_with_visibility: extra-namespace gate check failed (fail-closed)"
1724                            );
1725                            return Err(crate::RuntimeError::Internal(format!(
1726                                "gate error: {}",
1727                                e.wire_reason()
1728                            )));
1729                        }
1730                    }
1731                }
1732                Ok(NamespaceToken::mint_with_visibility(
1733                    primary,
1734                    extra_visible,
1735                    actor,
1736                ))
1737            }
1738            Ok(khive_gate::GateDecision::Deny { reason }) => {
1739                Err(crate::RuntimeError::permission_denied("authorize", reason))
1740            }
1741            Ok(_) => Err(crate::RuntimeError::permission_denied(
1742                "authorize",
1743                "gate denied",
1744            )),
1745            Err(e) => {
1746                tracing::warn!(
1747                    namespace = %primary.as_str(),
1748                    error = %crate::secret_gate::bounded_masked_log_text(&e.to_string()),
1749                    "authorize_with_visibility: gate check failed (fail-closed)"
1750                );
1751                Err(crate::RuntimeError::Internal(format!(
1752                    "gate error: {}",
1753                    e.wire_reason()
1754                )))
1755            }
1756        }
1757    }
1758
1759    /// Install the pack-aggregated edge endpoint rules.
1760    ///
1761    /// Called by the transport layer after the `VerbRegistry` is built so
1762    /// that runtime-layer edge validation can consult pack rules. Idempotent:
1763    /// later calls overwrite the previous rule set.
1764    pub fn install_edge_rules(&self, rules: Vec<EdgeEndpointRule>) {
1765        if let Ok(mut guard) = self.edge_rules.write() {
1766            *guard = rules;
1767        }
1768    }
1769
1770    /// Install an already-paired blob hydrator into this runtime.
1771    ///
1772    /// Reinstalling the exact same `Arc` is idempotent. A different pair is
1773    /// rejected: replacing it would split or reset the aggregate admission
1774    /// budget while requests may still hold leases.
1775    pub fn install_blob_hydrator(
1776        &self,
1777        hydrator: Arc<crate::blob::BlobHydrator>,
1778    ) -> RuntimeResult<()> {
1779        if hydrator.budget_bytes() != self.config.blob_hydration_bytes {
1780            return Err(RuntimeError::InvalidInput(format!(
1781                "blob hydrator budget {} does not match this runtime's resolved budget {}",
1782                hydrator.budget_bytes(),
1783                self.config.blob_hydration_bytes
1784            )));
1785        }
1786        // The mode gate must sit on THIS seam, not only on `install_blob_store`:
1787        // `BlobHydrator::new` is public, so without it a caller pairs a writable
1788        // store, installs it here, and `blob_store()` hands mutating pack paths
1789        // a writable store on a runtime whose declared mode is read-only.
1790        // Boot paths installing one hydrator across handles of MIXED modes —
1791        // where the blob pack's own backend mode, not each receiving
1792        // handle's, governs mutability — go through
1793        // [`Self::install_shared_blob_hydrator`] instead.
1794        if self.is_read_only() && !hydrator.enforces_read_only() {
1795            return Err(RuntimeError::InvalidInput(
1796                "this runtime is read-only: install the raw store with install_blob_store, \
1797                 which wraps it so every physical mutator refuses"
1798                    .to_string(),
1799            ));
1800        }
1801        self.install_blob_hydrator_slot(hydrator)
1802    }
1803
1804    /// Install a boot-shared hydrator whose mutability is governed by the
1805    /// blob runtime's own mode, not this handle's domain-store mode.
1806    ///
1807    /// ADR-160 D3 installs one hydrator `Arc` on every runtime handle a boot
1808    /// produces, and the documented multi-backend matrix includes a writable
1809    /// blob secondary beside a read-only main: there the shared hydrator is
1810    /// legitimately writable on a read-only domain handle. The mode decision
1811    /// must therefore already be encoded in the hydrator, and it must have
1812    /// been DERIVED, not declared: only hydrators built through
1813    /// [`crate::BlobHydrator::resolve_for_governing_backend`] — whose mode
1814    /// comes from the governing backend's own access mode — are accepted
1815    /// here. A hand-paired hydrator (`BlobHydrator::new` / `for_mode`) is
1816    /// refused so a safe downstream caller cannot use this seam to put a
1817    /// writable store on a read-only runtime; such callers use
1818    /// [`Self::install_blob_hydrator`], which holds hydrator mode against
1819    /// this runtime's own.
1820    ///
1821    /// This gate is a wrong-wiring guard, not an in-process sandbox: which
1822    /// backend governs is the boot host's topology assertion, and a caller
1823    /// who deliberately selects an unrelated writable backend as governing
1824    /// is outside what any runtime seam can enforce (see the trust-model
1825    /// note on [`crate::BlobHydrator::resolve_for_governing_backend`]).
1826    pub fn install_shared_blob_hydrator(
1827        &self,
1828        hydrator: Arc<crate::blob::BlobHydrator>,
1829    ) -> RuntimeResult<()> {
1830        if !hydrator.is_governed() {
1831            return Err(RuntimeError::InvalidInput(
1832                "the shared install seam accepts only hydrators whose mode was derived from a \
1833                 governing backend (BlobHydrator::resolve_for_governing_backend); use \
1834                 install_blob_hydrator for a hand-paired hydrator"
1835                    .to_string(),
1836            ));
1837        }
1838        if hydrator.budget_bytes() != self.config.blob_hydration_bytes {
1839            return Err(RuntimeError::InvalidInput(format!(
1840                "blob hydrator budget {} does not match this runtime's resolved budget {}",
1841                hydrator.budget_bytes(),
1842                self.config.blob_hydration_bytes
1843            )));
1844        }
1845        self.install_blob_hydrator_slot(hydrator)
1846    }
1847
1848    /// Whether `candidate` is the same pairing as the installed `current`:
1849    /// the exact `Arc`, or a distinct hydrator allocation over the same raw
1850    /// store with the same budget and mode. The latter arises when two
1851    /// concurrent first installs each construct a hydrator from one raw
1852    /// store — the `OnceLock` loser must read as idempotent, not as a
1853    /// conflicting install.
1854    fn is_same_blob_pairing(
1855        current: &Arc<crate::blob::BlobHydrator>,
1856        candidate: &Arc<crate::blob::BlobHydrator>,
1857    ) -> bool {
1858        Arc::ptr_eq(current, candidate)
1859            || (Arc::ptr_eq(&current.raw_store(), &candidate.raw_store())
1860                && current.budget_bytes() == candidate.budget_bytes()
1861                && current.enforces_read_only() == candidate.enforces_read_only())
1862    }
1863
1864    /// One-shot slot semantics shared by both install seams: first install
1865    /// wins, an equivalent pairing is idempotent, a different pairing is
1866    /// refused (replacing it would split or reset the aggregate admission
1867    /// budget while requests may still hold leases).
1868    fn install_blob_hydrator_slot(
1869        &self,
1870        hydrator: Arc<crate::blob::BlobHydrator>,
1871    ) -> RuntimeResult<()> {
1872        if let Some(current) = self.blob_hydrator.get() {
1873            return if Self::is_same_blob_pairing(current, &hydrator) {
1874                Ok(())
1875            } else {
1876                Err(RuntimeError::InvalidInput(
1877                    "a different blob hydrator is already installed".to_string(),
1878                ))
1879            };
1880        }
1881
1882        match self.blob_hydrator.set(hydrator) {
1883            Ok(()) => Ok(()),
1884            Err(candidate) => {
1885                let current = self.blob_hydrator.get().ok_or_else(|| {
1886                    RuntimeError::Internal(
1887                        "blob hydrator install raced without a visible winner".to_string(),
1888                    )
1889                })?;
1890                if Self::is_same_blob_pairing(current, &candidate) {
1891                    Ok(())
1892                } else {
1893                    Err(RuntimeError::InvalidInput(
1894                        "a different blob hydrator is already installed".to_string(),
1895                    ))
1896                }
1897            }
1898        }
1899    }
1900
1901    /// Pair and install a store using this runtime's resolved hydration budget.
1902    ///
1903    /// Boot paths that own multiple runtimes should instead construct one
1904    /// [`crate::BlobHydrator`] and call [`Self::install_blob_hydrator`] with
1905    /// the same `Arc` on every handle.
1906    pub fn install_blob_store(
1907        &self,
1908        store: Arc<dyn khive_storage::BlobStore>,
1909    ) -> RuntimeResult<()> {
1910        if let Some(current) = self.blob_hydrator.get() {
1911            // Reinstalling the same raw store is idempotent in BOTH modes:
1912            // on a read-only runtime the installed hydrator wraps the raw
1913            // store, so identity is checked against the raw handle the
1914            // hydrator remembers, not only the (possibly wrapped) paired one.
1915            let current_store = current.store();
1916            if Arc::ptr_eq(&current_store, &store) || Arc::ptr_eq(&current.raw_store(), &store) {
1917                return Ok(());
1918            }
1919        }
1920        // A read-only runtime holds its mode at this seam, not only during
1921        // boot resolution: an arbitrary store installed after launch is
1922        // wrapped so every physical mutator refuses while the bounded read
1923        // surface stays available. Without this, post-boot installation is a
1924        // writable bypass of the runtime's declared mode.
1925        let hydrator = if self.is_read_only() {
1926            crate::blob::BlobHydrator::new_read_only(store, self.config.blob_hydration_bytes)?
1927        } else {
1928            crate::blob::BlobHydrator::new(store, self.config.blob_hydration_bytes)?
1929        };
1930        self.install_blob_hydrator(Arc::new(hydrator))
1931    }
1932
1933    /// Return the installed shared blob hydrator, if boot configured one.
1934    pub fn blob_hydrator(&self) -> Option<Arc<crate::blob::BlobHydrator>> {
1935        self.blob_hydrator.get().cloned()
1936    }
1937
1938    /// Return the installed `BlobStore`, if the boot path resolved and
1939    /// installed one. `None` when no `[storage.blob]` selection was ever
1940    /// installed — e.g. a bare/test runtime constructed without going
1941    /// through the `khive-mcp` boot path.
1942    pub fn blob_store(&self) -> Option<Arc<dyn khive_storage::BlobStore>> {
1943        self.blob_hydrator.get().map(|hydrator| hydrator.store())
1944    }
1945
1946    /// Return the installed `BlobStore`, or `RuntimeError::Unconfigured` with
1947    /// the operator-facing message when none is installed.
1948    ///
1949    /// Packs share this conversion so an operator sees one error wherever the
1950    /// missing `[storage.blob]` configuration is first hit.
1951    pub fn require_blob_store(&self) -> RuntimeResult<Arc<dyn khive_storage::BlobStore>> {
1952        self.blob_store().ok_or_else(Self::no_blob_store)
1953    }
1954
1955    /// Return the installed shared blob hydrator, or `RuntimeError::Unconfigured`
1956    /// with the same message as [`Self::require_blob_store`] when none is
1957    /// installed. The store is read from the hydrator, so both accessors are
1958    /// unset under the same condition.
1959    pub fn require_blob_hydrator(&self) -> RuntimeResult<Arc<crate::blob::BlobHydrator>> {
1960        self.blob_hydrator().ok_or_else(Self::no_blob_store)
1961    }
1962
1963    fn no_blob_store() -> RuntimeError {
1964        RuntimeError::Unconfigured(
1965            "no BlobStore installed on this server (configure [storage.blob] in khive.toml, or \
1966             KHIVE_BLOB_ROOT)"
1967                .to_string(),
1968        )
1969    }
1970
1971    /// Install the pack-aggregated valid entity and note kinds.
1972    ///
1973    /// Called by the transport layer after the `VerbRegistry` is built so that
1974    /// runtime-layer entity/note creation and import validate kind strings against
1975    /// the merged pack vocabulary. Idempotent: later calls overwrite previous sets.
1976    ///
1977    /// When no kinds are installed (empty lists), kind validation is skipped at
1978    /// the runtime layer. The pack handler layer remains the primary enforcement
1979    /// point; this provides defense-in-depth for direct Rust callers and import.
1980    pub fn install_kind_registry(&self, entity_kinds: Vec<String>, note_kinds: Vec<String>) {
1981        if let Ok(mut guard) = self.valid_entity_kinds.write() {
1982            *guard = entity_kinds;
1983        }
1984        if let Ok(mut guard) = self.valid_note_kinds.write() {
1985            let prior = std::mem::take(&mut *guard);
1986            *guard = note_kinds
1987                .into_iter()
1988                .map(|name| NoteKindEntry {
1989                    embedding_policy: prior
1990                        .iter()
1991                        .find(|entry| entry.name == name)
1992                        .map(|entry| entry.embedding_policy)
1993                        .unwrap_or_default(),
1994                    name,
1995                    registered: true,
1996                })
1997                .collect();
1998        }
1999    }
2000
2001    /// Install pack-declared embedding policy on the note-kind registry.
2002    /// The transport calls this for every runtime after pack registration.
2003    pub fn install_note_embedding_policies(&self, policies: &[crate::NoteEmbeddingPolicySpec]) {
2004        if let Ok(mut guard) = self.valid_note_kinds.write() {
2005            for spec in policies {
2006                if let Some(entry) = guard.iter_mut().find(|entry| entry.name == spec.kind) {
2007                    entry.embedding_policy = spec.policy;
2008                } else {
2009                    guard.push(NoteKindEntry {
2010                        name: spec.kind.to_owned(),
2011                        embedding_policy: spec.policy,
2012                        registered: false,
2013                    });
2014                }
2015            }
2016        }
2017    }
2018
2019    /// Registered models selected by the installed embedding policy for a note kind.
2020    /// Unknown kinds retain the all-models default.
2021    pub fn embedding_models_for_note_kind(&self, kind: &str) -> Vec<String> {
2022        let policy = self
2023            .valid_note_kinds
2024            .read()
2025            .ok()
2026            .and_then(|guard| {
2027                guard
2028                    .iter()
2029                    .find(|entry| entry.name == kind)
2030                    .map(|entry| entry.embedding_policy)
2031            })
2032            .unwrap_or_default();
2033        let models = self.registered_embedding_model_names();
2034        match policy {
2035            crate::NoteEmbeddingPolicy::AllModels => models,
2036            crate::NoteEmbeddingPolicy::DefaultModel => {
2037                let default = self.default_embedder_name();
2038                models
2039                    .into_iter()
2040                    .filter(|name| name.as_str() == default)
2041                    .collect()
2042            }
2043        }
2044    }
2045
2046    /// Install the pack-owned note kinds aggregated from the pack registry.
2047    ///
2048    /// Called by the transport after the `VerbRegistry` is built, same timing
2049    /// as [`install_kind_registry`](Self::install_kind_registry).
2050    pub fn install_pack_owned_note_kinds(&self, kinds: Vec<String>) {
2051        if let Ok(mut guard) = self.pack_owned_note_kinds.write() {
2052            *guard = kinds;
2053        }
2054    }
2055
2056    /// Whether `kind` is a note kind owned by a pack (see
2057    /// [`install_pack_owned_note_kinds`](Self::install_pack_owned_note_kinds)).
2058    ///
2059    /// Always `false` before the transport installs the list — a bare runtime
2060    /// has no packs, so no kind is pack-owned there.
2061    pub fn is_pack_owned_note_kind(&self, kind: &str) -> bool {
2062        self.pack_owned_note_kinds
2063            .read()
2064            .map(|g| g.iter().any(|k| k == kind))
2065            .unwrap_or(false)
2066    }
2067
2068    /// Validate that `kind` is a pack-registered entity kind.
2069    ///
2070    /// Returns `Ok(())` when no kinds are installed (bare runtime without packs).
2071    /// Returns `InvalidInput` when kinds are installed and `kind` is not among them.
2072    pub(crate) fn validate_entity_kind(&self, kind: &str) -> crate::RuntimeResult<()> {
2073        let guard = self.valid_entity_kinds.read().map_err(|_| {
2074            crate::RuntimeError::Internal("entity kind registry lock poisoned".into())
2075        })?;
2076        if guard.is_empty() {
2077            return Ok(());
2078        }
2079        if guard.iter().any(|k| k == kind) {
2080            Ok(())
2081        } else {
2082            Err(crate::RuntimeError::InvalidInput(format!(
2083                "unknown entity kind {kind:?}; valid: {}",
2084                guard.join(", ")
2085            )))
2086        }
2087    }
2088
2089    /// Validate that `kind` is a pack-registered note kind.
2090    ///
2091    /// Returns `Ok(())` when no kinds are installed (bare runtime without packs).
2092    /// Returns `InvalidInput` when kinds are installed and `kind` is not among them.
2093    pub(crate) fn validate_note_kind(&self, kind: &str) -> crate::RuntimeResult<()> {
2094        let guard = self.valid_note_kinds.read().map_err(|_| {
2095            crate::RuntimeError::Internal("note kind registry lock poisoned".into())
2096        })?;
2097        if !guard.iter().any(|entry| entry.registered) {
2098            return Ok(());
2099        }
2100        if guard
2101            .iter()
2102            .any(|entry| entry.registered && entry.name == kind)
2103        {
2104            Ok(())
2105        } else {
2106            let valid = guard
2107                .iter()
2108                .filter(|entry| entry.registered)
2109                .map(|entry| entry.name.as_str())
2110                .collect::<Vec<_>>()
2111                .join(", ");
2112            Err(crate::RuntimeError::InvalidInput(format!(
2113                "unknown note kind {kind:?}; valid: {}",
2114                valid
2115            )))
2116        }
2117    }
2118
2119    /// Install a pack-supplied entity-type validator.
2120    ///
2121    /// Called by the `KgPack` during registration so that `create_many` can validate
2122    /// `entity_type` values at the runtime layer, closing the hole where direct Rust
2123    /// callers bypass the handler-layer `validate_entity_type` check.
2124    ///
2125    /// The callback receives `(kind, entity_type)` and returns the normalised type
2126    /// string, or `RuntimeError::InvalidInput` if the type is not registered for that
2127    /// kind. Passing `entity_type = None` must return `Ok(None)`.
2128    pub fn install_entity_type_validator(&self, f: EntityTypeValidatorFn) {
2129        if let Ok(mut guard) = self.entity_type_validator.write() {
2130            *guard = Some(f);
2131        }
2132    }
2133
2134    /// Validate and normalise `entity_type` through the pack-installed validator.
2135    ///
2136    /// Returns `Ok(entity_type)` when no validator is installed (bare runtime).
2137    /// Returns `InvalidInput` when a validator is installed and rejects the type.
2138    pub(crate) fn validate_entity_type_for_kind(
2139        &self,
2140        kind: &str,
2141        entity_type: Option<&str>,
2142    ) -> crate::RuntimeResult<Option<String>> {
2143        let guard = self.entity_type_validator.read().map_err(|_| {
2144            crate::RuntimeError::Internal("entity type validator lock poisoned".into())
2145        })?;
2146        match guard.as_ref() {
2147            None => Ok(entity_type.map(str::to_string)),
2148            Some(validate) => validate(kind, entity_type),
2149        }
2150    }
2151
2152    /// Install a pack-owned note-mutation hook.
2153    ///
2154    /// Overwrites any previously-installed hook, same single-slot semantics
2155    /// as [`install_entity_type_validator`](Self::install_entity_type_validator).
2156    /// In practice only one pack (`khive-pack-memory`) installs one today;
2157    /// if a second pack ever needs this, the slot should be widened to a
2158    /// `Vec` at that point rather than silently overwritten.
2159    pub fn install_note_mutation_hook(&self, f: NoteMutationHookFn) {
2160        if let Ok(mut guard) = self.note_mutation_hook.write() {
2161            *guard = Some(f);
2162        }
2163    }
2164
2165    /// Clone read-side backend and embedder handles without retaining this
2166    /// runtime's pack callback slots. Providers stored in one of those slots
2167    /// must not hold a clone that points back to their own installation Arc,
2168    /// including indirectly through another hook's captured runtime.
2169    pub fn detached_for_note_search_ann_provider(&self) -> Self {
2170        let mut detached = self.clone();
2171        detached.note_search_ann_provider = Arc::new(RwLock::new(None));
2172        detached.note_mutation_hook = Arc::new(RwLock::new(None));
2173        detached.entity_type_validator = Arc::new(RwLock::new(None));
2174        detached.note_write_validator = Arc::new(RwLock::new(None));
2175        detached.entity_kind_hooks = Arc::new(RwLock::new(Vec::new()));
2176        detached.fusion_executors = Arc::new(RwLock::new(HashMap::new()));
2177        detached
2178    }
2179
2180    /// Install the memory pack's note-search provider only on its own opened
2181    /// backend. A same-named but separate store keeps the exact search route.
2182    pub fn install_note_search_ann_provider(&self, provider: Arc<dyn NoteSearchAnnProvider>) {
2183        if !provider.serves_backend(self) {
2184            return;
2185        }
2186        if let Ok(mut guard) = self.note_search_ann_provider.write() {
2187            *guard = Some(provider);
2188        }
2189    }
2190
2191    pub(crate) fn note_search_ann_provider(
2192        &self,
2193    ) -> RuntimeResult<Option<Arc<dyn NoteSearchAnnProvider>>> {
2194        self.note_search_ann_provider
2195            .read()
2196            .map(|guard| {
2197                guard
2198                    .as_ref()
2199                    .filter(|provider| provider.serves_backend(self))
2200                    .cloned()
2201            })
2202            .map_err(|_| RuntimeError::Internal("note-search ANN provider lock poisoned".into()))
2203    }
2204
2205    /// Install the pack-aggregated entity-kind update hooks (issue #2943).
2206    ///
2207    /// Called by the transport after the `VerbRegistry` is built, same
2208    /// timing as [`install_kind_registry`](Self::install_kind_registry) —
2209    /// pass `registry.entity_kind_hooks()`. Idempotent: a later call
2210    /// replaces the set.
2211    pub fn install_entity_kind_hooks(&self, hooks: EntityKindHooks) {
2212        if let Ok(mut guard) = self.entity_kind_hooks.write() {
2213            *guard = hooks;
2214        }
2215    }
2216
2217    /// The installed `KindHook` for entity `kind`, if its owning pack
2218    /// registered one via [`install_entity_kind_hooks`](Self::install_entity_kind_hooks).
2219    ///
2220    /// `None` before the transport installs the aggregate (bare runtime) or
2221    /// when no pack registered a hook for this entity kind — the caller
2222    /// treats this the same as a hook whose `validate_entity_update`
2223    /// inherited the trait's `Ok(())` default.
2224    pub(crate) fn entity_kind_hook(&self, kind: &str) -> Option<Arc<dyn KindHook>> {
2225        self.entity_kind_hooks.read().ok().and_then(|guard| {
2226            guard
2227                .iter()
2228                .find(|(k, _)| k == kind)
2229                .map(|(_, hook)| hook.clone())
2230        })
2231    }
2232
2233    /// Install a pack-owned note-write validator.
2234    ///
2235    /// Called during pack registration (`PackRuntime::register_note_write_validator`)
2236    /// so that the covered note-write sites carrying caller-supplied
2237    /// `properties` derive the owning pack's identity properties from the
2238    /// authorization token, closing the gap where a direct Rust caller, the
2239    /// generic `create` verb, or the proposal-apply path (which dispatches no
2240    /// pack hooks) writes them unchecked. Single-slot semantics, same as
2241    /// [`install_note_mutation_hook`](Self::install_note_mutation_hook): a
2242    /// second installing pack overwrites the first, so a validator must return
2243    /// kinds it does not own unchanged.
2244    ///
2245    /// Covered sites — each calls `derive_note_write_properties`
2246    /// before the write: `create_note_inner` (`operations.rs`, the generic
2247    /// `create` verb funnel and every other public `create_note*` variant),
2248    /// `atomic_prepare::prepare_add_note` (the proposal-apply add-note path),
2249    /// and `atomic_message::create_notes_atomic_with_report` (the atomic
2250    /// multi-note writer).
2251    ///
2252    /// NOT covered by this validator: `try_create_note` (`operations.rs`).
2253    /// `try_create_note` is deliberately excluded — its only caller path is
2254    /// `comm.ingest`, where `properties.from_actor` is the external transport
2255    /// sender named by the `from` parameter, not the authenticated caller,
2256    /// and where transport-owned quarantine/channel properties are
2257    /// legitimately established. Running the generic validator there would
2258    /// stamp every inbound message as the ingesting daemon and reject the
2259    /// evidence the trusted ingest handler just derived. `try_create_note`
2260    /// instead runs its own narrower reserved-transport-property check
2261    /// inline (`operations.rs`'s `try_create_note_impl`), which allows the
2262    /// `message`-kind transport properties only when called through
2263    /// [`Self::try_create_note_as_trusted_ingest`] with a
2264    /// [`crate::pack::ChannelIngestCapability`].
2265    ///
2266    /// The `NoteStore` returned by [`notes`](Self::notes) is covered by a
2267    /// different, narrower mechanism: it is wrapped in
2268    /// `note_store_guard::PolicyEnforcingNoteStore`, which refuses
2269    /// `upsert_note` / `upsert_notes` / `try_insert_note` /
2270    /// `replace_note_if_unchanged` calls that would write a `kind = "message"`
2271    /// note carrying `quarantined` / `channel_kind` / `channel_slug`, and
2272    /// refuses `set_note_property` / `try_patch_note_property` /
2273    /// `patch_note_property_atomic` / `update_note_properties` calls that
2274    /// would patch any of those keys onto any note — unconditionally, since
2275    /// that public accessor has no way to see a trust decision. `try_create_note_impl` itself reaches storage through
2276    /// `Self::raw_notes`, the unwrapped accessor, so its own inline check
2277    /// (which can legitimately allow those properties for trusted ingest)
2278    /// is not double-enforced or contradicted by the wrapper.
2279    /// Register a pack-defined custom fusion strategy under `name` (ADR-012).
2280    ///
2281    /// Unlike `install_entity_type_validator`/`install_note_mutation_hook`,
2282    /// this slot is keyed rather than single-occupancy: multiple packs each
2283    /// register their own named strategy, and a second registration under an
2284    /// already-used `name` replaces the first. Looked up by
2285    /// `FusionStrategy::Custom { name, .. }` at the hybrid-search dispatch
2286    /// boundary in `crate::fusion`; an unregistered name fails closed with
2287    /// `RuntimeError::UnknownFusionStrategy` rather than silently falling
2288    /// back to RRF.
2289    pub fn register_fusion_strategy(
2290        &self,
2291        name: impl Into<String>,
2292        executor: Arc<dyn crate::fusion::FusionExecutor>,
2293    ) {
2294        if let Ok(mut guard) = self.fusion_executors.write() {
2295            guard.insert(name.into(), executor);
2296        }
2297    }
2298
2299    /// Resolve a registered custom fusion executor by name.
2300    ///
2301    /// Returns `RuntimeError::UnknownFusionStrategy` when no pack has
2302    /// registered `name` — callers must invoke this before any
2303    /// empty-input/zero-limit short circuit so a misconfigured name errors
2304    /// on every call, including zero-result ones.
2305    pub(crate) fn fusion_executor(
2306        &self,
2307        name: &str,
2308    ) -> RuntimeResult<Arc<dyn crate::fusion::FusionExecutor>> {
2309        let guard = self
2310            .fusion_executors
2311            .read()
2312            .map_err(|_| RuntimeError::Internal("fusion executor registry lock poisoned".into()))?;
2313        guard
2314            .get(name)
2315            .cloned()
2316            .ok_or_else(|| RuntimeError::UnknownFusionStrategy(name.to_string()))
2317    }
2318
2319    pub fn install_note_write_validator(&self, f: NoteWriteValidatorFn) {
2320        if let Ok(mut guard) = self.note_write_validator.write() {
2321            *guard = Some(f);
2322        }
2323    }
2324
2325    /// Whether a note-write validator is installed on this runtime.
2326    ///
2327    /// Exists so a transport's own tests can assert, per boot path, that the
2328    /// documented startup sequence actually filled the slot. A missing install
2329    /// fails open and silently — an empty slot passes caller-supplied
2330    /// properties straight through, which no write site can distinguish from a
2331    /// validator that approved them — so occupancy is asserted, never assumed.
2332    pub fn has_note_write_validator(&self) -> bool {
2333        self.note_write_validator
2334            .read()
2335            .map(|g| g.is_some())
2336            .unwrap_or(false)
2337    }
2338
2339    /// Run caller-supplied note `properties` through the installed note-write
2340    /// validator, returning the properties to store.
2341    ///
2342    /// Returns them unchanged when no validator is installed (bare runtime).
2343    pub(crate) fn derive_note_write_properties(
2344        &self,
2345        kind: &str,
2346        token: &NamespaceToken,
2347        properties: Option<serde_json::Value>,
2348    ) -> RuntimeResult<Option<serde_json::Value>> {
2349        let validator = self
2350            .note_write_validator
2351            .read()
2352            .map_err(|_| RuntimeError::Internal("note write validator lock poisoned".into()))?
2353            .clone();
2354        match validator {
2355            None => Ok(properties),
2356            Some(validate) => validate(kind, &token.actor().id, properties),
2357        }
2358    }
2359
2360    /// Invoke the pack-installed note-mutation hook, if any.
2361    ///
2362    /// `kind` is the note's `kind` string (e.g. `"memory"`); `id` is the
2363    /// note's UUID. No-op when no hook is installed (bare runtime, or no
2364    /// pack cares). Errors inside the hook are the hook's own concern to
2365    /// handle/log — this call site cannot propagate a failure without
2366    /// changing `update_note`/`delete_note`'s already-committed success
2367    /// return value.
2368    pub(crate) async fn fire_note_mutation_hook(&self, kind: &str, id: uuid::Uuid) {
2369        let hook = self
2370            .note_mutation_hook
2371            .read()
2372            .ok()
2373            .and_then(|guard| guard.clone());
2374        if let Some(hook) = hook {
2375            hook(kind.to_string(), id).await;
2376        }
2377    }
2378
2379    /// Snapshot of currently-installed pack edge rules.
2380    ///
2381    /// This is the same composed rule set `validate_edge_relation_endpoints`
2382    /// consults via `pack_rule_allows` when accepting/rejecting an edge. Public
2383    /// so pack-layer error-hint code (e.g. `khive-pack-kg`'s
2384    /// `valid_relations_for_entity_pair`) can derive hints from the exact
2385    /// source the validator uses, rather than maintaining a separate
2386    /// hand-authored table that can drift out of sync.
2387    pub fn pack_edge_rules(&self) -> Vec<EdgeEndpointRule> {
2388        self.edge_rules
2389            .read()
2390            .map(|g| g.clone())
2391            .unwrap_or_default()
2392    }
2393
2394    /// Borrow the installed pack edge rules for a synchronous calculation.
2395    pub(crate) fn with_pack_edge_rules<T>(&self, f: impl FnOnce(&[EdgeEndpointRule]) -> T) -> T {
2396        match self.edge_rules.read() {
2397            Ok(rules) => f(&rules),
2398            Err(_) => f(&[]),
2399        }
2400    }
2401
2402    /// Return the name of the default embedding model (empty string if none configured).
2403    pub fn default_embedder_name(&self) -> &str {
2404        self.default_embedder_name.as_ref()
2405    }
2406
2407    /// Resolve a model name (or `None` for the default) to an `EmbeddingModel`.
2408    ///
2409    /// Returns `UnknownModel` if the name is not in the registry, or
2410    /// `Unconfigured` if `None` is passed and no default model is set.
2411    pub fn resolve_embedding_model(&self, name: Option<&str>) -> RuntimeResult<EmbeddingModel> {
2412        let model = match name {
2413            Some(raw) => parse_embedding_model_alias(raw)
2414                .ok_or_else(|| crate::RuntimeError::UnknownModel(raw.to_string()))?,
2415            None => self
2416                .config
2417                .embedding_model
2418                .ok_or_else(|| crate::RuntimeError::Unconfigured("embedding_model".into()))?,
2419        };
2420        let key = model.to_string();
2421        if request_excludes_embedder(&key) {
2422            return Err(crate::RuntimeError::UnknownModel(
2423                name.unwrap_or_else(|| self.default_embedder_name())
2424                    .to_string(),
2425            ));
2426        }
2427        let contains = self
2428            .embedder_registry
2429            .read()
2430            .map(|reg| reg.contains(&key))
2431            .unwrap_or(false);
2432        if contains {
2433            Ok(model)
2434        } else {
2435            Err(crate::RuntimeError::UnknownModel(
2436                name.unwrap_or_else(|| self.default_embedder_name())
2437                    .to_string(),
2438            ))
2439        }
2440    }
2441
2442    /// Names of all registered embedding models in this runtime.
2443    ///
2444    /// Includes both built-in lattice models and any custom embedders
2445    /// registered by packs via [`register_embedder`](Self::register_embedder).
2446    /// Useful for operations that must touch every model's storage (e.g.,
2447    /// scoped vector deletion on note delete). The default model is included.
2448    pub fn registered_embedding_model_names(&self) -> Vec<String> {
2449        self.embedder_registry
2450            .read()
2451            .map(|reg| {
2452                reg.names()
2453                    .into_iter()
2454                    .filter(|name| !request_excludes_embedder(name))
2455                    .collect()
2456            })
2457            .unwrap_or_default()
2458    }
2459
2460    /// Get the lazily-initialized embedding service for the named model.
2461    ///
2462    /// Accepts both built-in lattice model names (e.g. `"all-minilm-l6-v2"`,
2463    /// `"paraphrase"`) and custom provider names registered via
2464    /// [`register_embedder`](Self::register_embedder).
2465    ///
2466    /// For lattice model names, aliases (e.g. `"paraphrase"`) are resolved to
2467    /// their canonical key before looking up the registry. For custom providers
2468    /// the name must match exactly as supplied during registration.
2469    ///
2470    /// First call for any name loads the underlying service (cold start cost);
2471    /// subsequent calls are cheap (registry caches the `Arc`).
2472    pub async fn embedder(&self, name: &str) -> RuntimeResult<Arc<dyn EmbeddingService>> {
2473        Ok(self.embedder_inner(name, None).await?.0)
2474    }
2475
2476    pub(crate) async fn embedder_with_token(
2477        &self,
2478        token: &NamespaceToken,
2479        name: &str,
2480    ) -> RuntimeResult<Arc<dyn EmbeddingService>> {
2481        Ok(self.embedder_inner(name, Some(token)).await?.0)
2482    }
2483
2484    /// Resolve the service and its document-preparation attestation from the
2485    /// same registry entry. A pack can replace a built-in name while a cold
2486    /// service is initializing, so a second registry lookup would be unsafe.
2487    pub(crate) async fn embedder_with_input_attestation(
2488        &self,
2489        name: &str,
2490        token: Option<&NamespaceToken>,
2491    ) -> RuntimeResult<(Arc<dyn EmbeddingService>, bool)> {
2492        self.embedder_inner(name, token).await
2493    }
2494
2495    /// Register a custom embedding provider with this runtime.
2496    ///
2497    /// The provider is added to the shared [`EmbedderRegistry`] so all clones
2498    /// of this runtime see the new provider immediately. If a provider with the
2499    /// same name already exists it is replaced (last-writer wins — see
2500    /// [`crate::EmbedderRegistry::register`] for the rationale).
2501    ///
2502    /// Packs should call this from [`crate::PackRuntime::register_embedders`] (the
2503    /// hook is invoked by the transport during pack initialisation, before the
2504    /// first verb dispatch).
2505    ///
2506    /// [`EmbedderRegistry`]: crate::embedder_registry::EmbedderRegistry
2507    pub fn register_embedder(
2508        &self,
2509        provider: impl crate::embedder_registry::EmbedderProvider + 'static,
2510    ) {
2511        if let Ok(mut registry) = self.embedder_registry.write() {
2512            registry.register(provider);
2513        } else {
2514            tracing::warn!(
2515                "embedder registry lock poisoned — embedder {} not registered",
2516                std::any::type_name::<dyn crate::embedder_registry::EmbedderProvider>()
2517            );
2518        }
2519    }
2520
2521    /// Install a deterministic backend for exact-input provenance tests.
2522    /// The test adapter, not the supplied backend, owns lattice passage
2523    /// prefixing; this API is absent unless `test-internals` is enabled.
2524    #[cfg(feature = "test-internals")]
2525    pub fn register_test_audited_embedder(
2526        &self,
2527        model: EmbeddingModel,
2528        provider: impl crate::embedder_registry::EmbedderProvider + 'static,
2529    ) {
2530        self.embedder_registry
2531            .write()
2532            .expect("test embedder registry lock")
2533            .register_test_audited(model, provider);
2534    }
2535
2536    /// List registered embedding models via `SqlAccess`, routing through the
2537    /// existing connection pool rather than opening a fresh `Connection` per call.
2538    ///
2539    /// Optionally filter by `engine_name`. Returns an empty vec when the
2540    /// `_embedding_models` table does not yet exist (e.g. no migrations have run
2541    /// or no models have been registered). All other SQL errors are propagated.
2542    pub async fn list_embedding_models(
2543        &self,
2544        engine_filter: Option<&str>,
2545    ) -> RuntimeResult<Vec<khive_db::EmbeddingModelRegistryRecord>> {
2546        use khive_storage::{SqlStatement, SqlValue};
2547
2548        let (sql_text, params) = if let Some(engine) = engine_filter {
2549            (
2550                "SELECT engine_name, model_id, key_version, dim, status, \
2551                 activated_at, superseded_at \
2552                 FROM _embedding_models WHERE engine_name = ?1 \
2553                 ORDER BY engine_name, activated_at IS NULL, activated_at"
2554                    .to_string(),
2555                vec![SqlValue::Text(engine.to_string())],
2556            )
2557        } else {
2558            (
2559                "SELECT engine_name, model_id, key_version, dim, status, \
2560                 activated_at, superseded_at \
2561                 FROM _embedding_models \
2562                 ORDER BY engine_name, activated_at IS NULL, activated_at"
2563                    .to_string(),
2564                vec![],
2565            )
2566        };
2567
2568        let stmt = SqlStatement {
2569            sql: sql_text,
2570            params,
2571            label: Some("list_embedding_models".into()),
2572        };
2573
2574        let mut reader = self
2575            .sql()
2576            .reader()
2577            .await
2578            .map_err(crate::RuntimeError::Storage)?;
2579
2580        let rows = match reader.query_all(stmt).await {
2581            Ok(rows) => rows,
2582            Err(e) if e.to_string().contains("no such table: _embedding_models") => {
2583                return Ok(Vec::new())
2584            }
2585            Err(e) => return Err(crate::RuntimeError::Storage(e)),
2586        };
2587
2588        let mut records = Vec::with_capacity(rows.len());
2589        for row in rows {
2590            macro_rules! required_text {
2591                ($col:expr) => {
2592                    match row.get($col) {
2593                        Some(SqlValue::Text(s)) => s.clone(),
2594                        other => {
2595                            tracing::warn!(column = $col, value = ?other, "skipping registry row: unexpected type");
2596                            continue;
2597                        }
2598                    }
2599                };
2600            }
2601            let engine_name = required_text!("engine_name");
2602            let model_id = required_text!("model_id");
2603            let key_version = required_text!("key_version");
2604            let dimensions = match row.get("dim") {
2605                Some(SqlValue::Integer(n)) => match u32::try_from(*n) {
2606                    Ok(d) => d,
2607                    Err(_) => {
2608                        tracing::warn!(dim = n, "skipping registry row: dim out of u32 range");
2609                        continue;
2610                    }
2611                },
2612                other => {
2613                    tracing::warn!(column = "dim", value = ?other, "skipping registry row: unexpected type");
2614                    continue;
2615                }
2616            };
2617            let status = required_text!("status");
2618            let activated_at = match row.get("activated_at") {
2619                Some(SqlValue::Integer(n)) => Some(*n),
2620                _ => None,
2621            };
2622            let superseded_at = match row.get("superseded_at") {
2623                Some(SqlValue::Integer(n)) => Some(*n),
2624                _ => None,
2625            };
2626            records.push(khive_db::EmbeddingModelRegistryRecord {
2627                engine_name,
2628                model_id,
2629                key_version,
2630                dimensions,
2631                status,
2632                activated_at,
2633                superseded_at,
2634            });
2635        }
2636
2637        Ok(records)
2638    }
2639}
2640
2641fn vector_dimensions_from_ddl(ddl: &str) -> Option<usize> {
2642    let lower = ddl.to_ascii_lowercase();
2643    let suffix = lower.split_once("embedding float[")?.1;
2644    let dimension = suffix.split_once(']')?.0;
2645    if dimension.is_empty() || !dimension.bytes().all(|byte| byte.is_ascii_digit()) {
2646        return None;
2647    }
2648    dimension.parse().ok()
2649}
2650
2651// IN-CRATE TEST JUSTIFICATION: tests here cover KhiveRuntime construction helpers
2652// (in-memory backend wiring, NamespaceToken::for_namespace) that are
2653// pub(crate)-only and cannot be called from the integration test crate.
2654#[cfg(test)]
2655#[path = "runtime_tests.rs"]
2656mod tests;