1use 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 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 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
68pub 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
79pub 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
102pub type EntityTypeValidatorFn =
108 Arc<dyn Fn(&str, Option<&str>) -> Result<Option<String>, RuntimeError> + Send + Sync>;
109
110pub type EntityKindHooks = Vec<(String, Arc<dyn KindHook>)>;
117
118pub 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
134pub 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#[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 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#[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#[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#[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 named_vector_stores: Arc<RwLock<HashMap<String, CachedNamedVectorStores>>>,
309 core_named_vector_stores: Option<Arc<RwLock<HashMap<String, CachedNamedVectorStores>>>>,
312 core_backend: Option<Arc<StorageBackend>>,
316 config: RuntimeConfig,
317 outbound_email_policy: crate::OutboundEmailPolicy,
318 declared_backend_db_paths: Arc<[PathBuf]>,
321 diagnostic_backends: Arc<[OpenedDiagnosticBackend]>,
324 late_diagnostic_backends: Arc<Mutex<Vec<LateOpenedDiagnosticBackend>>>,
327 ann_fresh_tail_enabled: bool,
331 embedder_registry: Arc<std::sync::RwLock<crate::embedder_registry::EmbedderRegistry>>,
338 default_embedder_name: Arc<str>,
339 core_embedders: Option<CoreEmbedderState>,
346 edge_rules: Arc<RwLock<Vec<EdgeEndpointRule>>>,
350 valid_entity_kinds: Arc<RwLock<Vec<String>>>,
358 valid_note_kinds: Arc<RwLock<Vec<NoteKindEntry>>>,
359 entity_type_validator: Arc<RwLock<Option<EntityTypeValidatorFn>>>,
366 note_mutation_hook: Arc<RwLock<Option<NoteMutationHookFn>>>,
378 note_search_ann_provider: Arc<RwLock<Option<Arc<dyn NoteSearchAnnProvider>>>>,
381 note_write_validator: Arc<RwLock<Option<NoteWriteValidatorFn>>>,
389 entity_kind_hooks: Arc<RwLock<EntityKindHooks>>,
403 pack_owned_note_kinds: Arc<RwLock<Vec<String>>>,
411 blob_hydrator: Arc<OnceLock<Arc<crate::blob::BlobHydrator>>>,
418 fusion_executors: Arc<RwLock<HashMap<String, Arc<dyn crate::fusion::FusionExecutor>>>>,
429}
430
431impl KhiveRuntime {
432 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 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 #[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 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 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 #[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 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 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 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 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 pub fn core(&self) -> KhiveRuntime {
764 match &self.core_backend {
765 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 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 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 pub fn backend_id(&self) -> &BackendId {
857 &self.config.backend_id
858 }
859
860 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 pub fn with_diagnostic_backends(mut self, backends: Arc<[OpenedDiagnosticBackend]>) -> Self {
886 self.diagnostic_backends = backends;
887 self
888 }
889
890 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 pub fn vector_arm_selected(&self) -> bool {
981 self.config.embedding_model.is_some()
982 }
983
984 pub fn backend(&self) -> &StorageBackend {
994 &self.backend
995 }
996
997 pub fn is_read_only(&self) -> bool {
1000 self.backend.is_read_only()
1001 }
1002
1003 pub fn backend_data_dir(&self) -> Option<std::path::PathBuf> {
1006 self.backend.data_dir()
1007 }
1008
1009 pub fn backend_ann_root(&self) -> Option<std::path::PathBuf> {
1014 self.backend.ann_root()
1015 }
1016
1017 pub async fn db_diagnostics(&self) -> RuntimeResult<khive_db::diagnostics::DbDiagnostics> {
1031 self.db_diagnostics_with_audit_metrics(None).await
1039 }
1040
1041 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 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 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 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 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 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 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 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 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 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 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 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 pub fn sql(&self) -> Arc<dyn SqlAccess> {
1296 self.backend.sql()
1297 }
1298
1299 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 pub fn blob_hydrator(&self) -> Option<Arc<crate::blob::BlobHydrator>> {
1935 self.blob_hydrator.get().cloned()
1936 }
1937
1938 pub fn blob_store(&self) -> Option<Arc<dyn khive_storage::BlobStore>> {
1943 self.blob_hydrator.get().map(|hydrator| hydrator.store())
1944 }
1945
1946 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 pub fn default_embedder_name(&self) -> &str {
2404 self.default_embedder_name.as_ref()
2405 }
2406
2407 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 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 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 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 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 #[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 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#[cfg(test)]
2655mod tests {
2656 use super::*;
2657 use khive_gate::GateRef;
2658 use serial_test::serial;
2659
2660 #[cfg(target_os = "macos")]
2661 #[test]
2662 fn in_process_runtime_tests_have_4096_open_file_slots() {
2663 let _runtime = KhiveRuntime::memory().expect("test runtime");
2664 let mut limits = libc::rlimit {
2665 rlim_cur: 0,
2666 rlim_max: 0,
2667 };
2668 assert_eq!(
2670 unsafe { libc::getrlimit(libc::RLIMIT_NOFILE, &mut limits) },
2671 0
2672 );
2673 assert!(
2674 limits.rlim_cur >= IN_PROCESS_TEST_NOFILE_LIMIT,
2675 "a parallel runtime suite needs at least 4096 open-file slots"
2676 );
2677 }
2678
2679 fn test_blob_hydrator() -> (tempfile::TempDir, Arc<crate::BlobHydrator>) {
2680 let root = tempfile::tempdir().expect("blob root");
2681 let store = Arc::new(
2682 khive_db::stores::blob::FsBlobStore::new(root.path().to_path_buf(), 0)
2683 .expect("fs blob store"),
2684 );
2685 let hydrator = Arc::new(
2686 crate::BlobHydrator::new(store, crate::DEFAULT_BLOB_HYDRATION_BYTES)
2687 .expect("blob hydrator"),
2688 );
2689 (root, hydrator)
2690 }
2691
2692 #[test]
2693 fn memory_runtime_creates_successfully() {
2694 let rt = KhiveRuntime::memory().expect("memory runtime should create");
2695 assert!(rt.config().db_path.is_none());
2696 }
2697
2698 include!("runtime_blob_hydrator_tests.rs");
2699 include!("runtime_volume_lock_dir_tests.rs");
2700
2701 #[test]
2702 fn fresh_tail_policy_is_instance_scoped_and_clone_stable() {
2703 let enabled = KhiveRuntime::memory()
2704 .expect("enabled memory runtime")
2705 .with_ann_fresh_tail_enabled(true);
2706 let disabled = KhiveRuntime::memory()
2707 .expect("disabled memory runtime")
2708 .with_ann_fresh_tail_enabled(false);
2709
2710 assert!(enabled.ann_fresh_tail_enabled());
2711 assert!(enabled.clone().ann_fresh_tail_enabled());
2712 assert!(!disabled.ann_fresh_tail_enabled());
2713 assert!(!disabled.clone().ann_fresh_tail_enabled());
2714 }
2715
2716 include!("runtime_db_diagnostics_counter_tests.rs");
2717
2718 #[test]
2719 fn diagnostics_tracks_late_events_sidecar_without_retaining_its_pool() {
2720 let dir = tempfile::tempdir().expect("diagnostics database directory");
2721 let guard = crate::events_split::TestRegistryGuard::new(dir.path());
2722 let sidecar_path = dir.path().join("main.db.events.db");
2723 let mut config = RuntimeConfig::no_embeddings();
2724 config.db_path = Some(dir.path().join("main.db"));
2725 config.events_split = Some(crate::events_split::EventsSplitConfig {
2726 db_path: sidecar_path.clone(),
2727 socket_path: None,
2728 });
2729 let runtime = KhiveRuntime::new_for_test(config).expect("main runtime");
2730 let clone = runtime.clone();
2731
2732 assert!(!sidecar_path.exists());
2733 assert_eq!(
2734 runtime
2735 .diagnostic_backends()
2736 .iter()
2737 .filter(|backend| backend.canonical_path.as_deref() == Some(sidecar_path.as_path()))
2738 .count(),
2739 0,
2740 "an unopened sidecar must not appear in diagnostics"
2741 );
2742 assert!(runtime
2743 .events_sidecar_sql_read_only()
2744 .expect("missing sidecar lookup")
2745 .is_none());
2746 assert!(
2747 !sidecar_path.exists(),
2748 "inspection must not create the sidecar"
2749 );
2750
2751 std::fs::File::create(&sidecar_path).expect("preexisting events sidecar");
2752 let sql = runtime
2753 .events_sidecar_sql_read_only()
2754 .expect("sidecar SQL lookup")
2755 .expect("preexisting sidecar opens");
2756 let canonical_sidecar = sidecar_path.canonicalize().expect("canonical sidecar path");
2757 let snapshot = clone.diagnostic_backends();
2758 let entries: Vec<_> = snapshot
2759 .iter()
2760 .filter(|backend| {
2761 backend.canonical_path.as_deref() == Some(canonical_sidecar.as_path())
2762 })
2763 .collect();
2764 assert_eq!(entries.len(), 1, "sidecar opens after runtime composition");
2765 assert_eq!(entries[0].backend_names, vec!["events".to_string()]);
2766 let weak_sidecar_pool = Arc::downgrade(&entries[0].pool);
2767
2768 let event_store = runtime
2769 .raw_events_for_namespace("local")
2770 .expect("direct events store");
2771 let repeated = runtime.diagnostic_backends();
2772 assert_eq!(
2773 repeated
2774 .iter()
2775 .filter(|backend| backend.canonical_path.as_deref()
2776 == Some(canonical_sidecar.as_path()))
2777 .count(),
2778 1,
2779 "the SQL and event-store paths must report one physical file"
2780 );
2781
2782 let boot_alias =
2783 StorageBackend::sqlite_for_test(dir.path().join(".").join("main.db.events.db"))
2784 .expect("second pool for canonical-file alias");
2785 let boot_alias_pool = boot_alias.pool_arc();
2786 let main_pool = runtime.backend().pool_arc();
2787 let with_boot_alias = runtime.clone().with_diagnostic_backends(
2788 vec![
2789 OpenedDiagnosticBackend {
2790 backend_names: vec!["main".into()],
2791 canonical_path: main_pool.canonical_path().map(PathBuf::from),
2792 pool: main_pool,
2793 },
2794 OpenedDiagnosticBackend {
2795 backend_names: vec!["boot_alias".into()],
2796 canonical_path: boot_alias_pool.canonical_path().map(PathBuf::from),
2797 pool: Arc::clone(&boot_alias_pool),
2798 },
2799 ]
2800 .into(),
2801 );
2802 let with_boot_alias_snapshot = with_boot_alias.diagnostic_backends();
2803 let merged: Vec<_> = with_boot_alias_snapshot
2804 .iter()
2805 .filter(|backend| {
2806 backend.canonical_path.as_deref() == Some(canonical_sidecar.as_path())
2807 })
2808 .collect();
2809 assert_eq!(merged.len(), 1, "two pools over one canonical file");
2810 assert_eq!(
2811 merged[0].backend_names,
2812 vec!["boot_alias".to_string(), "events".to_string()]
2813 );
2814 assert!(Arc::ptr_eq(&merged[0].pool, &boot_alias_pool));
2815
2816 drop(with_boot_alias_snapshot);
2817 drop(with_boot_alias);
2818 drop(boot_alias_pool);
2819 drop(boot_alias);
2820 drop(repeated);
2821 drop(snapshot);
2822 drop(guard);
2823 assert!(weak_sidecar_pool.upgrade().is_some());
2824 assert!(runtime.diagnostic_backends().iter().any(|backend| {
2825 backend.canonical_path.as_deref() == Some(canonical_sidecar.as_path())
2826 }));
2827
2828 drop(event_store);
2829 drop(sql);
2830 assert!(weak_sidecar_pool.upgrade().is_none());
2831 assert_eq!(
2832 runtime
2833 .diagnostic_backends()
2834 .iter()
2835 .filter(|backend| backend.canonical_path.as_deref()
2836 == Some(canonical_sidecar.as_path()))
2837 .count(),
2838 0,
2839 "diagnostics must not retain a pool after its owner releases it"
2840 );
2841 }
2842
2843 #[test]
2844 fn diagnostics_tracks_direct_events_open_after_runtime_composition() {
2845 let dir = tempfile::tempdir().expect("diagnostics database directory");
2846 let guard = crate::events_split::TestRegistryGuard::new(dir.path());
2847 let sidecar_path = dir.path().join("main.db.events.db");
2848 let mut config = RuntimeConfig::no_embeddings();
2849 config.db_path = Some(dir.path().join("main.db"));
2850 config.events_split = Some(crate::events_split::EventsSplitConfig {
2851 db_path: sidecar_path.clone(),
2852 socket_path: None,
2853 });
2854 let runtime = KhiveRuntime::new_for_test(config.clone()).expect("main runtime");
2855
2856 assert!(!sidecar_path.exists());
2857 let events = runtime
2858 .raw_events_for_namespace("local")
2859 .expect("direct event store opens sidecar");
2860 let canonical_sidecar = sidecar_path.canonicalize().expect("canonical sidecar path");
2861 let opened = runtime.clone().diagnostic_backends();
2862 let matching: Vec<_> = opened
2863 .iter()
2864 .filter(|backend| {
2865 backend.canonical_path.as_deref() == Some(canonical_sidecar.as_path())
2866 })
2867 .collect();
2868 assert_eq!(matching.len(), 1);
2869 assert_eq!(matching[0].backend_names, vec!["events".to_string()]);
2870 let weak_sidecar_pool = Arc::downgrade(&matching[0].pool);
2871
2872 let pack_runtime = KhiveRuntime::memory()
2873 .expect("pack runtime")
2874 .with_diagnostic_observer_from(&runtime);
2875 assert!(pack_runtime.diagnostic_backends().iter().any(|backend| {
2876 backend.canonical_path.as_deref() == Some(canonical_sidecar.as_path())
2877 }));
2878
2879 let unrelated = KhiveRuntime::new_for_test(config).expect("independent runtime");
2880 assert_eq!(
2881 unrelated.diagnostic_backends().len(),
2882 1,
2883 "an independent runtime must not inherit another runtime's events pool"
2884 );
2885
2886 drop(opened);
2887 drop(guard);
2888 assert!(weak_sidecar_pool.upgrade().is_some());
2889 assert!(runtime.diagnostic_backends().iter().any(|backend| {
2890 backend.canonical_path.as_deref() == Some(canonical_sidecar.as_path())
2891 }));
2892
2893 drop(events);
2894 assert!(weak_sidecar_pool.upgrade().is_none());
2895 assert!(runtime
2896 .diagnostic_backends()
2897 .iter()
2898 .all(|backend| backend.canonical_path.as_deref() != Some(canonical_sidecar.as_path())));
2899 }
2900
2901 #[test]
2902 fn backend_data_dir_returns_none_for_memory_backend() {
2903 let rt = KhiveRuntime::memory().expect("memory runtime");
2904 assert!(rt.backend_data_dir().is_none());
2905 }
2906
2907 #[test]
2908 fn backend_data_dir_returns_parent_dir_for_file_backend() {
2909 let dir = tempfile::tempdir().unwrap();
2910 let path = dir.path().join("test.db");
2911 let config = RuntimeConfig {
2912 web: Default::default(),
2913 telemetry: Default::default(),
2914 mounts: Vec::new(),
2915 brain: Default::default(),
2916 git_write: Default::default(),
2917 display_timezone: chrono_tz::Tz::UTC,
2918 events_split: None,
2919 db_path: Some(path),
2920 wal_ceiling_bytes: 0,
2921 wal_ceiling_configured_bytes: 0,
2922 wal_ceiling_source: khive_db::WalCeilingSource::Default,
2923 wal_ceiling_env_raw: None,
2924 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2925 default_namespace: Namespace::local(),
2926 embedding_model: None,
2927 additional_embedding_models: vec![],
2928 gate: Arc::new(AllowAllGate),
2929 packs: vec!["kg".to_string()],
2930 backend_id: BackendId::main(),
2931 brain_profile: None,
2932 visible_namespaces: vec![],
2933 allowed_outbound_namespaces: vec![],
2934 actor_id: None,
2935 exec: Default::default(),
2936 ..crate::RuntimeConfig::no_embeddings()
2937 };
2938 let rt = KhiveRuntime::new_for_test(config).expect("file runtime");
2939 let data_dir = rt
2940 .backend_data_dir()
2941 .expect("file backend must return Some");
2942 assert_eq!(data_dir, dir.path());
2943 }
2944
2945 #[tokio::test]
2951 async fn resolve_prefix_finds_sidecar_only_event() {
2952 let dir = tempfile::tempdir().unwrap();
2953 let _registry_guard = crate::events_split::TestRegistryGuard::new(dir.path());
2954 let sidecar_path = dir.path().join("main.db.events.db");
2955 let config = RuntimeConfig {
2956 web: Default::default(),
2957 telemetry: Default::default(),
2958 mounts: Vec::new(),
2959 brain: Default::default(),
2960 git_write: Default::default(),
2961 display_timezone: chrono_tz::Tz::UTC,
2962 events_split: Some(crate::events_split::EventsSplitConfig {
2963 db_path: sidecar_path.clone(),
2964 socket_path: None,
2965 }),
2966 db_path: Some(dir.path().join("main.db")),
2967 wal_ceiling_bytes: 0,
2968 wal_ceiling_configured_bytes: 0,
2969 wal_ceiling_source: khive_db::WalCeilingSource::Default,
2970 wal_ceiling_env_raw: None,
2971 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2972 default_namespace: Namespace::local(),
2973 embedding_model: None,
2974 additional_embedding_models: vec![],
2975 gate: Arc::new(AllowAllGate),
2976 packs: vec!["kg".to_string()],
2977 backend_id: BackendId::main(),
2978 brain_profile: None,
2979 visible_namespaces: vec![],
2980 allowed_outbound_namespaces: vec![],
2981 actor_id: None,
2982 exec: Default::default(),
2983 ..crate::RuntimeConfig::no_embeddings()
2984 };
2985 let rt = KhiveRuntime::new_for_test(config).expect("file runtime");
2986
2987 let event = khive_storage::Event::new(
2988 "local",
2989 "memory.recall",
2990 khive_types::EventKind::RecallExecuted,
2991 khive_types::SubstrateKind::Note,
2992 "agent:test",
2993 );
2994 let event_id = event.id;
2995 let prefix = event_id.to_string()[..8].to_string();
2996
2997 assert_eq!(
2998 rt.resolve_prefix_unfiltered(&prefix)
2999 .await
3000 .expect("pre-insert resolve"),
3001 None,
3002 "control: prefix must miss before the lane row exists"
3003 );
3004
3005 let lane = crate::events_split::direct_backend_with_policies(
3008 &sidecar_path,
3009 false,
3010 None,
3011 rt.events_wal_ceiling_policy(),
3012 Some(rt.events_disk_guard_policy().expect("events disk policy")),
3013 rt.events_volume_lock_dir(),
3014 )
3015 .expect("lane backend")
3016 .events_for_namespace("local")
3017 .expect("lane store");
3018 lane.append_event(event).await.expect("lane append");
3019
3020 assert_eq!(
3021 rt.resolve_prefix_unfiltered(&prefix)
3022 .await
3023 .expect("post-insert resolve"),
3024 Some(event_id),
3025 "a sidecar-only event id must resolve by hex prefix"
3026 );
3027 }
3028
3029 #[test]
3030 fn backend_data_dir_returns_none_for_from_backend_with_memory() {
3031 let backend = Arc::new(StorageBackend::memory().expect("memory backend"));
3032 let config = RuntimeConfig {
3033 web: Default::default(),
3034 telemetry: Default::default(),
3035 mounts: Vec::new(),
3036 brain: Default::default(),
3037 git_write: Default::default(),
3038 display_timezone: chrono_tz::Tz::UTC,
3039 events_split: None,
3040 db_path: None,
3041 wal_ceiling_bytes: 0,
3042 wal_ceiling_configured_bytes: 0,
3043 wal_ceiling_source: khive_db::WalCeilingSource::Default,
3044 wal_ceiling_env_raw: None,
3045 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3046 default_namespace: Namespace::local(),
3047 embedding_model: None,
3048 additional_embedding_models: vec![],
3049 gate: Arc::new(AllowAllGate),
3050 packs: vec!["kg".to_string()],
3051 backend_id: BackendId::main(),
3052 brain_profile: None,
3053 visible_namespaces: vec![],
3054 allowed_outbound_namespaces: vec![],
3055 actor_id: None,
3056 exec: Default::default(),
3057 ..crate::RuntimeConfig::no_embeddings()
3058 };
3059 let rt = KhiveRuntime::from_backend(backend, config);
3060 assert!(rt.backend_data_dir().is_none());
3061 }
3062
3063 #[test]
3064 fn direct_memory_runtime_rejects_nonzero_wal_ceiling() {
3065 let config = RuntimeConfig {
3066 db_path: None,
3067 wal_ceiling_bytes: 4152,
3068 wal_ceiling_configured_bytes: 4152,
3069 ..RuntimeConfig::no_embeddings()
3070 };
3071 let error = match KhiveRuntime::new(config) {
3072 Ok(_) => panic!("memory has no WAL extent"),
3073 Err(error) => error,
3074 };
3075 assert!(matches!(
3076 error,
3077 RuntimeError::Sqlite(khive_db::SqliteError::InvalidConfig(_))
3078 ));
3079 }
3080
3081 #[test]
3082 fn file_runtime_creates_successfully() {
3083 let dir = tempfile::tempdir().unwrap();
3084 let path = dir.path().join("test.db");
3085 let config = RuntimeConfig {
3086 web: Default::default(),
3087 telemetry: Default::default(),
3088 mounts: Vec::new(),
3089 brain: Default::default(),
3090 git_write: Default::default(),
3091 display_timezone: chrono_tz::Tz::UTC,
3092 events_split: None,
3093 db_path: Some(path.clone()),
3094 wal_ceiling_bytes: 0,
3095 wal_ceiling_configured_bytes: 0,
3096 wal_ceiling_source: khive_db::WalCeilingSource::Default,
3097 wal_ceiling_env_raw: None,
3098 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3099 default_namespace: Namespace::parse("test").unwrap(),
3100 embedding_model: None,
3101 additional_embedding_models: vec![],
3102 gate: Arc::new(AllowAllGate),
3103 packs: vec!["kg".to_string()],
3104 backend_id: BackendId::main(),
3105 brain_profile: None,
3106 visible_namespaces: vec![],
3107 allowed_outbound_namespaces: vec![],
3108 actor_id: None,
3109 exec: Default::default(),
3110 ..crate::RuntimeConfig::no_embeddings()
3111 };
3112 let rt = KhiveRuntime::new_for_test(config).expect("file runtime should create");
3113 assert!(path.exists());
3114 assert_eq!(rt.config().default_namespace.as_str(), "test");
3115 }
3116
3117 #[cfg(unix)]
3118 #[tokio::test]
3119 async fn normal_boot_detects_read_only_snapshot_and_skips_model_registration() {
3120 use std::os::unix::fs::PermissionsExt;
3121
3122 let dir = tempfile::tempdir().unwrap();
3123 let path = dir.path().join("read_only_runtime.db");
3124 let base = RuntimeConfig {
3125 web: Default::default(),
3126 telemetry: Default::default(),
3127 mounts: Vec::new(),
3128 brain: Default::default(),
3129 git_write: Default::default(),
3130 display_timezone: chrono_tz::Tz::UTC,
3131 events_split: None,
3132 db_path: Some(path.clone()),
3133 wal_ceiling_bytes: 0,
3134 wal_ceiling_configured_bytes: 0,
3135 wal_ceiling_source: khive_db::WalCeilingSource::Default,
3136 wal_ceiling_env_raw: None,
3137 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3138 default_namespace: Namespace::local(),
3139 embedding_model: None,
3140 additional_embedding_models: vec![],
3141 gate: Arc::new(AllowAllGate),
3142 packs: vec!["kg".to_string()],
3143 backend_id: BackendId::main(),
3144 brain_profile: None,
3145 visible_namespaces: vec![],
3146 allowed_outbound_namespaces: vec![],
3147 actor_id: None,
3148 exec: Default::default(),
3149 ..crate::RuntimeConfig::no_embeddings()
3150 };
3151 {
3152 let writable =
3153 KhiveRuntime::new_for_test(base.clone()).expect("create migrated snapshot");
3154 assert!(writable
3155 .list_embedding_models(None)
3156 .await
3157 .expect("registry query")
3158 .is_empty());
3159 }
3160
3161 let mut permissions = std::fs::metadata(&path).unwrap().permissions();
3162 permissions.set_mode(0o444);
3163 std::fs::set_permissions(&path, permissions).unwrap();
3164 khive_storage::test_support::freeze_snapshot_sidecars(&path);
3168
3169 let read_only_config = RuntimeConfig {
3170 embedding_model: Some(EmbeddingModel::AllMiniLmL6V2),
3171 ..base
3172 };
3173 let runtime = KhiveRuntime::new_for_test(read_only_config)
3174 .expect("read-only boot must validate instead of migrating/registering");
3175 assert!(runtime.is_read_only());
3176 assert_eq!(
3177 runtime.backend().pool().writer_acquisition_snapshot(),
3178 khive_db::pool::WriterAcquisitionSnapshot::default(),
3179 "the construction-inclusive acquisition baseline must stay at zero"
3180 );
3181 assert!(
3182 runtime
3183 .list_embedding_models(None)
3184 .await
3185 .expect("read-only registry query")
3186 .is_empty(),
3187 "configured models must remain in-memory only during read-only boot"
3188 );
3189 }
3190
3191 #[test]
3192 fn explicit_readonly_constructor_uses_read_only_pool_even_on_writable_file_mode() {
3193 let dir = tempfile::tempdir().unwrap();
3194 let path = dir.path().join("explicit_read_only_runtime.db");
3195 let config = RuntimeConfig {
3196 web: Default::default(),
3197 telemetry: Default::default(),
3198 mounts: Vec::new(),
3199 brain: Default::default(),
3200 git_write: Default::default(),
3201 display_timezone: chrono_tz::Tz::UTC,
3202 events_split: None,
3203 db_path: Some(path.clone()),
3204 wal_ceiling_bytes: 0,
3205 wal_ceiling_configured_bytes: 0,
3206 wal_ceiling_source: khive_db::WalCeilingSource::Default,
3207 wal_ceiling_env_raw: None,
3208 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3209 default_namespace: Namespace::local(),
3210 embedding_model: None,
3211 additional_embedding_models: vec![],
3212 gate: Arc::new(AllowAllGate),
3213 packs: vec!["kg".to_string()],
3214 backend_id: BackendId::main(),
3215 brain_profile: None,
3216 visible_namespaces: vec![],
3217 allowed_outbound_namespaces: vec![],
3218 actor_id: None,
3219 exec: Default::default(),
3220 ..crate::RuntimeConfig::no_embeddings()
3221 };
3222 KhiveRuntime::new_for_test(config.clone()).expect("create migrated database");
3223 #[cfg(unix)]
3224 khive_storage::test_support::freeze_snapshot_sidecars(&path);
3225
3226 let runtime = KhiveRuntime::new_readonly_for_test(config).expect("explicit read-only boot");
3227 assert!(runtime.is_read_only());
3228 assert_eq!(
3229 runtime.backend().pool().writer_acquisition_snapshot(),
3230 khive_db::pool::WriterAcquisitionSnapshot::default(),
3231 "explicit read-only construction must validate through a reader without ever \
3232 acquiring the writer"
3233 );
3234 }
3235
3236 #[derive(Debug)]
3240 struct ReadOnlyExtraGate {
3241 primary: &'static str,
3242 extra: &'static str,
3243 }
3244
3245 impl khive_gate::Gate for ReadOnlyExtraGate {
3246 fn check(
3247 &self,
3248 req: &khive_gate::GateRequest,
3249 ) -> Result<khive_gate::GateDecision, khive_gate::GateError> {
3250 let allowed = match req.verb.as_str() {
3251 "authorize" => req.namespace.as_str() == self.primary,
3252 "authorize.visible" => req.namespace.as_str() == self.extra,
3253 _ => false,
3254 };
3255 if allowed {
3256 Ok(khive_gate::GateDecision::allow())
3257 } else {
3258 Ok(khive_gate::GateDecision::Deny {
3259 reason: format!(
3260 "{} denied for namespace {:?}",
3261 req.verb,
3262 req.namespace.as_str()
3263 ),
3264 })
3265 }
3266 }
3267 }
3268
3269 #[test]
3270 fn authorize_with_visibility_allows_read_only_extra_namespace() {
3271 let primary = Namespace::parse("lambda:caller").expect("primary");
3272 let extra = Namespace::parse("lambda:read-only").expect("extra");
3273 let config = RuntimeConfig {
3274 db_path: None,
3275 packs: vec!["kg".to_string()],
3276 brain_profile: None,
3277 actor_id: None,
3278 gate: Arc::new(ReadOnlyExtraGate {
3279 primary: "lambda:caller",
3280 extra: "lambda:read-only",
3281 }),
3282 ..RuntimeConfig::no_embeddings()
3283 };
3284 let rt = KhiveRuntime::new(config).expect("memory runtime");
3285
3286 let token = rt
3287 .authorize_with_visibility(primary, vec![extra.clone()])
3288 .expect("Write on primary and Read on extra must mint");
3289 assert!(token.visible_namespaces().contains(&extra));
3290 }
3291
3292 #[derive(Debug)]
3296 struct DenyNamespaceGate {
3297 deny: &'static str,
3298 }
3299
3300 impl khive_gate::Gate for DenyNamespaceGate {
3301 fn check(
3302 &self,
3303 req: &khive_gate::GateRequest,
3304 ) -> Result<khive_gate::GateDecision, khive_gate::GateError> {
3305 if req.namespace.as_str() == self.deny {
3306 Ok(khive_gate::GateDecision::Deny {
3307 reason: "namespace denied by policy".to_string(),
3308 })
3309 } else {
3310 Ok(khive_gate::GateDecision::allow())
3311 }
3312 }
3313 }
3314
3315 #[test]
3316 fn authorize_with_visibility_denies_missing_extra_read() {
3317 let config = RuntimeConfig {
3318 db_path: None,
3319 packs: vec!["kg".to_string()],
3320 brain_profile: None,
3321 actor_id: None,
3322 gate: Arc::new(DenyNamespaceGate {
3323 deny: "lambda:secret",
3324 }),
3325 ..RuntimeConfig::no_embeddings()
3326 };
3327 let rt = KhiveRuntime::new(config).expect("memory runtime");
3328 let primary = Namespace::parse("lambda:caller").expect("primary");
3329 let denied = Namespace::parse("lambda:secret").expect("denied");
3330 let allowed = Namespace::parse("lambda:open").expect("allowed");
3331
3332 rt.authorize_with_visibility(primary.clone(), vec![allowed.clone()])
3335 .expect("mint with only allowed extras");
3336
3337 let err = rt
3338 .authorize_with_visibility(primary, vec![allowed, denied])
3339 .expect_err("a denied extra namespace must refuse the whole mint");
3340 let msg = err.to_string();
3341 assert!(
3342 msg.contains("lambda:secret"),
3343 "refusal must name the offending namespace: {msg}"
3344 );
3345 }
3346
3347 fn make_read_only_runtime() -> (tempfile::TempDir, KhiveRuntime) {
3350 let dir = tempfile::tempdir().unwrap();
3351 let path = dir.path().join("read_only_blob_seam.db");
3352 let config = RuntimeConfig {
3353 web: Default::default(),
3354 telemetry: Default::default(),
3355 mounts: Vec::new(),
3356 brain: Default::default(),
3357 git_write: Default::default(),
3358 display_timezone: chrono_tz::Tz::UTC,
3359 db_path: Some(path.clone()),
3360 wal_ceiling_bytes: 0,
3361 wal_ceiling_configured_bytes: 0,
3362 wal_ceiling_source: khive_db::WalCeilingSource::Default,
3363 wal_ceiling_env_raw: None,
3364 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3365 default_namespace: Namespace::local(),
3366 embedding_model: None,
3367 additional_embedding_models: vec![],
3368 gate: Arc::new(AllowAllGate),
3369 packs: vec!["kg".to_string()],
3370 backend_id: BackendId::main(),
3371 brain_profile: None,
3372 visible_namespaces: vec![],
3373 allowed_outbound_namespaces: vec![],
3374 actor_id: None,
3375 events_split: None,
3376 exec: Default::default(),
3377 ..crate::RuntimeConfig::no_embeddings()
3378 };
3379 KhiveRuntime::new_for_test(config.clone()).expect("create migrated database");
3380 #[cfg(unix)]
3381 khive_storage::test_support::freeze_snapshot_sidecars(&path);
3382 let runtime = KhiveRuntime::new_readonly_for_test(config).expect("read-only boot");
3383 assert!(runtime.is_read_only());
3384 (dir, runtime)
3385 }
3386
3387 #[tokio::test]
3388 async fn install_blob_store_on_read_only_runtime_refuses_mutators() {
3389 let (_dir, runtime) = make_read_only_runtime();
3390
3391 use khive_storage::BlobStore as _;
3394 let blob_root = tempfile::tempdir().unwrap();
3395 let writable = Arc::new(
3396 khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
3397 .expect("fs blob store"),
3398 );
3399 let seeded = writable
3400 .put(b"seed".to_vec())
3401 .await
3402 .expect("seed put through the raw store");
3403 runtime
3404 .install_blob_store(writable.clone())
3405 .expect("read-only install wraps rather than refusing");
3406
3407 let installed = runtime.blob_store().expect("installed store");
3408 assert!(
3409 installed.exists(&seeded).await.expect("exists"),
3410 "bounded read surface must stay available"
3411 );
3412 let err = installed
3413 .put(b"post-boot".to_vec())
3414 .await
3415 .expect_err("put must refuse on a read-only runtime");
3416 assert!(
3417 err.to_string().contains("read-only"),
3418 "refusal must name the mode: {err}"
3419 );
3420
3421 runtime
3425 .install_blob_store(writable.clone())
3426 .expect("reinstalling the same raw store must be idempotent");
3427
3428 let bypass_root = tempfile::tempdir().unwrap();
3432 let bypass_store = Arc::new(
3433 khive_db::stores::blob::FsBlobStore::new(bypass_root.path().to_path_buf(), 0)
3434 .expect("fs blob store"),
3435 );
3436 let writable_hydrator = Arc::new(
3437 crate::BlobHydrator::new(bypass_store, crate::DEFAULT_BLOB_HYDRATION_BYTES)
3438 .expect("construct writable hydrator"),
3439 );
3440 let err = runtime
3441 .install_blob_hydrator(writable_hydrator)
3442 .expect_err("a writable hydrator must be refused on a read-only runtime");
3443 assert!(
3444 err.to_string().contains("read-only"),
3445 "refusal must name the mode: {err}"
3446 );
3447 }
3448
3449 #[tokio::test]
3450 async fn mode_aware_hydrator_constructor_satisfies_read_only_install() {
3451 let (_dir, runtime) = make_read_only_runtime();
3457 let blob_root = tempfile::tempdir().unwrap();
3458 let raw = Arc::new(
3459 khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
3460 .expect("fs blob store"),
3461 ) as Arc<dyn khive_storage::BlobStore>;
3462 let seeded = raw.put(b"seed".to_vec()).await.expect("seed put");
3463 let hydrator = Arc::new(
3464 crate::BlobHydrator::for_mode(
3465 Arc::clone(&raw),
3466 crate::DEFAULT_BLOB_HYDRATION_BYTES,
3467 true,
3468 )
3469 .expect("mode-aware read-only construction"),
3470 );
3471 runtime
3472 .install_blob_hydrator(hydrator)
3473 .expect("a for_mode(read_only) hydrator must pass the read-only gate");
3474 let installed = runtime.blob_store().expect("installed store");
3475 assert!(installed.exists(&seeded).await.expect("exists"));
3476 installed
3477 .put(b"post".to_vec())
3478 .await
3479 .expect_err("mutation must refuse through the wrapped store");
3480
3481 let twin = Arc::new(
3485 crate::BlobHydrator::for_mode(
3486 Arc::clone(&raw),
3487 crate::DEFAULT_BLOB_HYDRATION_BYTES,
3488 true,
3489 )
3490 .expect("twin construction"),
3491 );
3492 runtime
3493 .install_blob_hydrator(twin)
3494 .expect("an equivalent pairing must read as idempotent");
3495 }
3496
3497 #[tokio::test]
3498 async fn shared_install_permits_governed_writable_hydrator_on_read_only_handle() {
3499 let (_dir, runtime) = make_read_only_runtime();
3508 let blob_root = tempfile::tempdir().unwrap();
3509 let raw = Arc::new(
3510 khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
3511 .expect("fs blob store"),
3512 ) as Arc<dyn khive_storage::BlobStore>;
3513 let hand_paired = Arc::new(
3514 crate::BlobHydrator::for_mode(raw, crate::DEFAULT_BLOB_HYDRATION_BYTES, false)
3515 .expect("writable construction"),
3516 );
3517 runtime
3518 .install_blob_hydrator(Arc::clone(&hand_paired))
3519 .expect_err("the plain seam must refuse a writable hydrator on a read-only handle");
3520 runtime
3521 .install_shared_blob_hydrator(hand_paired)
3522 .expect_err("the shared seam must refuse a hand-paired (ungoverned) hydrator");
3523
3524 let governing = khive_db::StorageBackend::memory().expect("memory backend");
3528 let cfg = crate::KhiveConfig {
3529 storage: crate::engine_config::StorageSectionConfig {
3530 blob: Some(crate::engine_config::BlobConfig::Fs {
3531 root: Some(blob_root.path().to_string_lossy().into_owned()),
3532 floor_bytes: Some(0),
3533 }),
3534 },
3535 ..crate::KhiveConfig::default()
3536 };
3537 let governed = Arc::new(
3538 crate::BlobHydrator::resolve_for_governing_backend(
3539 &cfg,
3540 &governing,
3541 &governing,
3542 crate::DEFAULT_BLOB_HYDRATION_BYTES,
3543 )
3544 .expect("governed construction"),
3545 );
3546 runtime
3547 .install_shared_blob_hydrator(governed)
3548 .expect("the shared seam accepts a governed hydrator");
3549 let installed = runtime.blob_store().expect("installed store");
3550 let put = installed
3551 .put(b"shared-write".to_vec())
3552 .await
3553 .expect("the governed writable capability must actually mutate");
3554 assert!(installed.exists(&put).await.expect("exists"));
3555 }
3556
3557 #[test]
3567 #[serial]
3568 fn tilde_prefixed_db_override_resolves_and_boots_like_the_absolute_equivalent() {
3569 if crate::test_process::run_in_child() {
3570 return;
3571 }
3572
3573 let original_home = std::env::var_os("HOME");
3574 let original_cwd = std::env::current_dir().expect("read cwd");
3575 let home_dir = tempfile::tempdir().expect("home tempdir");
3576 let work_dir = tempfile::tempdir().expect("work tempdir");
3577 std::env::set_var("HOME", home_dir.path());
3578 std::env::set_current_dir(work_dir.path()).expect("chdir into isolated work dir");
3579
3580 let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3581 let tilde_anchor = crate::config::resolve_db_anchor(Some("~/data.db"))
3582 .expect("an explicit path always anchors");
3583 let expected = home_dir.path().join("data.db");
3584 assert_eq!(
3585 tilde_anchor, expected,
3586 "resolve_db_anchor must expand a leading ~ to $HOME before it ever \
3587 reaches RuntimeConfig.db_path"
3588 );
3589
3590 let absolute_anchor = crate::config::resolve_db_anchor(Some(
3591 expected.to_str().expect("utf8 tempdir path"),
3592 ))
3593 .expect("an explicit path always anchors");
3594 assert_eq!(
3595 tilde_anchor, absolute_anchor,
3596 "a ~-prefixed override and its equivalent absolute path must resolve to \
3597 the identical anchor"
3598 );
3599
3600 let make_config = |db_path: std::path::PathBuf| RuntimeConfig {
3601 web: Default::default(),
3602 telemetry: Default::default(),
3603 mounts: Vec::new(),
3604 brain: Default::default(),
3605 git_write: Default::default(),
3606 display_timezone: chrono_tz::Tz::UTC,
3607 events_split: None,
3608 db_path: Some(db_path),
3609 wal_ceiling_bytes: 0,
3610 wal_ceiling_configured_bytes: 0,
3611 wal_ceiling_source: khive_db::WalCeilingSource::Default,
3612 wal_ceiling_env_raw: None,
3613 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3614 default_namespace: Namespace::local(),
3615 embedding_model: None,
3616 additional_embedding_models: vec![],
3617 gate: Arc::new(AllowAllGate),
3618 packs: vec!["kg".to_string()],
3619 backend_id: BackendId::main(),
3620 brain_profile: None,
3621 visible_namespaces: vec![],
3622 allowed_outbound_namespaces: vec![],
3623 actor_id: None,
3624 exec: Default::default(),
3625 ..crate::RuntimeConfig::no_embeddings()
3626 };
3627
3628 let tilde_cfg = make_config(tilde_anchor.clone());
3629
3630 let rt =
3631 KhiveRuntime::new_for_test(tilde_cfg).expect("boot must open the expanded path");
3632 assert_eq!(
3633 rt.backend_data_dir().expect("file backend"),
3634 home_dir.path(),
3635 "single-backend boot must open the file under the expanded $HOME \
3636 directory, not a literal ~ path relative to cwd"
3637 );
3638 assert!(
3639 expected.exists(),
3640 "the database file must be created at the expanded $HOME path"
3641 );
3642 assert!(
3643 !work_dir.path().join("~").exists(),
3644 "boot must never create a literal '~' directory under the process cwd"
3645 );
3646 }));
3647
3648 match &original_home {
3649 Some(h) => std::env::set_var("HOME", h),
3650 None => std::env::remove_var("HOME"),
3651 }
3652 let _ = std::env::set_current_dir(&original_cwd);
3653 outcome.expect("test body panicked");
3654 }
3655
3656 #[test]
3657 fn from_backend_uses_provided_backend() {
3658 let backend = Arc::new(StorageBackend::memory().expect("memory backend"));
3659 let config = RuntimeConfig {
3660 web: Default::default(),
3661 telemetry: Default::default(),
3662 mounts: Vec::new(),
3663 brain: Default::default(),
3664 git_write: Default::default(),
3665 display_timezone: chrono_tz::Tz::UTC,
3666 events_split: None,
3667 db_path: None,
3668 wal_ceiling_bytes: 0,
3669 wal_ceiling_configured_bytes: 0,
3670 wal_ceiling_source: khive_db::WalCeilingSource::Default,
3671 wal_ceiling_env_raw: None,
3672 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3673 default_namespace: Namespace::local(),
3674 embedding_model: None,
3675 additional_embedding_models: vec![],
3676 gate: Arc::new(AllowAllGate),
3677 packs: vec!["kg".to_string()],
3678 backend_id: BackendId::parse("lore").expect("valid backend id"),
3679 brain_profile: None,
3680 visible_namespaces: vec![],
3681 allowed_outbound_namespaces: vec![],
3682 actor_id: None,
3683 exec: Default::default(),
3684 ..crate::RuntimeConfig::no_embeddings()
3685 };
3686 let rt = KhiveRuntime::from_backend(backend, config);
3687 assert_eq!(rt.backend_id().as_str(), "lore");
3688 assert!(rt.config().db_path.is_none());
3689 }
3690
3691 #[test]
3692 fn backend_id_defaults_to_main() {
3693 let rt = KhiveRuntime::memory().unwrap();
3694 assert_eq!(rt.backend_id().as_str(), BackendId::MAIN);
3695 }
3696
3697 #[cfg(unix)]
3698 #[test]
3699 fn storage_identity_accepts_hard_links_but_rejects_distinct_files() {
3700 let dir = tempfile::tempdir().expect("temporary database directory");
3701 let first_path = dir.path().join("first.db");
3702 let alias_path = dir.path().join("alias.db");
3703 let distinct_path = dir.path().join("distinct.db");
3704 let open = |path: &std::path::Path| {
3705 Arc::new(
3706 StorageBackend::sqlite_for_test_with_journal_mode(
3707 path,
3708 false,
3709 std::time::Duration::from_secs(1),
3710 )
3711 .expect("test backend"),
3712 )
3713 };
3714 let first = open(&first_path);
3715 std::fs::hard_link(&first_path, &alias_path).expect("hard-link alias");
3716 let alias = open(&alias_path);
3717 let distinct = open(&distinct_path);
3718 let first = KhiveRuntime::from_backend(first, RuntimeConfig::no_embeddings());
3719 let alias = KhiveRuntime::from_backend(alias, RuntimeConfig::no_embeddings());
3720 let distinct = KhiveRuntime::from_backend(distinct, RuntimeConfig::no_embeddings());
3721
3722 assert!(first.shares_backend_storage_with(&alias));
3723 assert!(!first.shares_backend_storage_with(&distinct));
3724 assert!(first.shares_backend_storage_with(&first.clone()));
3725 }
3726
3727 #[test]
3728 fn store_accessors_return_ok() {
3729 let rt = KhiveRuntime::memory().unwrap();
3730 let tok = NamespaceToken::local();
3731 assert!(rt.entities(&tok).is_ok());
3732 assert!(rt.graph(&tok).is_ok());
3733 assert!(rt.notes(&tok).is_ok());
3734 assert!(rt.events(&tok).is_ok());
3735 }
3736
3737 fn attributed_event_runtime() -> (KhiveRuntime, NamespaceToken) {
3738 let runtime = KhiveRuntime::new(RuntimeConfig {
3739 db_path: None,
3740 actor_id: Some("lambda:enrolled".to_string()),
3741 ..RuntimeConfig::no_embeddings()
3742 })
3743 .expect("memory runtime");
3744 let token = runtime
3745 .authorize(Namespace::local())
3746 .expect("configured actor is allowed by the default gate");
3747 (runtime, token)
3748 }
3749
3750 fn forged_event(verb: &str) -> Event {
3751 Event::new(
3752 "caller-selected-namespace",
3753 verb,
3754 EventKind::Audit,
3755 SubstrateKind::Event,
3756 "caller-selected-actor",
3757 )
3758 }
3759
3760 fn assert_token_attribution(event: &Event) {
3761 assert_eq!(event.namespace, "local");
3762 assert_eq!(event.actor, "actor:lambda:enrolled");
3763 }
3764
3765 #[tokio::test]
3766 async fn token_scoped_event_store_stamps_resolved_attribution_on_single_append() {
3767 let (runtime, token) = attributed_event_runtime();
3768 let store = runtime.events(&token).expect("event store");
3769 let event = forged_event("single");
3770 let id = event.id;
3771
3772 store.append_event(event).await.expect("append");
3773
3774 let stored = store
3775 .get_event(id)
3776 .await
3777 .expect("read")
3778 .expect("the token-stamped event remains visible to the token");
3779 assert_token_attribution(&stored);
3780 assert_eq!(stored.verb, "single", "non-attribution fields survive");
3781 }
3782
3783 #[tokio::test]
3784 async fn token_scoped_event_store_stamps_resolved_attribution_on_batch_paths() {
3785 let (runtime, token) = attributed_event_runtime();
3786 let store = runtime.events(&token).expect("event store");
3787 let ordinary = forged_event("batch");
3788 let ordinary_id = ordinary.id;
3789
3790 store
3791 .append_events(vec![ordinary])
3792 .await
3793 .expect("ordinary batch append");
3794 let stored = store
3795 .get_event(ordinary_id)
3796 .await
3797 .expect("read")
3798 .expect("ordinary batch event remains visible to the token");
3799 assert_token_attribution(&stored);
3800
3801 let idempotent = forged_event("idempotent_batch");
3802 let idempotent_id = idempotent.id;
3803 let outcome = store
3804 .append_events_idempotent(vec![idempotent])
3805 .await
3806 .expect("idempotent batch append");
3807 assert_eq!(
3808 outcome.rows,
3809 vec![khive_storage::event::EventAppendDisposition::Inserted]
3810 );
3811 let stored = store
3812 .get_event(idempotent_id)
3813 .await
3814 .expect("read")
3815 .expect("idempotent batch event remains visible to the token");
3816 assert_token_attribution(&stored);
3817 }
3818
3819 #[test]
3820 fn vectors_returns_unconfigured_without_model() {
3821 let rt = KhiveRuntime::memory().unwrap();
3822 let tok = NamespaceToken::local();
3823 match rt.vectors(&tok) {
3824 Err(crate::RuntimeError::Unconfigured(s)) => assert_eq!(s, "embedding_model"),
3825 Err(other) => panic!("expected Unconfigured, got {:?}", other),
3826 Ok(_) => panic!("expected Err, got Ok"),
3827 }
3828 }
3829
3830 #[test]
3831 fn vec_model_key_sanitizes_dots_and_dashes() {
3832 assert_eq!(
3833 vec_model_key(EmbeddingModel::BgeSmallEnV15),
3834 "bge_small_en_v1_5"
3835 );
3836 assert_eq!(
3837 vec_model_key(EmbeddingModel::BgeBaseEnV15),
3838 "bge_base_en_v1_5"
3839 );
3840 assert_eq!(
3841 vec_model_key(EmbeddingModel::AllMiniLmL6V2),
3842 "all_minilm_l6_v2"
3843 );
3844 }
3845
3846 #[test]
3847 fn default_config_uses_allow_all_gate() {
3848 let cfg = RuntimeConfig::default();
3849 assert_eq!(cfg.default_namespace.as_str(), "local");
3850 let _: GateRef = cfg.gate.clone();
3851 }
3852
3853 #[test]
3854 fn parse_pack_list_handles_comma_and_whitespace() {
3855 assert_eq!(parse_pack_list("kg"), vec!["kg".to_string()]);
3856 assert_eq!(
3857 parse_pack_list("kg,gtd"),
3858 vec!["kg".to_string(), "gtd".to_string()]
3859 );
3860 assert_eq!(
3861 parse_pack_list(" kg , gtd "),
3862 vec!["kg".to_string(), "gtd".to_string()]
3863 );
3864 assert_eq!(
3865 parse_pack_list("kg gtd"),
3866 vec!["kg".to_string(), "gtd".to_string()]
3867 );
3868 assert_eq!(parse_pack_list(",,"), Vec::<String>::new());
3869 assert_eq!(parse_pack_list(""), Vec::<String>::new());
3870 }
3871
3872 #[test]
3873 fn default_config_packs_loads_production_set() {
3874 let prior = std::env::var("KHIVE_PACKS").ok();
3875 unsafe {
3877 std::env::remove_var("KHIVE_PACKS");
3878 }
3879 let cfg = RuntimeConfig::default();
3881 assert_eq!(cfg.packs, RuntimeConfig::built_in_packs());
3882 assert!(cfg.packs.contains(&"kg".to_string()));
3883 assert!(cfg.packs.contains(&"gtd".to_string()));
3884 assert!(cfg.packs.contains(&"memory".to_string()));
3885 assert!(cfg.packs.contains(&"brain".to_string()));
3886 assert!(cfg.packs.contains(&"comm".to_string()));
3887 assert!(cfg.packs.contains(&"schedule".to_string()));
3888 assert!(cfg.packs.contains(&"knowledge".to_string()));
3889 assert!(cfg.packs.contains(&"session".to_string()));
3892 assert!(cfg.packs.contains(&"git".to_string()));
3893 assert!(cfg.packs.contains(&"code".to_string()));
3894 assert!(cfg.packs.contains(&"workspace".to_string()));
3895 assert!(cfg.packs.contains(&"blob".to_string()));
3900 assert!(cfg.packs.contains(&"tool".to_string()));
3903 assert!(cfg.packs.contains(&"exec".to_string()));
3904 assert_eq!(cfg.packs.len(), 14);
3905 if let Some(v) = prior {
3906 unsafe {
3908 std::env::set_var("KHIVE_PACKS", v);
3909 }
3910 }
3911 }
3912
3913 #[test]
3914 fn default_config_uses_minilm_when_env_unset() {
3915 let prior = std::env::var("KHIVE_EMBEDDING_MODEL").ok();
3916 unsafe {
3919 std::env::remove_var("KHIVE_EMBEDDING_MODEL");
3920 }
3921 let cfg = RuntimeConfig::default();
3922 assert_eq!(cfg.embedding_model, Some(EmbeddingModel::AllMiniLmL6V2));
3923 if let Some(v) = prior {
3924 unsafe {
3926 std::env::set_var("KHIVE_EMBEDDING_MODEL", v);
3927 }
3928 }
3929 }
3930
3931 use crate::engine_config::{ActorConfig, KhiveConfig};
3934
3935 fn khive_cfg_with_actor(id: &str) -> KhiveConfig {
3936 KhiveConfig {
3937 engines: vec![],
3938 actor: ActorConfig {
3939 id: Some(id.to_string()),
3940 display_name: None,
3941 ..Default::default()
3942 },
3943 ..KhiveConfig::default()
3944 }
3945 }
3946
3947 #[test]
3948 fn runtime_config_from_khive_config_actor_id_does_not_override_default_namespace() {
3949 let base = RuntimeConfig {
3954 web: Default::default(),
3955 telemetry: Default::default(),
3956 mounts: Vec::new(),
3957 brain: Default::default(),
3958 git_write: Default::default(),
3959 display_timezone: chrono_tz::Tz::UTC,
3960 events_split: None,
3961 db_path: None,
3962 wal_ceiling_bytes: 0,
3963 wal_ceiling_configured_bytes: 0,
3964 wal_ceiling_source: khive_db::WalCeilingSource::Default,
3965 wal_ceiling_env_raw: None,
3966 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3967 default_namespace: Namespace::local(),
3968 embedding_model: None,
3969 additional_embedding_models: vec![],
3970 gate: Arc::new(AllowAllGate),
3971 packs: vec!["kg".to_string()],
3972 backend_id: BackendId::main(),
3973 brain_profile: None,
3974 visible_namespaces: vec![],
3975 allowed_outbound_namespaces: vec![],
3976 actor_id: None,
3977 exec: Default::default(),
3978 ..crate::RuntimeConfig::no_embeddings()
3979 };
3980 let cfg = khive_cfg_with_actor("lambda:khive");
3981 let result = runtime_config_from_khive_config(&cfg, base);
3982 assert_eq!(
3983 result.default_namespace.as_str(),
3984 "local",
3985 "actor.id must not become default_namespace (ADR-007 Rev 4 Rule 0); writes pin to local"
3986 );
3987 }
3988
3989 #[test]
3990 fn runtime_config_from_khive_config_empty_actor_id_keeps_base_namespace() {
3991 let base = RuntimeConfig {
3992 web: Default::default(),
3993 telemetry: Default::default(),
3994 mounts: Vec::new(),
3995 brain: Default::default(),
3996 git_write: Default::default(),
3997 display_timezone: chrono_tz::Tz::UTC,
3998 events_split: None,
3999 db_path: None,
4000 wal_ceiling_bytes: 0,
4001 wal_ceiling_configured_bytes: 0,
4002 wal_ceiling_source: khive_db::WalCeilingSource::Default,
4003 wal_ceiling_env_raw: None,
4004 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
4005 default_namespace: Namespace::parse("lambda:base").unwrap(),
4006 embedding_model: None,
4007 additional_embedding_models: vec![],
4008 gate: Arc::new(AllowAllGate),
4009 packs: vec!["kg".to_string()],
4010 backend_id: BackendId::main(),
4011 brain_profile: None,
4012 visible_namespaces: vec![],
4013 allowed_outbound_namespaces: vec![],
4014 actor_id: None,
4015 exec: Default::default(),
4016 ..crate::RuntimeConfig::no_embeddings()
4017 };
4018 let cfg = KhiveConfig {
4019 engines: vec![],
4020 actor: ActorConfig {
4021 id: Some(String::new()),
4022 display_name: None,
4023 ..Default::default()
4024 },
4025 ..KhiveConfig::default()
4026 };
4027 let result = runtime_config_from_khive_config(&cfg, base);
4028 assert_eq!(
4029 result.default_namespace.as_str(),
4030 "lambda:base",
4031 "empty actor.id must not override base namespace"
4032 );
4033 }
4034
4035 #[test]
4036 fn runtime_config_from_khive_config_absent_actor_id_keeps_base_namespace() {
4037 let base = RuntimeConfig {
4038 web: Default::default(),
4039 telemetry: Default::default(),
4040 mounts: Vec::new(),
4041 brain: Default::default(),
4042 git_write: Default::default(),
4043 display_timezone: chrono_tz::Tz::UTC,
4044 events_split: None,
4045 db_path: None,
4046 wal_ceiling_bytes: 0,
4047 wal_ceiling_configured_bytes: 0,
4048 wal_ceiling_source: khive_db::WalCeilingSource::Default,
4049 wal_ceiling_env_raw: None,
4050 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
4051 default_namespace: Namespace::parse("lambda:base").unwrap(),
4052 embedding_model: None,
4053 additional_embedding_models: vec![],
4054 gate: Arc::new(AllowAllGate),
4055 packs: vec!["kg".to_string()],
4056 backend_id: BackendId::main(),
4057 brain_profile: None,
4058 visible_namespaces: vec![],
4059 allowed_outbound_namespaces: vec![],
4060 actor_id: None,
4061 exec: Default::default(),
4062 ..crate::RuntimeConfig::no_embeddings()
4063 };
4064 let cfg = KhiveConfig::default(); let result = runtime_config_from_khive_config(&cfg, base);
4066 assert_eq!(
4067 result.default_namespace.as_str(),
4068 "lambda:base",
4069 "absent actor.id must not override base namespace"
4070 );
4071 }
4072
4073 #[test]
4074 fn runtime_config_from_khive_config_actor_id_with_engines() {
4075 let base = RuntimeConfig {
4076 web: Default::default(),
4077 telemetry: Default::default(),
4078 mounts: Vec::new(),
4079 brain: Default::default(),
4080 git_write: Default::default(),
4081 display_timezone: chrono_tz::Tz::UTC,
4082 events_split: None,
4083 db_path: None,
4084 wal_ceiling_bytes: 0,
4085 wal_ceiling_configured_bytes: 0,
4086 wal_ceiling_source: khive_db::WalCeilingSource::Default,
4087 wal_ceiling_env_raw: None,
4088 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
4089 default_namespace: Namespace::local(),
4090 embedding_model: None,
4091 additional_embedding_models: vec![],
4092 gate: Arc::new(AllowAllGate),
4093 packs: vec!["kg".to_string()],
4094 backend_id: BackendId::main(),
4095 brain_profile: None,
4096 visible_namespaces: vec![],
4097 allowed_outbound_namespaces: vec![],
4098 actor_id: None,
4099 exec: Default::default(),
4100 ..crate::RuntimeConfig::no_embeddings()
4101 };
4102 let cfg = KhiveConfig {
4103 engines: vec![crate::engine_config::EngineConfig {
4104 name: "default".to_string(),
4105 model: "all-minilm-l6-v2".to_string(),
4106 default: true,
4107 fusion_weight: None,
4108 dims: None,
4109 }],
4110 actor: ActorConfig {
4111 id: Some("lambda:test".to_string()),
4112 display_name: None,
4113 ..Default::default()
4114 },
4115 ..KhiveConfig::default()
4116 };
4117 let result = runtime_config_from_khive_config(&cfg, base);
4118 assert_eq!(
4119 result.default_namespace.as_str(),
4120 "local",
4121 "actor.id must not override default_namespace (ADR-007 Rev 4 Rule 0); \
4122 writes pin to local; engine config is still applied"
4123 );
4124 assert!(result.embedding_model.is_some());
4125 }
4126
4127 #[test]
4130 fn runtime_config_from_khive_config_display_timezone_overrides_base() {
4131 let base = RuntimeConfig {
4132 web: Default::default(),
4133 telemetry: Default::default(),
4134 mounts: Vec::new(),
4135 brain: Default::default(),
4136 git_write: Default::default(),
4137 display_timezone: chrono_tz::Tz::UTC,
4138 events_split: None,
4139 db_path: None,
4140 wal_ceiling_bytes: 0,
4141 wal_ceiling_configured_bytes: 0,
4142 wal_ceiling_source: khive_db::WalCeilingSource::Default,
4143 wal_ceiling_env_raw: None,
4144 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
4145 default_namespace: Namespace::local(),
4146 embedding_model: None,
4147 additional_embedding_models: vec![],
4148 gate: Arc::new(AllowAllGate),
4149 packs: vec!["kg".to_string()],
4150 backend_id: BackendId::main(),
4151 brain_profile: None,
4152 visible_namespaces: vec![],
4153 allowed_outbound_namespaces: vec![],
4154 actor_id: None,
4155 exec: Default::default(),
4156 ..crate::RuntimeConfig::no_embeddings()
4157 };
4158 let cfg = KhiveConfig {
4159 display: crate::engine_config::DisplaySectionConfig {
4160 timezone: Some("America/New_York".to_string()),
4161 },
4162 ..KhiveConfig::default()
4163 };
4164 let result = runtime_config_from_khive_config(&cfg, base);
4165 assert_eq!(
4166 result.display_timezone,
4167 "America/New_York".parse::<chrono_tz::Tz>().unwrap(),
4168 "[display] timezone in khive.toml must override base.display_timezone"
4169 );
4170 }
4171
4172 #[test]
4173 fn runtime_config_from_khive_config_absent_display_timezone_keeps_base() {
4174 let base = RuntimeConfig {
4175 web: Default::default(),
4176 telemetry: Default::default(),
4177 mounts: Vec::new(),
4178 brain: Default::default(),
4179 git_write: Default::default(),
4180 display_timezone: "Asia/Tokyo".parse().unwrap(),
4181 events_split: None,
4182 db_path: None,
4183 wal_ceiling_bytes: 0,
4184 wal_ceiling_configured_bytes: 0,
4185 wal_ceiling_source: khive_db::WalCeilingSource::Default,
4186 wal_ceiling_env_raw: None,
4187 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
4188 default_namespace: Namespace::local(),
4189 embedding_model: None,
4190 additional_embedding_models: vec![],
4191 gate: Arc::new(AllowAllGate),
4192 packs: vec!["kg".to_string()],
4193 backend_id: BackendId::main(),
4194 brain_profile: None,
4195 visible_namespaces: vec![],
4196 allowed_outbound_namespaces: vec![],
4197 actor_id: None,
4198 exec: Default::default(),
4199 ..crate::RuntimeConfig::no_embeddings()
4200 };
4201 let cfg = KhiveConfig::default(); let result = runtime_config_from_khive_config(&cfg, base);
4203 assert_eq!(
4204 result.display_timezone,
4205 "Asia/Tokyo".parse::<chrono_tz::Tz>().unwrap(),
4206 "absent [display] timezone must preserve base.display_timezone unchanged"
4207 );
4208 }
4209
4210 #[test]
4219 #[serial]
4220 fn runtime_config_from_khive_config_engines_present_preserves_env_actor_when_toml_has_none() {
4221 let prior = std::env::var("KHIVE_ACTOR").ok();
4222 unsafe {
4224 std::env::set_var("KHIVE_ACTOR", "lambda:test-env-actor");
4225 }
4226 let base = RuntimeConfig::default();
4227 assert_eq!(base.actor_id.as_deref(), Some("lambda:test-env-actor"));
4228
4229 let cfg = KhiveConfig {
4230 engines: vec![crate::engine_config::EngineConfig {
4231 name: "default".to_string(),
4232 model: "all-minilm-l6-v2".to_string(),
4233 default: true,
4234 fusion_weight: None,
4235 dims: None,
4236 }],
4237 actor: ActorConfig::default(), ..KhiveConfig::default()
4239 };
4240 let result = runtime_config_from_khive_config(&cfg, base);
4241 assert_eq!(
4242 result.actor_id.as_deref(),
4243 Some("lambda:test-env-actor"),
4244 "engines-present arm must preserve base.actor_id (env actor) when TOML has no [actor] id"
4245 );
4246
4247 unsafe {
4249 match prior {
4250 Some(v) => std::env::set_var("KHIVE_ACTOR", v),
4251 None => std::env::remove_var("KHIVE_ACTOR"),
4252 }
4253 }
4254 }
4255
4256 #[test]
4257 #[serial]
4258 fn runtime_config_from_khive_config_engines_empty_preserves_env_actor_when_toml_has_none() {
4259 let prior = std::env::var("KHIVE_ACTOR").ok();
4260 unsafe {
4262 std::env::set_var("KHIVE_ACTOR", "lambda:test-env-actor");
4263 }
4264 let base = RuntimeConfig::default();
4265 assert_eq!(base.actor_id.as_deref(), Some("lambda:test-env-actor"));
4266
4267 let cfg = KhiveConfig {
4268 engines: vec![],
4269 actor: ActorConfig::default(), ..KhiveConfig::default()
4271 };
4272 let result = runtime_config_from_khive_config(&cfg, base);
4273 assert_eq!(
4274 result.actor_id.as_deref(),
4275 Some("lambda:test-env-actor"),
4276 "engines-empty early-return arm must preserve base.actor_id (env actor) when TOML has no [actor] id"
4277 );
4278
4279 unsafe {
4281 match prior {
4282 Some(v) => std::env::set_var("KHIVE_ACTOR", v),
4283 None => std::env::remove_var("KHIVE_ACTOR"),
4284 }
4285 }
4286 }
4287
4288 #[test]
4289 #[serial]
4290 fn runtime_config_from_khive_config_toml_actor_wins_over_env_actor() {
4291 let prior = std::env::var("KHIVE_ACTOR").ok();
4292 unsafe {
4294 std::env::set_var("KHIVE_ACTOR", "lambda:test-env-actor");
4295 }
4296 let base = RuntimeConfig::default();
4297 assert_eq!(base.actor_id.as_deref(), Some("lambda:test-env-actor"));
4298
4299 let cfg = khive_cfg_with_actor("lambda:toml-actor");
4300 let result = runtime_config_from_khive_config(&cfg, base);
4301 assert_eq!(
4302 result.actor_id.as_deref(),
4303 Some("lambda:toml-actor"),
4304 "TOML [actor] id must win over the env-resolved base.actor_id"
4305 );
4306
4307 unsafe {
4309 match prior {
4310 Some(v) => std::env::set_var("KHIVE_ACTOR", v),
4311 None => std::env::remove_var("KHIVE_ACTOR"),
4312 }
4313 }
4314 }
4315
4316 fn migrated_memory_backend() -> Arc<StorageBackend> {
4322 let backend = StorageBackend::memory().expect("memory backend");
4323 {
4324 let mut writer = backend.pool().try_writer().expect("writer");
4325 khive_db::run_migrations(writer.conn_mut()).expect("migrations");
4326 }
4327 Arc::new(backend)
4328 }
4329
4330 fn secondary_config() -> RuntimeConfig {
4331 RuntimeConfig {
4332 web: Default::default(),
4333 telemetry: Default::default(),
4334 mounts: Vec::new(),
4335 brain: Default::default(),
4336 git_write: Default::default(),
4337 exec: Default::default(),
4338 display_timezone: chrono_tz::Tz::UTC,
4339 events_split: None,
4340 db_path: None,
4341 wal_ceiling_bytes: 0,
4342 wal_ceiling_configured_bytes: 0,
4343 wal_ceiling_source: khive_db::WalCeilingSource::Default,
4344 wal_ceiling_env_raw: None,
4345 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
4346 default_namespace: Namespace::local(),
4347 embedding_model: None,
4348 additional_embedding_models: vec![],
4349 gate: Arc::new(AllowAllGate),
4350 packs: vec!["kg".to_string()],
4351 backend_id: BackendId::parse("lore").expect("valid backend id"),
4352 brain_profile: None,
4353 visible_namespaces: vec![],
4354 allowed_outbound_namespaces: vec![],
4355 actor_id: None,
4356 ..crate::RuntimeConfig::no_embeddings()
4357 }
4358 }
4359
4360 #[test]
4361 fn core_on_main_runtime_returns_same_backend_id() {
4362 let rt = KhiveRuntime::memory().unwrap();
4364 assert_eq!(rt.backend_id().as_str(), BackendId::MAIN);
4365 let core_rt = rt.core();
4366 assert_eq!(core_rt.backend_id().as_str(), BackendId::MAIN);
4367 }
4368
4369 #[tokio::test]
4370 async fn core_on_main_runtime_round_trips_note() {
4371 let rt = KhiveRuntime::memory().unwrap();
4374 let tok = NamespaceToken::local();
4375
4376 let note = rt
4377 .core()
4378 .create_note(
4379 &tok,
4380 "observation",
4381 None,
4382 "adr073-main-round-trip",
4383 None,
4384 None,
4385 vec![],
4386 )
4387 .await
4388 .expect("create_note via core()");
4389
4390 let found = rt
4391 .notes(&tok)
4392 .expect("notes store")
4393 .get_note(note.id)
4394 .await
4395 .expect("get_note");
4396
4397 assert!(
4398 found.is_some(),
4399 "note written via core() must be visible through original rt"
4400 );
4401 }
4402
4403 #[tokio::test]
4415 async fn cross_backend_split_note_to_main_aux_to_secondary() {
4416 use khive_storage::{SqlStatement, SqlValue};
4417
4418 let main_arc = migrated_memory_backend();
4419 let secondary_arc = migrated_memory_backend();
4420
4421 let main_config = RuntimeConfig {
4422 web: Default::default(),
4423 telemetry: Default::default(),
4424 mounts: Vec::new(),
4425 brain: Default::default(),
4426 git_write: Default::default(),
4427 display_timezone: chrono_tz::Tz::UTC,
4428 events_split: None,
4429 db_path: None,
4430 wal_ceiling_bytes: 0,
4431 wal_ceiling_configured_bytes: 0,
4432 wal_ceiling_source: khive_db::WalCeilingSource::Default,
4433 wal_ceiling_env_raw: None,
4434 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
4435 default_namespace: Namespace::local(),
4436 embedding_model: None,
4437 additional_embedding_models: vec![],
4438 gate: Arc::new(AllowAllGate),
4439 packs: vec!["kg".to_string()],
4440 backend_id: BackendId::main(),
4441 brain_profile: None,
4442 visible_namespaces: vec![],
4443 allowed_outbound_namespaces: vec![],
4444 actor_id: None,
4445 exec: Default::default(),
4446 ..crate::RuntimeConfig::no_embeddings()
4447 };
4448
4449 let rt_main = KhiveRuntime::from_backend(main_arc.clone(), main_config);
4450 let rt_secondary = KhiveRuntime::from_backend(secondary_arc, secondary_config())
4451 .with_core_backend(main_arc.clone());
4452
4453 let tok = NamespaceToken::local();
4454
4455 let note = rt_secondary
4458 .core()
4459 .create_note(
4460 &tok,
4461 "observation",
4462 None,
4463 "adr073-split-test",
4464 None,
4465 None,
4466 vec![],
4467 )
4468 .await
4469 .expect("create_note via core()");
4470 let note_id = note.id;
4471
4472 let in_main = rt_main
4474 .notes(&tok)
4475 .expect("main notes store")
4476 .get_note(note_id)
4477 .await
4478 .expect("get_note from main");
4479 assert!(
4480 in_main.is_some(),
4481 "note written via core() must appear in main backend A"
4482 );
4483
4484 let in_secondary = rt_secondary
4486 .notes(&tok)
4487 .expect("secondary notes store")
4488 .get_note(note_id)
4489 .await
4490 .expect("get_note from secondary");
4491 assert!(
4492 in_secondary.is_none(),
4493 "note written to main via core() must NOT appear in secondary backend B"
4494 );
4495
4496 {
4499 let mut writer = rt_secondary.sql().writer().await.expect("secondary writer");
4500 writer
4501 .execute(SqlStatement {
4502 sql: "CREATE TABLE IF NOT EXISTS _test_adr073_aux \
4503 (marker TEXT PRIMARY KEY)"
4504 .into(),
4505 params: vec![],
4506 label: None,
4507 })
4508 .await
4509 .expect("create aux table in B");
4510 writer
4511 .execute(SqlStatement {
4512 sql: "INSERT INTO _test_adr073_aux VALUES (?1)".into(),
4513 params: vec![SqlValue::Text("b-side-sentinel".into())],
4514 label: None,
4515 })
4516 .await
4517 .expect("insert into aux table in B");
4518 }
4519
4520 let mut reader_b = rt_secondary.sql().reader().await.expect("secondary reader");
4522 let rows_b = reader_b
4523 .query_all(SqlStatement {
4524 sql: "SELECT marker FROM _test_adr073_aux".into(),
4525 params: vec![],
4526 label: None,
4527 })
4528 .await
4529 .expect("select from B");
4530 assert_eq!(rows_b.len(), 1, "aux row must exist in B");
4531 match rows_b[0].get("marker") {
4532 Some(SqlValue::Text(s)) => {
4533 assert_eq!(s, "b-side-sentinel", "sentinel value must match")
4534 }
4535 other => panic!("expected Text('b-side-sentinel'), got {other:?}"),
4536 }
4537
4538 let mut reader_a = rt_main.sql().reader().await.expect("main reader");
4540 let result_a = reader_a
4541 .query_all(SqlStatement {
4542 sql: "SELECT marker FROM _test_adr073_aux".into(),
4543 params: vec![],
4544 label: None,
4545 })
4546 .await;
4547 match result_a {
4549 Err(e) => assert!(
4550 e.to_string().contains("no such table"),
4551 "expected 'no such table' error from A, got: {e}"
4552 ),
4553 Ok(rows) => assert!(
4554 rows.is_empty(),
4555 "aux table must not have rows in A, got {} rows",
4556 rows.len()
4557 ),
4558 }
4559 }
4560
4561 #[test]
4562 fn constructors_leave_core_backend_none_by_behavior() {
4563 let rt_mem = KhiveRuntime::memory().unwrap();
4566 assert_eq!(rt_mem.core().backend_id().as_str(), BackendId::MAIN);
4567
4568 let backend = migrated_memory_backend();
4569 let rt_from = KhiveRuntime::from_backend(
4570 backend,
4571 RuntimeConfig {
4572 web: Default::default(),
4573 telemetry: Default::default(),
4574 mounts: Vec::new(),
4575 brain: Default::default(),
4576 git_write: Default::default(),
4577 display_timezone: chrono_tz::Tz::UTC,
4578 events_split: None,
4579 db_path: None,
4580 wal_ceiling_bytes: 0,
4581 wal_ceiling_configured_bytes: 0,
4582 wal_ceiling_source: khive_db::WalCeilingSource::Default,
4583 wal_ceiling_env_raw: None,
4584 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
4585 default_namespace: Namespace::local(),
4586 embedding_model: None,
4587 additional_embedding_models: vec![],
4588 gate: Arc::new(AllowAllGate),
4589 packs: vec!["kg".to_string()],
4590 backend_id: BackendId::parse("lore").expect("valid backend id"),
4591 brain_profile: None,
4592 visible_namespaces: vec![],
4593 allowed_outbound_namespaces: vec![],
4594 actor_id: None,
4595 exec: Default::default(),
4596 ..crate::RuntimeConfig::no_embeddings()
4597 },
4598 );
4599 assert_eq!(rt_from.core().backend_id().as_str(), "lore");
4602 }
4603
4604 #[test]
4605 fn with_core_backend_sets_core_then_core_returns_main_id() {
4606 let main_arc = migrated_memory_backend();
4608 let secondary_arc = migrated_memory_backend();
4609
4610 let rt_secondary = KhiveRuntime::from_backend(secondary_arc, secondary_config())
4611 .with_core_backend(main_arc);
4612
4613 assert_eq!(rt_secondary.backend_id().as_str(), "lore");
4614 assert_eq!(
4615 rt_secondary.core().backend_id().as_str(),
4616 BackendId::MAIN,
4617 "core() on a secondary runtime must return a main-bound handle"
4618 );
4619 }
4620
4621 #[test]
4622 fn attachment_store_rejects_secondary_handle_and_accepts_its_core_projection() {
4623 let main_arc = migrated_memory_backend();
4624 let secondary_arc = migrated_memory_backend();
4625 let rt_secondary = KhiveRuntime::from_backend(secondary_arc, secondary_config())
4626 .with_core_backend(main_arc);
4627
4628 let error = match rt_secondary.attachments() {
4629 Ok(_) => panic!("a secondary runtime must not expose attachment mutation"),
4630 Err(error) => error,
4631 };
4632 assert!(matches!(error, RuntimeError::InvalidInput(_)));
4633 assert!(
4634 error.to_string().contains("canonical main backend"),
4635 "secondary refusal must explain the liveness authority: {error}"
4636 );
4637 rt_secondary
4638 .core()
4639 .attachments()
4640 .expect("core projection must expose the main attachment store");
4641 }
4642
4643 #[tokio::test]
4644 async fn record_plus_attachment_publication_rejects_a_secondary_runtime() {
4645 use khive_storage::{BlobStore as _, NewAttachment};
4646
4647 let main_arc = migrated_memory_backend();
4648 let secondary_arc = migrated_memory_backend();
4649 let rt_secondary =
4650 KhiveRuntime::from_backend(Arc::clone(&secondary_arc), secondary_config())
4651 .with_core_backend(Arc::clone(&main_arc));
4652 let blob_root = tempfile::tempdir().expect("blob root");
4653 let blob_store = Arc::new(
4654 khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
4655 .expect("blob store"),
4656 );
4657 let content_ref = blob_store.put(b"secondary-ref".to_vec()).await.unwrap();
4658 rt_secondary
4659 .install_blob_store(blob_store.clone())
4660 .expect("shared blob store");
4661 let token = rt_secondary.authorize(Namespace::local()).unwrap();
4662
4663 let error = rt_secondary
4664 .create_entity_with_attachments(
4665 &token,
4666 "artifact",
4667 Some("visual_asset"),
4668 "must route through core",
4669 None,
4670 None,
4671 vec![],
4672 vec![NewAttachment {
4673 role: "content".to_string(),
4674 content_ref: content_ref.clone(),
4675 media_type: None,
4676 size_bytes: Some(13),
4677 }],
4678 )
4679 .await
4680 .expect_err("secondary attachment publication must fail closed");
4681 assert!(error.to_string().contains("canonical main backend"));
4682 assert!(rt_secondary
4683 .list_entities(&token, None, None, 10, 0)
4684 .await
4685 .unwrap()
4686 .is_empty());
4687 assert!(rt_secondary
4688 .core()
4689 .list_entities(&token, None, None, 10, 0)
4690 .await
4691 .unwrap()
4692 .is_empty());
4693 assert!(
4694 blob_store.exists(&content_ref).await.unwrap(),
4695 "refusal must not mutate the already-published object"
4696 );
4697 }
4698
4699 #[tokio::test]
4700 async fn list_embedding_models_returns_empty_when_table_absent() {
4701 let rt = KhiveRuntime::memory().expect("memory runtime");
4704 let records = rt
4705 .list_embedding_models(None)
4706 .await
4707 .expect("list ok on empty table");
4708 assert!(records.is_empty());
4709 }
4710
4711 #[tokio::test]
4712 async fn list_embedding_models_returns_row_after_insert() {
4713 use khive_storage::{SqlStatement, SqlValue};
4714
4715 let rt = KhiveRuntime::memory().expect("memory runtime");
4716 let sql = rt.sql();
4717
4718 let now = 1_000_000i64;
4719 let id = uuid::Uuid::new_v4();
4720 let canonical_key = b"test_engine:test-model-v1:v1:384".to_vec();
4721
4722 let mut writer = sql.writer().await.expect("writer");
4723 writer
4724 .execute(SqlStatement {
4725 sql: "INSERT INTO _embedding_models \
4726 (id, engine_name, model_id, key_version, dim, output_dim, status, \
4727 activated_at, superseded_at, superseded_by, canonical_key, created_at) \
4728 VALUES (?1, ?2, ?3, ?4, ?5, NULL, ?6, ?7, NULL, NULL, ?8, ?9)"
4729 .into(),
4730 params: vec![
4731 SqlValue::Blob(id.as_bytes().to_vec()),
4732 SqlValue::Text("test_engine".into()),
4733 SqlValue::Text("test-model-v1".into()),
4734 SqlValue::Text("v1".into()),
4735 SqlValue::Integer(384),
4736 SqlValue::Text("active".into()),
4737 SqlValue::Integer(now),
4738 SqlValue::Blob(canonical_key),
4739 SqlValue::Integer(now),
4740 ],
4741 label: None,
4742 })
4743 .await
4744 .expect("insert row");
4745 drop(writer);
4746
4747 let records = rt.list_embedding_models(None).await.expect("list ok");
4748 assert_eq!(records.len(), 1);
4749 assert_eq!(records[0].engine_name, "test_engine");
4750 assert_eq!(records[0].model_id, "test-model-v1");
4751 assert_eq!(records[0].key_version, "v1");
4752 assert_eq!(records[0].dimensions, 384);
4753 assert_eq!(records[0].status, "active");
4754
4755 let filtered = rt
4757 .list_embedding_models(Some("test_engine"))
4758 .await
4759 .expect("filter ok");
4760 assert_eq!(filtered.len(), 1);
4761
4762 let no_match = rt
4764 .list_embedding_models(Some("other_engine"))
4765 .await
4766 .expect("no-match ok");
4767 assert!(no_match.is_empty());
4768 }
4769
4770 #[test]
4771 fn named_vector_identity_rejects_ambiguous_or_unsafe_values() {
4772 assert!(NamedVectorIdentity::new("", "model", 4).is_err());
4773 assert!(NamedVectorIdentity::new("bad-key", "model", 4).is_err());
4774 assert!(NamedVectorIdentity::new("valid_key", " model", 4).is_err());
4775 assert!(NamedVectorIdentity::new("valid_key", "model", 0).is_err());
4776 assert!(NamedVectorIdentity::new("valid_key", "model", 8193).is_err());
4777 assert!(NamedVectorIdentity::new("k".repeat(128), "m".repeat(512), 4).is_ok());
4778 assert!(NamedVectorIdentity::new("k".repeat(129), "model", 4).is_err());
4779 assert!(NamedVectorIdentity::new("valid_key", "m".repeat(513), 4).is_err());
4780 assert_eq!(
4781 NamedVectorIdentity::new("valid_key", "model", 4)
4782 .expect("valid identity")
4783 .dimensions(),
4784 4
4785 );
4786 }
4787
4788 #[tokio::test]
4789 async fn named_vector_store_rejects_dimension_or_model_key_rebinding() {
4790 let rt = KhiveRuntime::memory().expect("memory runtime");
4791 let token = rt.authorize(Namespace::local()).expect("authorize");
4792 let original = NamedVectorIdentity::new("visual_contract", "model-a", 4).unwrap();
4793 rt.vectors_for_named_identity(&token, &original)
4794 .await
4795 .expect("create named vector store");
4796 let registered = rt
4797 .list_embedding_models(Some("visual_contract"))
4798 .await
4799 .expect("list model registry");
4800 assert!(registered.iter().any(|record| {
4801 record.model_id == "model-a"
4802 && record.key_version == "visual_contract"
4803 && record.dimensions == 4
4804 }));
4805 let wrong_dimensions = NamedVectorIdentity::new("visual_contract", "model-a", 5).unwrap();
4806 let Err(dimension_error) = rt
4807 .vectors_for_named_identity(&token, &wrong_dimensions)
4808 .await
4809 else {
4810 panic!("same key cannot change dimensions");
4811 };
4812 assert!(dimension_error.to_string().contains("dimensions"));
4813
4814 let wrong_model = NamedVectorIdentity::new("visual_contract", "model-b", 4).unwrap();
4815 let Err(model_error) = rt.vectors_for_named_identity(&token, &wrong_model).await else {
4816 panic!("same key cannot change model identity");
4817 };
4818 assert!(model_error.to_string().contains("already bound"));
4819 }
4820
4821 #[tokio::test]
4822 async fn repeated_named_vector_lookup_avoids_writer_acquisition() {
4823 let runtime = KhiveRuntime::memory().expect("memory runtime");
4824 let token = runtime.authorize(Namespace::local()).expect("authorize");
4825 let identity = NamedVectorIdentity::new("visual_cached", "model-a", 4).unwrap();
4826 runtime
4827 .vectors_for_named_identity(&token, &identity)
4828 .await
4829 .expect("first lookup validates and registers");
4830 let writer_before = runtime.backend().pool().writer_acquisition_snapshot();
4831
4832 runtime
4833 .clone()
4834 .vectors_for_named_identity(&token, &identity)
4835 .await
4836 .expect("clone reuses verified store");
4837 assert_eq!(
4838 runtime.backend().pool().writer_acquisition_snapshot(),
4839 writer_before,
4840 "repeated reads must not reach vector-table setup or model registration"
4841 );
4842 }
4843
4844 #[tokio::test]
4845 async fn core_projection_reuses_main_named_vector_cache() {
4846 let main_backend = migrated_memory_backend();
4847 let main =
4848 KhiveRuntime::from_backend(Arc::clone(&main_backend), RuntimeConfig::no_embeddings());
4849 let secondary = KhiveRuntime::from_backend(migrated_memory_backend(), secondary_config())
4850 .with_core_embedders_from(&main)
4851 .with_core_backend(Arc::clone(&main_backend));
4852 let core = secondary.core();
4853 let token = core.authorize(Namespace::local()).expect("authorize");
4854 let identity = NamedVectorIdentity::new("core_visual_cached", "model-a", 4).unwrap();
4855 core.vectors_for_named_identity(&token, &identity)
4856 .await
4857 .expect("first lookup validates on main");
4858 let writer_before = main_backend.pool().writer_acquisition_snapshot();
4859
4860 secondary
4861 .core()
4862 .vectors_for_named_identity(&token, &identity)
4863 .await
4864 .expect("new core projection reuses main store");
4865 assert_eq!(
4866 main_backend.pool().writer_acquisition_snapshot(),
4867 writer_before
4868 );
4869 }
4870
4871 #[tokio::test]
4872 async fn rebound_core_named_vector_lookup_uses_new_backend() {
4873 let first_backend = migrated_memory_backend();
4874 let second_backend = migrated_memory_backend();
4875 let secondary = KhiveRuntime::from_backend(migrated_memory_backend(), secondary_config())
4876 .with_core_backend(Arc::clone(&first_backend));
4877 let identity = NamedVectorIdentity::new("rebound_visual", "model-a", 4).unwrap();
4878 let first_core = secondary.core();
4879 let first_token = first_core.authorize(Namespace::local()).expect("authorize");
4880 let first_store = first_core
4881 .vectors_for_named_identity(&first_token, &identity)
4882 .await
4883 .expect("first backend store");
4884
4885 let rebound = secondary.with_core_backend(Arc::clone(&second_backend));
4886 let second_core = rebound.core();
4887 let second_token = second_core
4888 .authorize(Namespace::local())
4889 .expect("authorize");
4890 let second_store = second_core
4891 .vectors_for_named_identity(&second_token, &identity)
4892 .await
4893 .expect("second backend store");
4894 assert!(
4895 !Arc::ptr_eq(&first_store, &second_store),
4896 "the second backend needs its own vector store"
4897 );
4898 let registered = second_core
4899 .list_embedding_models(Some("rebound_visual"))
4900 .await
4901 .expect("second backend model registry");
4902 assert!(registered.iter().any(|record| {
4903 record.model_id == "model-a"
4904 && record.key_version == "rebound_visual"
4905 && record.dimensions == 4
4906 }));
4907 }
4908
4909 #[test]
4910 fn rebound_core_discards_previous_main_embedder_wiring() {
4911 let first_backend = migrated_memory_backend();
4912 let second_backend = migrated_memory_backend();
4913 let first_main =
4914 KhiveRuntime::from_backend(Arc::clone(&first_backend), RuntimeConfig::no_embeddings());
4915 let second_main =
4916 KhiveRuntime::from_backend(Arc::clone(&second_backend), RuntimeConfig::no_embeddings());
4917 let secondary = KhiveRuntime::from_backend(migrated_memory_backend(), secondary_config())
4918 .with_core_embedders_from(&first_main)
4919 .with_core_backend(Arc::clone(&first_backend));
4920 assert!(secondary.core_embedders.is_some());
4921
4922 let rebound = secondary.with_core_backend(Arc::clone(&second_backend));
4923 assert!(
4924 rebound.core_embedders.is_none(),
4925 "the second backend cannot use the first main runtime's embedders"
4926 );
4927 assert!(Arc::ptr_eq(
4928 &rebound.core().embedder_registry,
4929 &rebound.embedder_registry
4930 ));
4931
4932 let rewired = rebound.with_core_embedders_from(&second_main);
4933 assert!(Arc::ptr_eq(
4934 &rewired.core().embedder_registry,
4935 &second_main.embedder_registry
4936 ));
4937 }
4938
4939 #[tokio::test]
4940 async fn concurrent_named_vector_first_bind_has_one_immutable_winner() {
4941 let rt = KhiveRuntime::memory().expect("memory runtime");
4942 let token = rt.authorize(Namespace::local()).expect("authorize");
4943 let first = NamedVectorIdentity::new("visual_race", "model-a", 4).unwrap();
4944 let second = NamedVectorIdentity::new("visual_race", "model-b", 4).unwrap();
4945
4946 let (first_result, second_result) = tokio::join!(
4947 rt.vectors_for_named_identity(&token, &first),
4948 rt.vectors_for_named_identity(&token, &second),
4949 );
4950 assert_ne!(
4951 first_result.is_ok(),
4952 second_result.is_ok(),
4953 "the active engine_name uniqueness rule must select exactly one first binding"
4954 );
4955
4956 let (winner, loser) = if first_result.is_ok() {
4957 (&first, &second)
4958 } else {
4959 (&second, &first)
4960 };
4961 rt.vectors_for_named_identity(&token, winner)
4962 .await
4963 .expect("winning identity remains idempotent");
4964 let error = match rt.vectors_for_named_identity(&token, loser).await {
4965 Ok(_) => panic!("losing identity cannot rebind the empty table"),
4966 Err(error) => error,
4967 };
4968 assert!(error.to_string().contains("already bound"));
4969
4970 let registered = rt
4971 .list_embedding_models(Some("visual_race"))
4972 .await
4973 .expect("list race registry");
4974 assert_eq!(registered.len(), 1);
4975 assert_eq!(registered[0].model_id, winner.model_name());
4976 }
4977
4978 #[tokio::test]
4979 async fn named_vector_registry_keeps_immutable_revisions_active_together() {
4980 let rt = KhiveRuntime::memory().expect("memory runtime");
4981 let token = rt.authorize(Namespace::local()).expect("authorize");
4982 let first = NamedVectorIdentity::new("visual_revision_a", "visual-model", 4).unwrap();
4983 let second = NamedVectorIdentity::new("visual_revision_b", "visual-model", 4).unwrap();
4984
4985 rt.vectors_for_named_identity(&token, &first)
4986 .await
4987 .expect("open first immutable space");
4988 rt.vectors_for_named_identity(&token, &second)
4989 .await
4990 .expect("open second immutable space");
4991
4992 let registered = rt.list_embedding_models(None).await.expect("list registry");
4993 assert!(registered.iter().any(|record| {
4994 record.engine_name == "visual_revision_a"
4995 && record.model_id == "visual-model"
4996 && record.key_version == "visual_revision_a"
4997 && record.status == "active"
4998 }));
4999 assert!(registered.iter().any(|record| {
5000 record.engine_name == "visual_revision_b"
5001 && record.model_id == "visual-model"
5002 && record.key_version == "visual_revision_b"
5003 && record.status == "active"
5004 }));
5005 }
5006}