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(¤t.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(¤t_store, &store) || Arc::ptr_eq(¤t.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;