1use std::collections::{HashMap, HashSet};
7use std::future::Future;
8use std::sync::{Arc, OnceLock, RwLock};
9
10use khive_db::StorageBackend;
11#[cfg(test)]
12use khive_gate::AllowAllGate;
13use khive_gate::GateRequest;
14use khive_storage::types::{SqlStatement, SqlValue};
15use khive_storage::{
16 AttachmentStore, EntityStore, Event, EventStore, GraphStore, NoteStore, SqlAccess, VectorStore,
17};
18use khive_types::{EdgeEndpointRule, EventKind, Namespace, SubstrateKind};
19use lattice_embed::{EmbeddingModel, EmbeddingService};
20
21use crate::config::{
22 build_embedder_registry, parse_embedding_model_alias, register_configured_embedding_models,
23 sanitize_key, vec_model_key,
24};
25use crate::error::{RuntimeError, RuntimeResult};
26use crate::pack::KindHook;
27
28#[cfg(all(test, target_os = "macos"))]
29const IN_PROCESS_TEST_NOFILE_LIMIT: libc::rlim_t = 4096;
30#[cfg(all(test, target_os = "macos"))]
31static IN_PROCESS_TEST_NOFILE_INIT: std::sync::Once = std::sync::Once::new();
32
33#[cfg(all(test, target_os = "macos"))]
34fn ensure_in_process_test_nofile_limit() {
35 IN_PROCESS_TEST_NOFILE_INIT.call_once(|| {
36 let mut limits = libc::rlimit {
37 rlim_cur: 0,
38 rlim_max: 0,
39 };
40 assert_eq!(unsafe { libc::getrlimit(libc::RLIMIT_NOFILE, &mut limits) }, 0);
43 assert!(
44 limits.rlim_max >= IN_PROCESS_TEST_NOFILE_LIMIT,
45 "in-process SQLite tests require a hard open-file limit of at least {IN_PROCESS_TEST_NOFILE_LIMIT}"
46 );
47 if limits.rlim_cur < IN_PROCESS_TEST_NOFILE_LIMIT {
48 limits.rlim_cur = IN_PROCESS_TEST_NOFILE_LIMIT;
49 assert_eq!(unsafe { libc::setrlimit(libc::RLIMIT_NOFILE, &limits) }, 0);
52 }
53 });
54}
55
56tokio::task_local! {
57 static REQUEST_EMBEDDER_EXCLUSIONS: Arc<HashSet<String>>;
58}
59
60pub fn scope_request_embedder_exclusions<F>(
62 excluded: Vec<String>,
63 future: F,
64) -> impl Future<Output = F::Output>
65where
66 F: Future,
67{
68 REQUEST_EMBEDDER_EXCLUSIONS.scope(Arc::new(excluded.into_iter().collect()), future)
69}
70
71pub fn inherit_request_embedder_scope<F>(future: F) -> impl Future<Output = F::Output>
73where
74 F: Future,
75{
76 let exclusions = REQUEST_EMBEDDER_EXCLUSIONS.try_with(Arc::clone).ok();
77 async move {
78 match exclusions {
79 Some(exclusions) => REQUEST_EMBEDDER_EXCLUSIONS.scope(exclusions, future).await,
80 None => future.await,
81 }
82 }
83}
84
85fn request_excludes_embedder(name: &str) -> bool {
86 let canonical = parse_embedding_model_alias(name)
87 .map(|model| model.to_string())
88 .unwrap_or_else(|| name.to_string());
89 REQUEST_EMBEDDER_EXCLUSIONS
90 .try_with(|excluded| excluded.contains(&canonical))
91 .unwrap_or(false)
92}
93
94pub type EntityTypeValidatorFn =
100 Arc<dyn Fn(&str, Option<&str>) -> Result<Option<String>, RuntimeError> + Send + Sync>;
101
102pub type EntityKindHooks = Vec<(String, Arc<dyn KindHook>)>;
109
110pub type NoteMutationHookFn = Arc<
121 dyn Fn(String, uuid::Uuid) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>
122 + Send
123 + Sync,
124>;
125
126pub type NoteWriteValidatorFn = Arc<
138 dyn Fn(&str, &str, Option<serde_json::Value>) -> Result<Option<serde_json::Value>, RuntimeError>
139 + Send
140 + Sync,
141>;
142
143#[derive(Clone, Debug, Eq, PartialEq)]
150pub struct NamedVectorIdentity {
151 model_key: String,
152 model_name: String,
153 dimensions: usize,
154}
155
156struct CachedNamedVectorStores {
157 identity: NamedVectorIdentity,
158 by_namespace: HashMap<String, Arc<dyn VectorStore>>,
159}
160
161fn check_cached_named_vector_identity(
162 cached: &NamedVectorIdentity,
163 requested: &NamedVectorIdentity,
164) -> RuntimeResult<()> {
165 if cached.dimensions() != requested.dimensions() {
166 return Err(RuntimeError::InvalidInput(format!(
167 "named vector model_key {:?} is already bound to {} dimensions, expected {}",
168 requested.model_key(),
169 cached.dimensions(),
170 requested.dimensions()
171 )));
172 }
173 if cached.model_name() != requested.model_name() {
174 return Err(RuntimeError::InvalidInput(format!(
175 "named vector model_key {:?} is already bound to a different active model identity",
176 requested.model_key()
177 )));
178 }
179 Ok(())
180}
181
182impl NamedVectorIdentity {
183 const MAX_MODEL_KEY_BYTES: usize = 128;
184 const MAX_MODEL_NAME_BYTES: usize = 512;
185
186 pub fn new(
188 model_key: impl Into<String>,
189 model_name: impl Into<String>,
190 dimensions: usize,
191 ) -> RuntimeResult<Self> {
192 let model_key = model_key.into();
193 let model_name = model_name.into();
194 if model_key.is_empty()
195 || model_key.len() > Self::MAX_MODEL_KEY_BYTES
196 || !model_key
197 .chars()
198 .all(|c| c.is_ascii_alphanumeric() || c == '_')
199 {
200 return Err(RuntimeError::InvalidInput(format!(
201 "named vector model_key must be 1..={} bytes of ASCII alphanumeric/underscore",
202 Self::MAX_MODEL_KEY_BYTES
203 )));
204 }
205 if model_name.trim().is_empty()
206 || model_name.trim() != model_name
207 || model_name.len() > Self::MAX_MODEL_NAME_BYTES
208 {
209 return Err(RuntimeError::InvalidInput(format!(
210 "named vector model_name must be 1..={} bytes with no surrounding whitespace",
211 Self::MAX_MODEL_NAME_BYTES
212 )));
213 }
214 if !(1..=8192).contains(&dimensions) {
215 return Err(RuntimeError::InvalidInput(format!(
216 "named vector dimensions must be in 1..=8192, got {dimensions}"
217 )));
218 }
219 Ok(Self {
220 model_key,
221 model_name,
222 dimensions,
223 })
224 }
225
226 pub fn model_key(&self) -> &str {
227 &self.model_key
228 }
229
230 pub fn model_name(&self) -> &str {
231 &self.model_name
232 }
233
234 pub fn dimensions(&self) -> usize {
235 self.dimensions
236 }
237}
238
239pub use crate::config::{
240 assert_captured_db_anchor_consistent, assert_db_anchor_consistent, expand_tilde,
241 parse_pack_list, resolve_db_anchor, resolve_project_actor_id, runtime_config_from_khive_config,
242 BackendId, BackendIdError, NamespaceToken, RuntimeConfig,
243};
244
245#[derive(Clone)]
254struct CoreEmbedderState {
255 registry: Arc<std::sync::RwLock<crate::embedder_registry::EmbedderRegistry>>,
256 default_embedder_name: Arc<str>,
257 embedding_model: Option<EmbeddingModel>,
258 additional_embedding_models: Vec<EmbeddingModel>,
259}
260
261#[derive(Clone)]
263pub struct KhiveRuntime {
264 backend: Arc<StorageBackend>,
265 named_vector_stores: Arc<RwLock<HashMap<String, CachedNamedVectorStores>>>,
269 core_named_vector_stores: Option<Arc<RwLock<HashMap<String, CachedNamedVectorStores>>>>,
272 core_backend: Option<Arc<StorageBackend>>,
276 config: RuntimeConfig,
277 ann_fresh_tail_enabled: bool,
281 embedder_registry: Arc<std::sync::RwLock<crate::embedder_registry::EmbedderRegistry>>,
288 default_embedder_name: Arc<str>,
289 core_embedders: Option<CoreEmbedderState>,
296 edge_rules: Arc<RwLock<Vec<EdgeEndpointRule>>>,
300 valid_entity_kinds: Arc<RwLock<Vec<String>>>,
308 valid_note_kinds: Arc<RwLock<Vec<String>>>,
309 entity_type_validator: Arc<RwLock<Option<EntityTypeValidatorFn>>>,
316 note_mutation_hook: Arc<RwLock<Option<NoteMutationHookFn>>>,
328 note_write_validator: Arc<RwLock<Option<NoteWriteValidatorFn>>>,
336 entity_kind_hooks: Arc<RwLock<EntityKindHooks>>,
350 pack_owned_note_kinds: Arc<RwLock<Vec<String>>>,
358 blob_hydrator: Arc<OnceLock<Arc<crate::blob::BlobHydrator>>>,
365 fusion_executors: Arc<RwLock<HashMap<String, Arc<dyn crate::fusion::FusionExecutor>>>>,
376}
377
378impl KhiveRuntime {
379 pub fn new(config: RuntimeConfig) -> RuntimeResult<Self> {
389 Self::new_with_file_backend(config, |path| StorageBackend::sqlite(path))
390 }
391
392 #[cfg(any(test, feature = "test-internals"))]
394 pub fn new_for_test(config: RuntimeConfig) -> RuntimeResult<Self> {
395 Self::new_with_file_backend(config, |path| StorageBackend::sqlite_for_test(path))
396 }
397
398 fn new_with_file_backend(
399 config: RuntimeConfig,
400 open_file: impl FnOnce(&std::path::Path) -> Result<StorageBackend, khive_db::SqliteError>,
401 ) -> RuntimeResult<Self> {
402 #[cfg(all(test, target_os = "macos"))]
403 ensure_in_process_test_nofile_limit();
404 let backend = match &config.db_path {
405 Some(path) => {
406 if let Some(parent) = path.parent() {
407 std::fs::create_dir_all(parent).ok();
408 }
409 open_file(path)?
410 }
411 None => StorageBackend::memory()?,
412 };
413 let schema_version = backend.prepare_core_schema()?;
417 if schema_version < khive_db::migrations::ATTACHMENT_CUTOVER_VERSION {
418 return Err(khive_db::SqliteError::InvalidData(
419 "database requires the host application-assisted V21 attachment cutover; \
420 start through khive-mcp/kkernel boot instead of constructing KhiveRuntime \
421 directly"
422 .into(),
423 )
424 .into());
425 }
426 if !backend.is_read_only() {
427 register_configured_embedding_models(&backend, &config)?;
428 }
429 Ok(Self::assemble_from_backend(Arc::new(backend), config))
430 }
431
432 pub fn new_readonly(config: RuntimeConfig) -> RuntimeResult<Self> {
439 Self::new_readonly_with_file_backend(config, |path| StorageBackend::sqlite_read_only(path))
440 }
441
442 #[cfg(any(test, feature = "test-internals"))]
444 pub fn new_readonly_for_test(config: RuntimeConfig) -> RuntimeResult<Self> {
445 Self::new_readonly_with_file_backend(config, |path| {
446 StorageBackend::sqlite_read_only_for_test(path)
447 })
448 }
449
450 fn new_readonly_with_file_backend(
451 config: RuntimeConfig,
452 open_file: impl FnOnce(&std::path::Path) -> Result<StorageBackend, khive_db::SqliteError>,
453 ) -> RuntimeResult<Self> {
454 #[cfg(all(test, target_os = "macos"))]
455 ensure_in_process_test_nofile_limit();
456 let backend = match &config.db_path {
457 Some(path) => open_file(path)?,
458 None => StorageBackend::memory()?,
459 };
460 backend.prepare_core_schema()?;
461 Ok(Self::assemble_from_backend(Arc::new(backend), config))
462 }
463
464 pub fn from_backend(backend: Arc<StorageBackend>, config: RuntimeConfig) -> Self {
477 if !backend.is_read_only() {
478 if let Err(err) = register_configured_embedding_models(&backend, &config) {
479 tracing::warn!(error = %err, "failed to register configured embedding models");
480 }
481 }
482 Self::assemble_from_backend(backend, config)
483 }
484
485 pub fn from_prepared_backend(
492 backend: Arc<StorageBackend>,
493 config: RuntimeConfig,
494 ) -> RuntimeResult<Self> {
495 if backend.attachment_cutover_status()?
496 != khive_db::migrations::AttachmentCutoverStatus::Complete
497 {
498 return Err(khive_db::SqliteError::InvalidData(
499 "from_prepared_backend requires a complete V21 attachment cutover".into(),
500 )
501 .into());
502 }
503 if !backend.is_read_only() {
504 register_configured_embedding_models(&backend, &config)?;
505 }
506 Ok(Self::assemble_from_backend(backend, config))
507 }
508
509 fn assemble_from_backend(backend: Arc<StorageBackend>, config: RuntimeConfig) -> Self {
510 if config.backend_id.as_str() == BackendId::MAIN {
511 backend.pool().main_pool_generation();
512 }
513 let ann_fresh_tail_enabled = crate::config::ann_fresh_tail_enabled_from_env();
514 let (registry, default_embedder_name) = build_embedder_registry(&config);
515 Self {
516 backend,
517 named_vector_stores: Arc::new(RwLock::new(HashMap::new())),
518 core_named_vector_stores: None,
519 core_backend: None,
520 config,
521 ann_fresh_tail_enabled,
522 embedder_registry: Arc::new(std::sync::RwLock::new(registry)),
523 default_embedder_name,
524 core_embedders: None,
525 edge_rules: Arc::new(RwLock::new(Vec::new())),
526 valid_entity_kinds: Arc::new(RwLock::new(Vec::new())),
527 valid_note_kinds: Arc::new(RwLock::new(Vec::new())),
528 entity_type_validator: Arc::new(RwLock::new(None)),
529 note_mutation_hook: Arc::new(RwLock::new(None)),
530 note_write_validator: Arc::new(RwLock::new(None)),
531 entity_kind_hooks: Arc::new(RwLock::new(Vec::new())),
532 pack_owned_note_kinds: Arc::new(RwLock::new(Vec::new())),
533 blob_hydrator: Arc::new(OnceLock::new()),
534 fusion_executors: Arc::new(RwLock::new(HashMap::new())),
535 }
536 }
537
538 pub fn with_core_backend(mut self, core: Arc<StorageBackend>) -> Self {
547 debug_assert_ne!(
548 self.config.backend_id.as_str(),
549 BackendId::MAIN,
550 "with_core_backend must not be called on the main runtime"
551 );
552 core.pool().main_pool_generation();
553 if self.core_named_vector_stores.is_none() {
554 self.core_named_vector_stores = Some(Arc::new(RwLock::new(HashMap::new())));
555 }
556 self.core_backend = Some(core);
557 self
558 }
559
560 pub fn with_core_embedders_from(mut self, main: &KhiveRuntime) -> Self {
567 debug_assert!(
568 main.core_backend.is_none(),
569 "with_core_embedders_from takes the MAIN runtime"
570 );
571 self.core_embedders = Some(CoreEmbedderState {
572 registry: main.embedder_registry.clone(),
573 default_embedder_name: main.default_embedder_name.clone(),
574 embedding_model: main.config.embedding_model,
575 additional_embedding_models: main.config.additional_embedding_models.clone(),
576 });
577 self.core_named_vector_stores = Some(main.named_vector_stores.clone());
578 if Arc::ptr_eq(&self.backend, &main.backend) {
579 self.named_vector_stores = main.named_vector_stores.clone();
580 }
581 self
582 }
583
584 pub fn core(&self) -> KhiveRuntime {
604 match &self.core_backend {
605 None => match &self.core_embedders {
610 None => self.clone(),
611 Some(core_embedders) => {
612 let mut core = self.clone();
613 core.config.embedding_model = core_embedders.embedding_model;
614 core.config.additional_embedding_models =
615 core_embedders.additional_embedding_models.clone();
616 core.embedder_registry = core_embedders.registry.clone();
617 core.default_embedder_name = core_embedders.default_embedder_name.clone();
618 core.core_embedders = None;
619 core
620 }
621 },
622 Some(main_arc) => {
623 let mut core_config = self.config.clone();
624 core_config.backend_id = BackendId::main();
625 let (embedder_registry, default_embedder_name) = match &self.core_embedders {
631 Some(core_embedders) => {
632 core_config.embedding_model = core_embedders.embedding_model;
633 core_config.additional_embedding_models =
634 core_embedders.additional_embedding_models.clone();
635 (
636 core_embedders.registry.clone(),
637 core_embedders.default_embedder_name.clone(),
638 )
639 }
640 None => (
641 self.embedder_registry.clone(),
642 self.default_embedder_name.clone(),
643 ),
644 };
645 KhiveRuntime {
646 backend: main_arc.clone(),
647 named_vector_stores: self
648 .core_named_vector_stores
649 .clone()
650 .unwrap_or_else(|| Arc::new(RwLock::new(HashMap::new()))),
651 core_named_vector_stores: None,
652 core_backend: None,
653 config: core_config,
654 ann_fresh_tail_enabled: self.ann_fresh_tail_enabled,
655 embedder_registry,
656 default_embedder_name,
657 core_embedders: None,
658 edge_rules: self.edge_rules.clone(),
659 valid_entity_kinds: self.valid_entity_kinds.clone(),
660 valid_note_kinds: self.valid_note_kinds.clone(),
661 entity_type_validator: self.entity_type_validator.clone(),
662 note_mutation_hook: self.note_mutation_hook.clone(),
663 note_write_validator: self.note_write_validator.clone(),
664 entity_kind_hooks: self.entity_kind_hooks.clone(),
665 pack_owned_note_kinds: self.pack_owned_note_kinds.clone(),
666 blob_hydrator: self.blob_hydrator.clone(),
667 fusion_executors: self.fusion_executors.clone(),
668 }
669 }
670 }
671 }
672
673 pub fn memory() -> RuntimeResult<Self> {
675 Self::new(RuntimeConfig {
676 db_path: None,
677 packs: vec!["kg".to_string()],
678 brain_profile: None,
679 actor_id: None,
680 ..RuntimeConfig::no_embeddings()
681 })
682 }
683
684 pub fn backend_id(&self) -> &BackendId {
689 &self.config.backend_id
690 }
691
692 pub fn visible_namespaces(&self) -> &[Namespace] {
699 &self.config.visible_namespaces
700 }
701
702 pub fn config(&self) -> &RuntimeConfig {
704 &self.config
705 }
706
707 pub fn vector_arm_selected(&self) -> bool {
712 self.config.embedding_model.is_some()
713 }
714
715 pub fn ann_fresh_tail_enabled(&self) -> bool {
718 self.ann_fresh_tail_enabled
719 }
720
721 pub fn with_ann_fresh_tail_enabled(mut self, enabled: bool) -> Self {
727 self.ann_fresh_tail_enabled = enabled;
728 self
729 }
730
731 pub fn backend(&self) -> &StorageBackend {
741 &self.backend
742 }
743
744 pub fn is_read_only(&self) -> bool {
747 self.backend.is_read_only()
748 }
749
750 pub fn backend_data_dir(&self) -> Option<std::path::PathBuf> {
753 self.backend.data_dir()
754 }
755
756 pub fn backend_ann_root(&self) -> Option<std::path::PathBuf> {
761 self.backend.ann_root()
762 }
763
764 pub async fn db_diagnostics(&self) -> RuntimeResult<khive_db::diagnostics::DbDiagnostics> {
778 self.db_diagnostics_with_audit_metrics(None).await
786 }
787
788 pub async fn db_diagnostics_with_audit_metrics(
794 &self,
795 runtime_audit_batch_metrics: Option<khive_db::diagnostics::RuntimeAuditBatchMetrics>,
796 ) -> RuntimeResult<khive_db::diagnostics::DbDiagnostics> {
797 let pool = self.core().backend.pool_arc();
798 let legacy_sweep_interval = khive_db::SessionSweepConfig::default().interval;
801 let build_hash = crate::build_info::BUILD_INFO
802 .is_stamped()
803 .then_some(crate::build_info::BUILD_INFO.source_revision);
804 let build = khive_db::diagnostics::BuildIdentity::from_env(
805 crate::build_info::PACKAGE_VERSION,
806 build_hash,
807 );
808
809 let mut report = khive_db::diagnostics::collect_with_runtime_audit_metrics_interruptibly(
810 pool,
811 build,
812 legacy_sweep_interval,
813 crate::pack::audit_append_failure_count(),
814 runtime_audit_batch_metrics,
815 )
816 .await
817 .map_err(RuntimeError::from)?;
818 report.writer_contention.audit_obligation_append_failures =
819 Some(crate::pack::audit_obligation_append_failure_count());
820 report
821 .writer_contention
822 .audit_obligation_append_failures_unavailable_reason = None;
823 Ok(report)
824 }
825
826 pub fn entities(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn EntityStore>> {
830 Ok(self
831 .backend
832 .entities_for_namespace(token.namespace().as_str())?)
833 }
834
835 pub fn graph(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn GraphStore>> {
837 Ok(self
838 .backend
839 .graph_for_namespace(token.namespace().as_str())?)
840 }
841
842 pub fn notes(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn NoteStore>> {
855 Ok(crate::note_store_guard::PolicyEnforcingNoteStore::wrap(
856 self.raw_notes(token)?,
857 ))
858 }
859
860 pub(crate) fn raw_notes(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn NoteStore>> {
870 Ok(self
871 .backend
872 .notes_for_namespace(token.namespace().as_str())?)
873 }
874
875 pub fn attachments(&self) -> RuntimeResult<Arc<dyn AttachmentStore>> {
882 if self.config.backend_id.as_str() != BackendId::MAIN {
883 return Err(RuntimeError::InvalidInput(format!(
884 "attachments are owned by the canonical main backend; runtime backend {:?} must route through KhiveRuntime::core()",
885 self.config.backend_id.as_str()
886 )));
887 }
888 Ok(self.backend.attachments()?)
889 }
890
891 pub fn events(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn EventStore>> {
906 Ok(crate::event_store_guard::AttributedEventStore::wrap(
907 self.raw_events_for_namespace(token.namespace().as_str())?,
908 token,
909 ))
910 }
911
912 pub(crate) fn raw_events_for_namespace(
916 &self,
917 namespace: &str,
918 ) -> RuntimeResult<Arc<dyn EventStore>> {
919 let legacy = self.backend.events_for_namespace(namespace)?;
920 match &self.config.events_split {
921 None => Ok(legacy),
922 Some(split) => {
923 if self.backend.is_read_only() {
935 if !split.db_path.exists() {
936 return Ok(legacy);
937 }
938 let lane = crate::events_split::direct_backend_with_max_readers(
939 &split.db_path,
940 true,
941 Some(self.backend.pool().config().max_readers),
942 )?
943 .events_for_namespace(namespace)?;
944 return Ok(Arc::new(crate::events_split::SplitEventStore::new(
945 legacy, lane,
946 )));
947 }
948 let lane: Arc<dyn EventStore> = match &split.socket_path {
949 #[cfg(unix)]
950 Some(socket) => {
951 let client = crate::events_split::client_for(socket)?;
952 Arc::new(crate::events_split::ForwardingEventStore::new(
953 namespace, client,
954 ))
955 }
956 #[cfg(not(unix))]
957 Some(_socket) => {
958 return Err(RuntimeError::InvalidInput(
959 "events-daemon socket forwarding requires a Unix platform; \
960 configure the events split in direct mode here"
961 .to_string(),
962 ));
963 }
964 None => crate::events_split::direct_backend_with_max_readers(
965 &split.db_path,
966 false,
967 Some(self.backend.pool().config().max_readers),
968 )?
969 .events_for_namespace(namespace)?,
970 };
971 Ok(Arc::new(crate::events_split::SplitEventStore::new(
972 legacy, lane,
973 )))
974 }
975 }
976 }
977
978 pub fn sql(&self) -> Arc<dyn SqlAccess> {
980 self.backend.sql()
981 }
982
983 pub fn events_sidecar_sql_read_only(&self) -> RuntimeResult<Option<Arc<dyn SqlAccess>>> {
1001 match &self.config.events_split {
1002 None => Ok(None),
1003 Some(split) => {
1004 if !split.db_path.exists() {
1005 return Ok(None);
1006 }
1007 let backend = if self.backend.is_read_only() {
1008 crate::events_split::direct_backend_with_max_readers(
1009 &split.db_path,
1010 true,
1011 Some(self.backend.pool().config().max_readers),
1012 )?
1013 } else {
1014 crate::events_split::direct_backend_with_max_readers(
1015 &split.db_path,
1016 false,
1017 Some(self.backend.pool().config().max_readers),
1018 )?
1019 };
1020 Ok(Some(backend.sql()))
1021 }
1022 }
1023 }
1024
1025 pub fn vectors(
1029 &self,
1030 token: &NamespaceToken,
1031 ) -> RuntimeResult<Arc<dyn khive_storage::VectorStore>> {
1032 let model = self.resolve_embedding_model(None)?;
1033 self.vectors_for_embedding_model(token, model)
1034 }
1035
1036 pub fn vectors_for_model(
1044 &self,
1045 token: &NamespaceToken,
1046 model_name: &str,
1047 ) -> RuntimeResult<Arc<dyn khive_storage::VectorStore>> {
1048 let (model_name, dims) = self.vector_model_metadata(model_name)?;
1049 Ok(self.backend.vectors_for_namespace(
1050 &sanitize_key(&model_name),
1051 &model_name,
1052 dims,
1053 token.namespace().as_str(),
1054 )?)
1055 }
1056
1057 pub(crate) fn vector_model_metadata(&self, model_name: &str) -> RuntimeResult<(String, usize)> {
1060 if request_excludes_embedder(model_name) {
1061 return Err(crate::RuntimeError::UnknownModel(model_name.to_string()));
1062 }
1063 let registry = self
1064 .embedder_registry
1065 .read()
1066 .map_err(|_| crate::RuntimeError::Internal("embedder registry lock poisoned".into()))?;
1067 if let Some(model) = parse_embedding_model_alias(model_name) {
1068 let key = model.to_string();
1071 if registry.contains(&key) {
1072 return Ok((key, model.dimensions()));
1073 }
1074 }
1075 registry
1076 .get_provider(model_name)
1077 .map(|provider| (model_name.to_owned(), provider.dimensions()))
1078 .ok_or_else(|| crate::RuntimeError::UnknownModel(model_name.to_string()))
1079 }
1080
1081 pub async fn vectors_for_named_identity(
1089 &self,
1090 token: &NamespaceToken,
1091 identity: &NamedVectorIdentity,
1092 ) -> RuntimeResult<Arc<dyn VectorStore>> {
1093 let namespace = token.namespace().as_str();
1094 {
1095 let cached = self.named_vector_stores.read().map_err(|_| {
1096 RuntimeError::Internal("named vector store cache lock poisoned".into())
1097 })?;
1098 if let Some(entry) = cached.get(identity.model_key()) {
1099 check_cached_named_vector_identity(&entry.identity, identity)?;
1100 if let Some(store) = entry.by_namespace.get(namespace) {
1101 return Ok(Arc::clone(store));
1102 }
1103 }
1104 }
1105 let store = self.backend.vectors_for_namespace(
1106 identity.model_key(),
1107 identity.model_name(),
1108 identity.dimensions(),
1109 namespace,
1110 )?;
1111
1112 let table = format!("vec_{}", identity.model_key());
1113 let mut reader = self.sql().reader().await?;
1114 let dimension_row = reader
1115 .query_row(SqlStatement {
1116 sql: "SELECT sql FROM sqlite_schema WHERE type = 'table' AND name = ?1".to_string(),
1117 params: vec![SqlValue::Text(table.clone())],
1118 label: Some("runtime_named_vector_dimension".to_string()),
1119 })
1120 .await?
1121 .ok_or_else(|| {
1122 RuntimeError::Internal(format!(
1123 "named vector table {table} has no sqlite_schema declaration"
1124 ))
1125 })?;
1126 let table_ddl = match dimension_row.get("sql") {
1127 Some(SqlValue::Text(value)) => value,
1128 other => {
1129 return Err(RuntimeError::Internal(format!(
1130 "named vector table {table} returned invalid schema metadata: {other:?}"
1131 )))
1132 }
1133 };
1134 let declared_dimensions = vector_dimensions_from_ddl(table_ddl).ok_or_else(|| {
1135 RuntimeError::Internal(format!(
1136 "named vector table {table} has no parseable embedding dimension"
1137 ))
1138 })?;
1139 if declared_dimensions != identity.dimensions() {
1140 return Err(RuntimeError::InvalidInput(format!(
1141 "named vector model_key {:?} is already bound to {declared_dimensions} dimensions, expected {}",
1142 identity.model_key(),
1143 identity.dimensions()
1144 )));
1145 }
1146
1147 let stored_models = reader
1148 .query_all(SqlStatement {
1149 sql: format!(
1150 "SELECT DISTINCT embedding_model FROM {table} ORDER BY embedding_model LIMIT 2"
1151 ),
1152 params: vec![],
1153 label: Some("runtime_named_vector_model_identity".to_string()),
1154 })
1155 .await?;
1156 for row in stored_models {
1157 let stored = match row.get("embedding_model") {
1158 Some(SqlValue::Text(value)) => value,
1159 other => {
1160 return Err(RuntimeError::Internal(format!(
1161 "named vector table {table} returned invalid model identity metadata: {other:?}"
1162 )))
1163 }
1164 };
1165 if stored != identity.model_name() {
1166 return Err(RuntimeError::InvalidInput(format!(
1167 "named vector model_key {:?} already contains model {stored:?}, cannot bind it to {:?}",
1168 identity.model_key(),
1169 identity.model_name()
1170 )));
1171 }
1172 }
1173
1174 self.backend
1175 .register_embedding_model(
1176 identity.model_key(),
1177 identity.model_name(),
1178 identity.model_key(),
1179 identity.dimensions() as u32,
1180 )
1181 .map_err(|error| {
1182 if matches!(
1183 &error,
1184 khive_db::SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
1185 if code.code == rusqlite::ErrorCode::ConstraintViolation
1186 ) {
1187 RuntimeError::InvalidInput(format!(
1188 "named vector model_key {:?} is already bound to a different active model identity",
1189 identity.model_key()
1190 ))
1191 } else {
1192 RuntimeError::Sqlite(error)
1193 }
1194 })?;
1195
1196 let mut cached = self
1197 .named_vector_stores
1198 .write()
1199 .map_err(|_| RuntimeError::Internal("named vector store cache lock poisoned".into()))?;
1200 let entry = cached
1201 .entry(identity.model_key().to_owned())
1202 .or_insert_with(|| CachedNamedVectorStores {
1203 identity: identity.clone(),
1204 by_namespace: HashMap::new(),
1205 });
1206 check_cached_named_vector_identity(&entry.identity, identity)?;
1207 entry
1208 .by_namespace
1209 .insert(namespace.to_owned(), Arc::clone(&store));
1210 Ok(store)
1211 }
1212
1213 pub fn embedder_dimensions(&self, model_name: &str) -> Option<usize> {
1220 if request_excludes_embedder(model_name) {
1221 return None;
1222 }
1223 if let Some(model) = parse_embedding_model_alias(model_name) {
1224 let key = model.to_string();
1225 let in_registry = self
1226 .embedder_registry
1227 .read()
1228 .map(|reg| reg.contains(&key))
1229 .unwrap_or(false);
1230 if in_registry {
1231 return Some(model.dimensions());
1232 }
1233 }
1234 self.embedder_registry
1235 .read()
1236 .ok()?
1237 .get_provider(model_name)
1238 .map(|p| p.dimensions())
1239 }
1240
1241 fn vectors_for_embedding_model(
1242 &self,
1243 token: &NamespaceToken,
1244 model: EmbeddingModel,
1245 ) -> RuntimeResult<Arc<dyn khive_storage::VectorStore>> {
1246 Ok(self.backend.vectors_for_namespace(
1247 &vec_model_key(model),
1248 &model.to_string(),
1249 model.dimensions(),
1250 token.namespace().as_str(),
1251 )?)
1252 }
1253
1254 pub fn text(
1256 &self,
1257 token: &NamespaceToken,
1258 ) -> RuntimeResult<Arc<dyn khive_storage::TextSearch>> {
1259 let _ = token;
1260 Ok(self.backend.text("entities")?)
1261 }
1262
1263 pub fn text_for_notes(
1265 &self,
1266 token: &NamespaceToken,
1267 ) -> RuntimeResult<Arc<dyn khive_storage::TextSearch>> {
1268 let _ = token;
1269 Ok(self.backend.text("notes")?)
1270 }
1271
1272 pub fn authorize(&self, ns: Namespace) -> RuntimeResult<NamespaceToken> {
1288 let actor = crate::actor_identity::resolve_actor(self.config.actor_id.as_deref());
1289 let req = GateRequest::new(
1290 actor.clone(),
1291 ns.clone(),
1292 "authorize",
1293 serde_json::Value::Null,
1294 );
1295 match self.config.gate.check(&req) {
1296 Ok(ref decision) if decision.is_allow() => {
1297 if let khive_gate::GateDecision::Allow { ref obligations } = decision {
1298 if !obligations.is_empty() {
1299 tracing::debug!(
1300 namespace = %ns.as_str(),
1301 "authorize: obligations={:?}",
1302 obligations
1303 );
1304 }
1305 }
1306 Ok(NamespaceToken::mint_authorized(ns, actor))
1307 }
1308 Ok(khive_gate::GateDecision::Deny { reason }) => {
1309 Err(crate::RuntimeError::permission_denied("authorize", reason))
1310 }
1311 Ok(_) => Err(crate::RuntimeError::permission_denied(
1312 "authorize",
1313 "gate denied",
1314 )),
1315 Err(e) => {
1316 tracing::warn!(
1317 namespace = %ns.as_str(),
1318 error = %crate::secret_gate::bounded_masked_log_text(&e.to_string()),
1319 "authorize: gate check failed (fail-closed)"
1320 );
1321 Err(crate::RuntimeError::Internal(format!(
1322 "gate error: {}",
1323 e.wire_reason()
1324 )))
1325 }
1326 }
1327 }
1328
1329 pub fn authorize_with_visibility(
1344 &self,
1345 primary: Namespace,
1346 extra_visible: Vec<Namespace>,
1347 ) -> RuntimeResult<NamespaceToken> {
1348 let actor = crate::actor_identity::resolve_actor(self.config.actor_id.as_deref());
1349 let req = GateRequest::new(
1350 actor.clone(),
1351 primary.clone(),
1352 "authorize",
1353 serde_json::Value::Null,
1354 );
1355 match self.config.gate.check(&req) {
1356 Ok(ref decision) if decision.is_allow() => {
1357 if let khive_gate::GateDecision::Allow { ref obligations } = decision {
1358 if !obligations.is_empty() {
1359 tracing::debug!(
1360 namespace = %primary.as_str(),
1361 "authorize_with_visibility: obligations={:?}",
1362 obligations
1363 );
1364 }
1365 }
1366 for extra in &extra_visible {
1373 let extra_req = GateRequest::new(
1374 actor.clone(),
1375 extra.clone(),
1376 "authorize.visible",
1377 serde_json::Value::Null,
1378 );
1379 match self.config.gate.check(&extra_req) {
1380 Ok(ref extra_decision) if extra_decision.is_allow() => {}
1381 Ok(khive_gate::GateDecision::Deny { reason }) => {
1382 return Err(crate::RuntimeError::permission_denied(
1383 "authorize",
1384 format!(
1385 "visibility namespace {:?} denied: {reason}",
1386 extra.as_str()
1387 ),
1388 ));
1389 }
1390 Ok(_) => {
1391 return Err(crate::RuntimeError::permission_denied(
1392 "authorize",
1393 format!("visibility namespace {:?} denied by gate", extra.as_str()),
1394 ));
1395 }
1396 Err(e) => {
1397 tracing::warn!(
1398 namespace = %extra.as_str(),
1399 error = %crate::secret_gate::bounded_masked_log_text(&e.to_string()),
1400 "authorize_with_visibility: extra-namespace gate check failed (fail-closed)"
1401 );
1402 return Err(crate::RuntimeError::Internal(format!(
1403 "gate error: {}",
1404 e.wire_reason()
1405 )));
1406 }
1407 }
1408 }
1409 Ok(NamespaceToken::mint_with_visibility(
1410 primary,
1411 extra_visible,
1412 actor,
1413 ))
1414 }
1415 Ok(khive_gate::GateDecision::Deny { reason }) => {
1416 Err(crate::RuntimeError::permission_denied("authorize", reason))
1417 }
1418 Ok(_) => Err(crate::RuntimeError::permission_denied(
1419 "authorize",
1420 "gate denied",
1421 )),
1422 Err(e) => {
1423 tracing::warn!(
1424 namespace = %primary.as_str(),
1425 error = %crate::secret_gate::bounded_masked_log_text(&e.to_string()),
1426 "authorize_with_visibility: gate check failed (fail-closed)"
1427 );
1428 Err(crate::RuntimeError::Internal(format!(
1429 "gate error: {}",
1430 e.wire_reason()
1431 )))
1432 }
1433 }
1434 }
1435
1436 pub fn install_edge_rules(&self, rules: Vec<EdgeEndpointRule>) {
1442 if let Ok(mut guard) = self.edge_rules.write() {
1443 *guard = rules;
1444 }
1445 }
1446
1447 pub fn install_blob_hydrator(
1453 &self,
1454 hydrator: Arc<crate::blob::BlobHydrator>,
1455 ) -> RuntimeResult<()> {
1456 if hydrator.budget_bytes() != self.config.blob_hydration_bytes {
1457 return Err(RuntimeError::InvalidInput(format!(
1458 "blob hydrator budget {} does not match this runtime's resolved budget {}",
1459 hydrator.budget_bytes(),
1460 self.config.blob_hydration_bytes
1461 )));
1462 }
1463 if self.is_read_only() && !hydrator.enforces_read_only() {
1472 return Err(RuntimeError::InvalidInput(
1473 "this runtime is read-only: install the raw store with install_blob_store, \
1474 which wraps it so every physical mutator refuses"
1475 .to_string(),
1476 ));
1477 }
1478 self.install_blob_hydrator_slot(hydrator)
1479 }
1480
1481 pub fn install_shared_blob_hydrator(
1504 &self,
1505 hydrator: Arc<crate::blob::BlobHydrator>,
1506 ) -> RuntimeResult<()> {
1507 if !hydrator.is_governed() {
1508 return Err(RuntimeError::InvalidInput(
1509 "the shared install seam accepts only hydrators whose mode was derived from a \
1510 governing backend (BlobHydrator::resolve_for_governing_backend); use \
1511 install_blob_hydrator for a hand-paired hydrator"
1512 .to_string(),
1513 ));
1514 }
1515 if hydrator.budget_bytes() != self.config.blob_hydration_bytes {
1516 return Err(RuntimeError::InvalidInput(format!(
1517 "blob hydrator budget {} does not match this runtime's resolved budget {}",
1518 hydrator.budget_bytes(),
1519 self.config.blob_hydration_bytes
1520 )));
1521 }
1522 self.install_blob_hydrator_slot(hydrator)
1523 }
1524
1525 fn is_same_blob_pairing(
1532 current: &Arc<crate::blob::BlobHydrator>,
1533 candidate: &Arc<crate::blob::BlobHydrator>,
1534 ) -> bool {
1535 Arc::ptr_eq(current, candidate)
1536 || (Arc::ptr_eq(¤t.raw_store(), &candidate.raw_store())
1537 && current.budget_bytes() == candidate.budget_bytes()
1538 && current.enforces_read_only() == candidate.enforces_read_only())
1539 }
1540
1541 fn install_blob_hydrator_slot(
1546 &self,
1547 hydrator: Arc<crate::blob::BlobHydrator>,
1548 ) -> RuntimeResult<()> {
1549 if let Some(current) = self.blob_hydrator.get() {
1550 return if Self::is_same_blob_pairing(current, &hydrator) {
1551 Ok(())
1552 } else {
1553 Err(RuntimeError::InvalidInput(
1554 "a different blob hydrator is already installed".to_string(),
1555 ))
1556 };
1557 }
1558
1559 match self.blob_hydrator.set(hydrator) {
1560 Ok(()) => Ok(()),
1561 Err(candidate) => {
1562 let current = self.blob_hydrator.get().ok_or_else(|| {
1563 RuntimeError::Internal(
1564 "blob hydrator install raced without a visible winner".to_string(),
1565 )
1566 })?;
1567 if Self::is_same_blob_pairing(current, &candidate) {
1568 Ok(())
1569 } else {
1570 Err(RuntimeError::InvalidInput(
1571 "a different blob hydrator is already installed".to_string(),
1572 ))
1573 }
1574 }
1575 }
1576 }
1577
1578 pub fn install_blob_store(
1584 &self,
1585 store: Arc<dyn khive_storage::BlobStore>,
1586 ) -> RuntimeResult<()> {
1587 if let Some(current) = self.blob_hydrator.get() {
1588 let current_store = current.store();
1593 if Arc::ptr_eq(¤t_store, &store) || Arc::ptr_eq(¤t.raw_store(), &store) {
1594 return Ok(());
1595 }
1596 }
1597 let hydrator = if self.is_read_only() {
1603 crate::blob::BlobHydrator::new_read_only(store, self.config.blob_hydration_bytes)?
1604 } else {
1605 crate::blob::BlobHydrator::new(store, self.config.blob_hydration_bytes)?
1606 };
1607 self.install_blob_hydrator(Arc::new(hydrator))
1608 }
1609
1610 pub fn blob_hydrator(&self) -> Option<Arc<crate::blob::BlobHydrator>> {
1612 self.blob_hydrator.get().cloned()
1613 }
1614
1615 pub fn blob_store(&self) -> Option<Arc<dyn khive_storage::BlobStore>> {
1620 self.blob_hydrator.get().map(|hydrator| hydrator.store())
1621 }
1622
1623 pub fn install_kind_registry(&self, entity_kinds: Vec<String>, note_kinds: Vec<String>) {
1633 if let Ok(mut guard) = self.valid_entity_kinds.write() {
1634 *guard = entity_kinds;
1635 }
1636 if let Ok(mut guard) = self.valid_note_kinds.write() {
1637 *guard = note_kinds;
1638 }
1639 }
1640
1641 pub fn install_pack_owned_note_kinds(&self, kinds: Vec<String>) {
1646 if let Ok(mut guard) = self.pack_owned_note_kinds.write() {
1647 *guard = kinds;
1648 }
1649 }
1650
1651 pub fn is_pack_owned_note_kind(&self, kind: &str) -> bool {
1657 self.pack_owned_note_kinds
1658 .read()
1659 .map(|g| g.iter().any(|k| k == kind))
1660 .unwrap_or(false)
1661 }
1662
1663 pub(crate) fn validate_entity_kind(&self, kind: &str) -> crate::RuntimeResult<()> {
1668 let guard = self.valid_entity_kinds.read().map_err(|_| {
1669 crate::RuntimeError::Internal("entity kind registry lock poisoned".into())
1670 })?;
1671 if guard.is_empty() {
1672 return Ok(());
1673 }
1674 if guard.iter().any(|k| k == kind) {
1675 Ok(())
1676 } else {
1677 Err(crate::RuntimeError::InvalidInput(format!(
1678 "unknown entity kind {kind:?}; valid: {}",
1679 guard.join(", ")
1680 )))
1681 }
1682 }
1683
1684 pub(crate) fn validate_note_kind(&self, kind: &str) -> crate::RuntimeResult<()> {
1689 let guard = self.valid_note_kinds.read().map_err(|_| {
1690 crate::RuntimeError::Internal("note kind registry lock poisoned".into())
1691 })?;
1692 if guard.is_empty() {
1693 return Ok(());
1694 }
1695 if guard.iter().any(|k| k == kind) {
1696 Ok(())
1697 } else {
1698 Err(crate::RuntimeError::InvalidInput(format!(
1699 "unknown note kind {kind:?}; valid: {}",
1700 guard.join(", ")
1701 )))
1702 }
1703 }
1704
1705 pub fn install_entity_type_validator(&self, f: EntityTypeValidatorFn) {
1715 if let Ok(mut guard) = self.entity_type_validator.write() {
1716 *guard = Some(f);
1717 }
1718 }
1719
1720 pub(crate) fn validate_entity_type_for_kind(
1725 &self,
1726 kind: &str,
1727 entity_type: Option<&str>,
1728 ) -> crate::RuntimeResult<Option<String>> {
1729 let guard = self.entity_type_validator.read().map_err(|_| {
1730 crate::RuntimeError::Internal("entity type validator lock poisoned".into())
1731 })?;
1732 match guard.as_ref() {
1733 None => Ok(entity_type.map(str::to_string)),
1734 Some(validate) => validate(kind, entity_type),
1735 }
1736 }
1737
1738 pub fn install_note_mutation_hook(&self, f: NoteMutationHookFn) {
1746 if let Ok(mut guard) = self.note_mutation_hook.write() {
1747 *guard = Some(f);
1748 }
1749 }
1750
1751 pub fn install_entity_kind_hooks(&self, hooks: EntityKindHooks) {
1758 if let Ok(mut guard) = self.entity_kind_hooks.write() {
1759 *guard = hooks;
1760 }
1761 }
1762
1763 pub(crate) fn entity_kind_hook(&self, kind: &str) -> Option<Arc<dyn KindHook>> {
1771 self.entity_kind_hooks.read().ok().and_then(|guard| {
1772 guard
1773 .iter()
1774 .find(|(k, _)| k == kind)
1775 .map(|(_, hook)| hook.clone())
1776 })
1777 }
1778
1779 pub fn register_fusion_strategy(
1836 &self,
1837 name: impl Into<String>,
1838 executor: Arc<dyn crate::fusion::FusionExecutor>,
1839 ) {
1840 if let Ok(mut guard) = self.fusion_executors.write() {
1841 guard.insert(name.into(), executor);
1842 }
1843 }
1844
1845 pub(crate) fn fusion_executor(
1852 &self,
1853 name: &str,
1854 ) -> RuntimeResult<Arc<dyn crate::fusion::FusionExecutor>> {
1855 let guard = self
1856 .fusion_executors
1857 .read()
1858 .map_err(|_| RuntimeError::Internal("fusion executor registry lock poisoned".into()))?;
1859 guard
1860 .get(name)
1861 .cloned()
1862 .ok_or_else(|| RuntimeError::UnknownFusionStrategy(name.to_string()))
1863 }
1864
1865 pub fn install_note_write_validator(&self, f: NoteWriteValidatorFn) {
1866 if let Ok(mut guard) = self.note_write_validator.write() {
1867 *guard = Some(f);
1868 }
1869 }
1870
1871 pub fn has_note_write_validator(&self) -> bool {
1879 self.note_write_validator
1880 .read()
1881 .map(|g| g.is_some())
1882 .unwrap_or(false)
1883 }
1884
1885 pub(crate) fn derive_note_write_properties(
1890 &self,
1891 kind: &str,
1892 token: &NamespaceToken,
1893 properties: Option<serde_json::Value>,
1894 ) -> RuntimeResult<Option<serde_json::Value>> {
1895 let validator = self
1896 .note_write_validator
1897 .read()
1898 .map_err(|_| RuntimeError::Internal("note write validator lock poisoned".into()))?
1899 .clone();
1900 match validator {
1901 None => Ok(properties),
1902 Some(validate) => validate(kind, &token.actor().id, properties),
1903 }
1904 }
1905
1906 pub(crate) async fn fire_note_mutation_hook(&self, kind: &str, id: uuid::Uuid) {
1915 let hook = self
1916 .note_mutation_hook
1917 .read()
1918 .ok()
1919 .and_then(|guard| guard.clone());
1920 if let Some(hook) = hook {
1921 hook(kind.to_string(), id).await;
1922 }
1923 }
1924
1925 pub fn pack_edge_rules(&self) -> Vec<EdgeEndpointRule> {
1934 self.edge_rules
1935 .read()
1936 .map(|g| g.clone())
1937 .unwrap_or_default()
1938 }
1939
1940 pub(crate) fn with_pack_edge_rules<T>(&self, f: impl FnOnce(&[EdgeEndpointRule]) -> T) -> T {
1942 match self.edge_rules.read() {
1943 Ok(rules) => f(&rules),
1944 Err(_) => f(&[]),
1945 }
1946 }
1947
1948 pub fn default_embedder_name(&self) -> &str {
1950 self.default_embedder_name.as_ref()
1951 }
1952
1953 pub fn resolve_embedding_model(&self, name: Option<&str>) -> RuntimeResult<EmbeddingModel> {
1958 let model = match name {
1959 Some(raw) => parse_embedding_model_alias(raw)
1960 .ok_or_else(|| crate::RuntimeError::UnknownModel(raw.to_string()))?,
1961 None => self
1962 .config
1963 .embedding_model
1964 .ok_or_else(|| crate::RuntimeError::Unconfigured("embedding_model".into()))?,
1965 };
1966 let key = model.to_string();
1967 if request_excludes_embedder(&key) {
1968 return Err(crate::RuntimeError::UnknownModel(
1969 name.unwrap_or_else(|| self.default_embedder_name())
1970 .to_string(),
1971 ));
1972 }
1973 let contains = self
1974 .embedder_registry
1975 .read()
1976 .map(|reg| reg.contains(&key))
1977 .unwrap_or(false);
1978 if contains {
1979 Ok(model)
1980 } else {
1981 Err(crate::RuntimeError::UnknownModel(
1982 name.unwrap_or_else(|| self.default_embedder_name())
1983 .to_string(),
1984 ))
1985 }
1986 }
1987
1988 pub fn registered_embedding_model_names(&self) -> Vec<String> {
1995 self.embedder_registry
1996 .read()
1997 .map(|reg| {
1998 reg.names()
1999 .into_iter()
2000 .filter(|name| !request_excludes_embedder(name))
2001 .collect()
2002 })
2003 .unwrap_or_default()
2004 }
2005
2006 pub async fn embedder(&self, name: &str) -> RuntimeResult<Arc<dyn EmbeddingService>> {
2019 self.embedder_inner(name, None).await
2020 }
2021
2022 pub(crate) async fn embedder_with_token(
2023 &self,
2024 token: &NamespaceToken,
2025 name: &str,
2026 ) -> RuntimeResult<Arc<dyn EmbeddingService>> {
2027 self.embedder_inner(name, Some(token)).await
2028 }
2029
2030 async fn embedder_inner(
2031 &self,
2032 name: &str,
2033 token: Option<&NamespaceToken>,
2034 ) -> RuntimeResult<Arc<dyn EmbeddingService>> {
2035 let canonical_key = match parse_embedding_model_alias(name) {
2038 Some(model) => model.to_string(),
2039 None => name.to_owned(),
2040 };
2041 if request_excludes_embedder(&canonical_key) {
2042 return Err(crate::RuntimeError::UnknownModel(name.to_string()));
2043 }
2044 let entry = {
2047 let registry = self.embedder_registry.read().map_err(|_| {
2048 crate::RuntimeError::Internal("embedder registry lock poisoned".into())
2049 })?;
2050 registry
2051 .get_entry(&canonical_key)
2052 .ok_or_else(|| crate::RuntimeError::UnknownModel(name.to_string()))?
2053 };
2054 let (service, init_duration_us) = entry.resolve().await?;
2055 if let Some(duration_us) = init_duration_us {
2056 if let Some(token) = token {
2057 self.emit_embedder_initialized(token, &canonical_key, duration_us)
2058 .await;
2059 } else if let Ok(token) = self.authorize(self.config.default_namespace.clone()) {
2060 self.emit_embedder_initialized(&token, &canonical_key, duration_us)
2061 .await;
2062 }
2063 }
2064 Ok(service)
2065 }
2066
2067 async fn emit_embedder_initialized(
2068 &self,
2069 token: &NamespaceToken,
2070 model_name: &str,
2071 duration_us: i64,
2072 ) {
2073 if self.is_read_only() {
2077 return;
2078 }
2079 let Ok(store) = self.events(token) else {
2080 return;
2081 };
2082 let event = Event::new(
2083 token.namespace().as_str(),
2084 "embedder.init",
2085 EventKind::EmbedderInitialized,
2086 SubstrateKind::Event,
2087 format!("{}:{}", token.actor().kind, token.actor().id),
2088 )
2089 .with_payload(serde_json::json!({
2090 "model_name": model_name,
2091 "duration_us": duration_us,
2092 }))
2093 .with_duration_us(duration_us);
2094 if let Err(err) = store.append_event(event).await {
2095 tracing::warn!(error = %err, model_name, "embedder initialization event append failed");
2096 }
2097 }
2098
2099 pub fn register_embedder(
2112 &self,
2113 provider: impl crate::embedder_registry::EmbedderProvider + 'static,
2114 ) {
2115 if let Ok(mut registry) = self.embedder_registry.write() {
2116 registry.register(provider);
2117 } else {
2118 tracing::warn!(
2119 "embedder registry lock poisoned — embedder {} not registered",
2120 std::any::type_name::<dyn crate::embedder_registry::EmbedderProvider>()
2121 );
2122 }
2123 }
2124
2125 pub async fn list_embedding_models(
2132 &self,
2133 engine_filter: Option<&str>,
2134 ) -> RuntimeResult<Vec<khive_db::EmbeddingModelRegistryRecord>> {
2135 use khive_storage::{SqlStatement, SqlValue};
2136
2137 let (sql_text, params) = if let Some(engine) = engine_filter {
2138 (
2139 "SELECT engine_name, model_id, key_version, dim, status, \
2140 activated_at, superseded_at \
2141 FROM _embedding_models WHERE engine_name = ?1 \
2142 ORDER BY engine_name, activated_at IS NULL, activated_at"
2143 .to_string(),
2144 vec![SqlValue::Text(engine.to_string())],
2145 )
2146 } else {
2147 (
2148 "SELECT engine_name, model_id, key_version, dim, status, \
2149 activated_at, superseded_at \
2150 FROM _embedding_models \
2151 ORDER BY engine_name, activated_at IS NULL, activated_at"
2152 .to_string(),
2153 vec![],
2154 )
2155 };
2156
2157 let stmt = SqlStatement {
2158 sql: sql_text,
2159 params,
2160 label: Some("list_embedding_models".into()),
2161 };
2162
2163 let mut reader = self
2164 .sql()
2165 .reader()
2166 .await
2167 .map_err(crate::RuntimeError::Storage)?;
2168
2169 let rows = match reader.query_all(stmt).await {
2170 Ok(rows) => rows,
2171 Err(e) if e.to_string().contains("no such table: _embedding_models") => {
2172 return Ok(Vec::new())
2173 }
2174 Err(e) => return Err(crate::RuntimeError::Storage(e)),
2175 };
2176
2177 let mut records = Vec::with_capacity(rows.len());
2178 for row in rows {
2179 macro_rules! required_text {
2180 ($col:expr) => {
2181 match row.get($col) {
2182 Some(SqlValue::Text(s)) => s.clone(),
2183 other => {
2184 tracing::warn!(column = $col, value = ?other, "skipping registry row: unexpected type");
2185 continue;
2186 }
2187 }
2188 };
2189 }
2190 let engine_name = required_text!("engine_name");
2191 let model_id = required_text!("model_id");
2192 let key_version = required_text!("key_version");
2193 let dimensions = match row.get("dim") {
2194 Some(SqlValue::Integer(n)) => match u32::try_from(*n) {
2195 Ok(d) => d,
2196 Err(_) => {
2197 tracing::warn!(dim = n, "skipping registry row: dim out of u32 range");
2198 continue;
2199 }
2200 },
2201 other => {
2202 tracing::warn!(column = "dim", value = ?other, "skipping registry row: unexpected type");
2203 continue;
2204 }
2205 };
2206 let status = required_text!("status");
2207 let activated_at = match row.get("activated_at") {
2208 Some(SqlValue::Integer(n)) => Some(*n),
2209 _ => None,
2210 };
2211 let superseded_at = match row.get("superseded_at") {
2212 Some(SqlValue::Integer(n)) => Some(*n),
2213 _ => None,
2214 };
2215 records.push(khive_db::EmbeddingModelRegistryRecord {
2216 engine_name,
2217 model_id,
2218 key_version,
2219 dimensions,
2220 status,
2221 activated_at,
2222 superseded_at,
2223 });
2224 }
2225
2226 Ok(records)
2227 }
2228}
2229
2230fn vector_dimensions_from_ddl(ddl: &str) -> Option<usize> {
2231 let lower = ddl.to_ascii_lowercase();
2232 let suffix = lower.split_once("embedding float[")?.1;
2233 let dimension = suffix.split_once(']')?.0;
2234 if dimension.is_empty() || !dimension.bytes().all(|byte| byte.is_ascii_digit()) {
2235 return None;
2236 }
2237 dimension.parse().ok()
2238}
2239
2240#[cfg(test)]
2244mod tests {
2245 use super::*;
2246 use khive_gate::GateRef;
2247 use serial_test::serial;
2248
2249 #[cfg(target_os = "macos")]
2250 #[test]
2251 fn in_process_runtime_tests_have_4096_open_file_slots() {
2252 let _runtime = KhiveRuntime::memory().expect("test runtime");
2253 let mut limits = libc::rlimit {
2254 rlim_cur: 0,
2255 rlim_max: 0,
2256 };
2257 assert_eq!(
2259 unsafe { libc::getrlimit(libc::RLIMIT_NOFILE, &mut limits) },
2260 0
2261 );
2262 assert!(
2263 limits.rlim_cur >= IN_PROCESS_TEST_NOFILE_LIMIT,
2264 "a parallel runtime suite needs at least 4096 open-file slots"
2265 );
2266 }
2267
2268 fn test_blob_hydrator() -> (tempfile::TempDir, Arc<crate::BlobHydrator>) {
2269 let root = tempfile::tempdir().expect("blob root");
2270 let store = Arc::new(
2271 khive_db::stores::blob::FsBlobStore::new(root.path().to_path_buf(), 0)
2272 .expect("fs blob store"),
2273 );
2274 let hydrator = Arc::new(
2275 crate::BlobHydrator::new(store, crate::DEFAULT_BLOB_HYDRATION_BYTES)
2276 .expect("blob hydrator"),
2277 );
2278 (root, hydrator)
2279 }
2280
2281 #[test]
2282 fn memory_runtime_creates_successfully() {
2283 let rt = KhiveRuntime::memory().expect("memory runtime should create");
2284 assert!(rt.config().db_path.is_none());
2285 }
2286
2287 #[test]
2288 fn installed_blob_hydrator_is_shared_by_clone_and_core_handles() {
2289 let main_backend = Arc::new(StorageBackend::memory().expect("main backend"));
2290 let pack_backend = Arc::new(StorageBackend::memory().expect("pack backend"));
2291 let mut config = RuntimeConfig::no_embeddings();
2292 config.backend_id = BackendId::parse("assets").expect("valid backend id");
2293 let runtime = KhiveRuntime::from_backend(pack_backend, config)
2294 .with_core_backend(Arc::clone(&main_backend));
2295 let (_root, hydrator) = test_blob_hydrator();
2296
2297 runtime
2298 .install_blob_hydrator(Arc::clone(&hydrator))
2299 .expect("first install");
2300
2301 for handle in [runtime.clone(), runtime.core()] {
2302 let installed = handle.blob_hydrator().expect("installed hydrator");
2303 assert!(Arc::ptr_eq(&installed, &hydrator));
2304 }
2305 }
2306
2307 #[test]
2308 fn blob_hydrator_install_is_idempotent_but_rejects_replacement() {
2309 let runtime = KhiveRuntime::memory().expect("runtime");
2310 let (_first_root, first) = test_blob_hydrator();
2311 let (_second_root, second) = test_blob_hydrator();
2312
2313 runtime
2314 .install_blob_hydrator(Arc::clone(&first))
2315 .expect("first install");
2316 runtime
2317 .install_blob_hydrator(Arc::clone(&first))
2318 .expect("same Arc reinstall is idempotent");
2319
2320 let error = runtime
2321 .install_blob_hydrator(second)
2322 .expect_err("a different hydrator must not replace the installed pair");
2323 assert!(error.to_string().contains("already installed"));
2324 assert!(Arc::ptr_eq(
2325 &runtime.blob_hydrator().expect("original remains"),
2326 &first
2327 ));
2328 }
2329
2330 #[test]
2331 fn blob_hydrator_install_rejects_a_budget_that_disagrees_with_runtime_identity() {
2332 let runtime = KhiveRuntime::memory().expect("runtime");
2333 let root = tempfile::tempdir().expect("blob root");
2334 let store = Arc::new(
2335 khive_db::stores::blob::FsBlobStore::new(root.path().to_path_buf(), 0)
2336 .expect("fs blob store"),
2337 );
2338 let mismatched = Arc::new(
2339 crate::BlobHydrator::new(store, khive_storage::MAX_BLOB_WHOLE_BYTES)
2340 .expect("minimum blob budget"),
2341 );
2342
2343 let error = runtime
2344 .install_blob_hydrator(mismatched)
2345 .expect_err("live admission must match the construction-baked config identity");
2346 assert!(matches!(error, RuntimeError::InvalidInput(_)));
2347 assert!(runtime.blob_hydrator().is_none());
2348 }
2349
2350 #[test]
2351 fn fresh_tail_policy_is_instance_scoped_and_clone_stable() {
2352 let enabled = KhiveRuntime::memory()
2353 .expect("enabled memory runtime")
2354 .with_ann_fresh_tail_enabled(true);
2355 let disabled = KhiveRuntime::memory()
2356 .expect("disabled memory runtime")
2357 .with_ann_fresh_tail_enabled(false);
2358
2359 assert!(enabled.ann_fresh_tail_enabled());
2360 assert!(enabled.clone().ann_fresh_tail_enabled());
2361 assert!(!disabled.ann_fresh_tail_enabled());
2362 assert!(!disabled.clone().ann_fresh_tail_enabled());
2363 }
2364
2365 #[tokio::test]
2366 async fn runtime_db_diagnostics_supplies_both_contention_counter_sources() {
2367 let rt = KhiveRuntime::memory().expect("memory runtime should create");
2368
2369 let report = rt.db_diagnostics().await.expect("diagnostics succeed");
2370
2371 assert!(
2372 report.writer_contention.writer_acquisitions >= 1,
2373 "runtime construction runs migrations through the finite-wait pooled writer"
2374 );
2375 assert_eq!(
2376 report.writer_contention.writer_acquisitions,
2377 report
2378 .writer_contention
2379 .pooled_writer_acquisitions
2380 .saturating_add(report.writer_contention.standalone_writer_acquisitions)
2381 .saturating_add(report.writer_contention.writer_task_acquisitions),
2382 "the public aggregate must equal the class-specific snapshot"
2383 );
2384 assert!(
2385 report.writer_contention.audit_append_failures.is_some(),
2386 "the runtime path must supply its process-wide swallowed-audit counter"
2387 );
2388 assert!(report
2389 .writer_contention
2390 .audit_obligation_append_failures
2391 .is_some());
2392 assert!(report
2393 .writer_contention
2394 .audit_obligation_append_failures_unavailable_reason
2395 .is_none());
2396 assert!(report
2397 .writer_contention
2398 .audit_append_failures_unavailable_reason
2399 .is_none());
2400 }
2401
2402 #[test]
2403 fn backend_data_dir_returns_none_for_memory_backend() {
2404 let rt = KhiveRuntime::memory().expect("memory runtime");
2405 assert!(rt.backend_data_dir().is_none());
2406 }
2407
2408 #[test]
2409 fn backend_data_dir_returns_parent_dir_for_file_backend() {
2410 let dir = tempfile::tempdir().unwrap();
2411 let path = dir.path().join("test.db");
2412 let config = RuntimeConfig {
2413 web: Default::default(),
2414 telemetry: Default::default(),
2415 mounts: Vec::new(),
2416 brain: Default::default(),
2417 git_write: Default::default(),
2418 display_timezone: chrono_tz::Tz::UTC,
2419 events_split: None,
2420 db_path: Some(path),
2421 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2422 default_namespace: Namespace::local(),
2423 embedding_model: None,
2424 additional_embedding_models: vec![],
2425 gate: Arc::new(AllowAllGate),
2426 packs: vec!["kg".to_string()],
2427 backend_id: BackendId::main(),
2428 brain_profile: None,
2429 visible_namespaces: vec![],
2430 allowed_outbound_namespaces: vec![],
2431 actor_id: None,
2432 exec: Default::default(),
2433 };
2434 let rt = KhiveRuntime::new_for_test(config).expect("file runtime");
2435 let data_dir = rt
2436 .backend_data_dir()
2437 .expect("file backend must return Some");
2438 assert_eq!(data_dir, dir.path());
2439 }
2440
2441 #[tokio::test]
2447 async fn resolve_prefix_finds_sidecar_only_event() {
2448 let dir = tempfile::tempdir().unwrap();
2449 let _registry_guard = crate::events_split::TestRegistryGuard::new(dir.path());
2450 let sidecar_path = dir.path().join("main.db.events.db");
2451 let config = RuntimeConfig {
2452 web: Default::default(),
2453 telemetry: Default::default(),
2454 mounts: Vec::new(),
2455 brain: Default::default(),
2456 git_write: Default::default(),
2457 display_timezone: chrono_tz::Tz::UTC,
2458 events_split: Some(crate::events_split::EventsSplitConfig {
2459 db_path: sidecar_path.clone(),
2460 socket_path: None,
2461 }),
2462 db_path: Some(dir.path().join("main.db")),
2463 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2464 default_namespace: Namespace::local(),
2465 embedding_model: None,
2466 additional_embedding_models: vec![],
2467 gate: Arc::new(AllowAllGate),
2468 packs: vec!["kg".to_string()],
2469 backend_id: BackendId::main(),
2470 brain_profile: None,
2471 visible_namespaces: vec![],
2472 allowed_outbound_namespaces: vec![],
2473 actor_id: None,
2474 exec: Default::default(),
2475 };
2476 let rt = KhiveRuntime::new_for_test(config).expect("file runtime");
2477
2478 let event = khive_storage::Event::new(
2479 "local",
2480 "memory.recall",
2481 khive_types::EventKind::RecallExecuted,
2482 khive_types::SubstrateKind::Note,
2483 "agent:test",
2484 );
2485 let event_id = event.id;
2486 let prefix = event_id.to_string()[..8].to_string();
2487
2488 assert_eq!(
2489 rt.resolve_prefix_unfiltered(&prefix)
2490 .await
2491 .expect("pre-insert resolve"),
2492 None,
2493 "control: prefix must miss before the lane row exists"
2494 );
2495
2496 let lane = crate::events_split::direct_backend_for(&sidecar_path)
2497 .expect("lane backend")
2498 .events_for_namespace("local")
2499 .expect("lane store");
2500 lane.append_event(event).await.expect("lane append");
2501
2502 assert_eq!(
2503 rt.resolve_prefix_unfiltered(&prefix)
2504 .await
2505 .expect("post-insert resolve"),
2506 Some(event_id),
2507 "a sidecar-only event id must resolve by hex prefix"
2508 );
2509 }
2510
2511 #[test]
2512 fn backend_data_dir_returns_none_for_from_backend_with_memory() {
2513 let backend = Arc::new(StorageBackend::memory().expect("memory backend"));
2514 let config = RuntimeConfig {
2515 web: Default::default(),
2516 telemetry: Default::default(),
2517 mounts: Vec::new(),
2518 brain: Default::default(),
2519 git_write: Default::default(),
2520 display_timezone: chrono_tz::Tz::UTC,
2521 events_split: None,
2522 db_path: None,
2523 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2524 default_namespace: Namespace::local(),
2525 embedding_model: None,
2526 additional_embedding_models: vec![],
2527 gate: Arc::new(AllowAllGate),
2528 packs: vec!["kg".to_string()],
2529 backend_id: BackendId::main(),
2530 brain_profile: None,
2531 visible_namespaces: vec![],
2532 allowed_outbound_namespaces: vec![],
2533 actor_id: None,
2534 exec: Default::default(),
2535 };
2536 let rt = KhiveRuntime::from_backend(backend, config);
2537 assert!(rt.backend_data_dir().is_none());
2538 }
2539
2540 #[test]
2541 fn file_runtime_creates_successfully() {
2542 let dir = tempfile::tempdir().unwrap();
2543 let path = dir.path().join("test.db");
2544 let config = RuntimeConfig {
2545 web: Default::default(),
2546 telemetry: Default::default(),
2547 mounts: Vec::new(),
2548 brain: Default::default(),
2549 git_write: Default::default(),
2550 display_timezone: chrono_tz::Tz::UTC,
2551 events_split: None,
2552 db_path: Some(path.clone()),
2553 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2554 default_namespace: Namespace::parse("test").unwrap(),
2555 embedding_model: None,
2556 additional_embedding_models: vec![],
2557 gate: Arc::new(AllowAllGate),
2558 packs: vec!["kg".to_string()],
2559 backend_id: BackendId::main(),
2560 brain_profile: None,
2561 visible_namespaces: vec![],
2562 allowed_outbound_namespaces: vec![],
2563 actor_id: None,
2564 exec: Default::default(),
2565 };
2566 let rt = KhiveRuntime::new_for_test(config).expect("file runtime should create");
2567 assert!(path.exists());
2568 assert_eq!(rt.config().default_namespace.as_str(), "test");
2569 }
2570
2571 #[cfg(unix)]
2572 #[tokio::test]
2573 async fn normal_boot_detects_read_only_snapshot_and_skips_model_registration() {
2574 use std::os::unix::fs::PermissionsExt;
2575
2576 let dir = tempfile::tempdir().unwrap();
2577 let path = dir.path().join("read_only_runtime.db");
2578 let base = RuntimeConfig {
2579 web: Default::default(),
2580 telemetry: Default::default(),
2581 mounts: Vec::new(),
2582 brain: Default::default(),
2583 git_write: Default::default(),
2584 display_timezone: chrono_tz::Tz::UTC,
2585 events_split: None,
2586 db_path: Some(path.clone()),
2587 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2588 default_namespace: Namespace::local(),
2589 embedding_model: None,
2590 additional_embedding_models: vec![],
2591 gate: Arc::new(AllowAllGate),
2592 packs: vec!["kg".to_string()],
2593 backend_id: BackendId::main(),
2594 brain_profile: None,
2595 visible_namespaces: vec![],
2596 allowed_outbound_namespaces: vec![],
2597 actor_id: None,
2598 exec: Default::default(),
2599 };
2600 {
2601 let writable =
2602 KhiveRuntime::new_for_test(base.clone()).expect("create migrated snapshot");
2603 assert!(writable
2604 .list_embedding_models(None)
2605 .await
2606 .expect("registry query")
2607 .is_empty());
2608 }
2609
2610 let mut permissions = std::fs::metadata(&path).unwrap().permissions();
2611 permissions.set_mode(0o444);
2612 std::fs::set_permissions(&path, permissions).unwrap();
2613 khive_storage::test_support::freeze_snapshot_sidecars(&path);
2617
2618 let read_only_config = RuntimeConfig {
2619 embedding_model: Some(EmbeddingModel::AllMiniLmL6V2),
2620 ..base
2621 };
2622 let runtime = KhiveRuntime::new_for_test(read_only_config)
2623 .expect("read-only boot must validate instead of migrating/registering");
2624 assert!(runtime.is_read_only());
2625 assert_eq!(
2626 runtime.backend().pool().writer_acquisition_snapshot(),
2627 khive_db::pool::WriterAcquisitionSnapshot::default(),
2628 "the construction-inclusive acquisition baseline must stay at zero"
2629 );
2630 assert!(
2631 runtime
2632 .list_embedding_models(None)
2633 .await
2634 .expect("read-only registry query")
2635 .is_empty(),
2636 "configured models must remain in-memory only during read-only boot"
2637 );
2638 }
2639
2640 #[test]
2641 fn explicit_readonly_constructor_uses_read_only_pool_even_on_writable_file_mode() {
2642 let dir = tempfile::tempdir().unwrap();
2643 let path = dir.path().join("explicit_read_only_runtime.db");
2644 let config = RuntimeConfig {
2645 web: Default::default(),
2646 telemetry: Default::default(),
2647 mounts: Vec::new(),
2648 brain: Default::default(),
2649 git_write: Default::default(),
2650 display_timezone: chrono_tz::Tz::UTC,
2651 events_split: None,
2652 db_path: Some(path.clone()),
2653 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2654 default_namespace: Namespace::local(),
2655 embedding_model: None,
2656 additional_embedding_models: vec![],
2657 gate: Arc::new(AllowAllGate),
2658 packs: vec!["kg".to_string()],
2659 backend_id: BackendId::main(),
2660 brain_profile: None,
2661 visible_namespaces: vec![],
2662 allowed_outbound_namespaces: vec![],
2663 actor_id: None,
2664 exec: Default::default(),
2665 };
2666 KhiveRuntime::new_for_test(config.clone()).expect("create migrated database");
2667 #[cfg(unix)]
2668 khive_storage::test_support::freeze_snapshot_sidecars(&path);
2669
2670 let runtime = KhiveRuntime::new_readonly_for_test(config).expect("explicit read-only boot");
2671 assert!(runtime.is_read_only());
2672 assert_eq!(
2673 runtime.backend().pool().writer_acquisition_snapshot(),
2674 khive_db::pool::WriterAcquisitionSnapshot::default(),
2675 "explicit read-only construction must validate through a reader without ever \
2676 acquiring the writer"
2677 );
2678 }
2679
2680 #[derive(Debug)]
2684 struct ReadOnlyExtraGate {
2685 primary: &'static str,
2686 extra: &'static str,
2687 }
2688
2689 impl khive_gate::Gate for ReadOnlyExtraGate {
2690 fn check(
2691 &self,
2692 req: &khive_gate::GateRequest,
2693 ) -> Result<khive_gate::GateDecision, khive_gate::GateError> {
2694 let allowed = match req.verb.as_str() {
2695 "authorize" => req.namespace.as_str() == self.primary,
2696 "authorize.visible" => req.namespace.as_str() == self.extra,
2697 _ => false,
2698 };
2699 if allowed {
2700 Ok(khive_gate::GateDecision::allow())
2701 } else {
2702 Ok(khive_gate::GateDecision::Deny {
2703 reason: format!(
2704 "{} denied for namespace {:?}",
2705 req.verb,
2706 req.namespace.as_str()
2707 ),
2708 })
2709 }
2710 }
2711 }
2712
2713 #[test]
2714 fn authorize_with_visibility_allows_read_only_extra_namespace() {
2715 let primary = Namespace::parse("lambda:caller").expect("primary");
2716 let extra = Namespace::parse("lambda:read-only").expect("extra");
2717 let config = RuntimeConfig {
2718 db_path: None,
2719 packs: vec!["kg".to_string()],
2720 brain_profile: None,
2721 actor_id: None,
2722 gate: Arc::new(ReadOnlyExtraGate {
2723 primary: "lambda:caller",
2724 extra: "lambda:read-only",
2725 }),
2726 ..RuntimeConfig::no_embeddings()
2727 };
2728 let rt = KhiveRuntime::new(config).expect("memory runtime");
2729
2730 let token = rt
2731 .authorize_with_visibility(primary, vec![extra.clone()])
2732 .expect("Write on primary and Read on extra must mint");
2733 assert!(token.visible_namespaces().contains(&extra));
2734 }
2735
2736 #[derive(Debug)]
2740 struct DenyNamespaceGate {
2741 deny: &'static str,
2742 }
2743
2744 impl khive_gate::Gate for DenyNamespaceGate {
2745 fn check(
2746 &self,
2747 req: &khive_gate::GateRequest,
2748 ) -> Result<khive_gate::GateDecision, khive_gate::GateError> {
2749 if req.namespace.as_str() == self.deny {
2750 Ok(khive_gate::GateDecision::Deny {
2751 reason: "namespace denied by policy".to_string(),
2752 })
2753 } else {
2754 Ok(khive_gate::GateDecision::allow())
2755 }
2756 }
2757 }
2758
2759 #[test]
2760 fn authorize_with_visibility_denies_missing_extra_read() {
2761 let config = RuntimeConfig {
2762 db_path: None,
2763 packs: vec!["kg".to_string()],
2764 brain_profile: None,
2765 actor_id: None,
2766 gate: Arc::new(DenyNamespaceGate {
2767 deny: "lambda:secret",
2768 }),
2769 ..RuntimeConfig::no_embeddings()
2770 };
2771 let rt = KhiveRuntime::new(config).expect("memory runtime");
2772 let primary = Namespace::parse("lambda:caller").expect("primary");
2773 let denied = Namespace::parse("lambda:secret").expect("denied");
2774 let allowed = Namespace::parse("lambda:open").expect("allowed");
2775
2776 rt.authorize_with_visibility(primary.clone(), vec![allowed.clone()])
2779 .expect("mint with only allowed extras");
2780
2781 let err = rt
2782 .authorize_with_visibility(primary, vec![allowed, denied])
2783 .expect_err("a denied extra namespace must refuse the whole mint");
2784 let msg = err.to_string();
2785 assert!(
2786 msg.contains("lambda:secret"),
2787 "refusal must name the offending namespace: {msg}"
2788 );
2789 }
2790
2791 fn make_read_only_runtime() -> (tempfile::TempDir, KhiveRuntime) {
2794 let dir = tempfile::tempdir().unwrap();
2795 let path = dir.path().join("read_only_blob_seam.db");
2796 let config = RuntimeConfig {
2797 web: Default::default(),
2798 telemetry: Default::default(),
2799 mounts: Vec::new(),
2800 brain: Default::default(),
2801 git_write: Default::default(),
2802 display_timezone: chrono_tz::Tz::UTC,
2803 db_path: Some(path.clone()),
2804 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
2805 default_namespace: Namespace::local(),
2806 embedding_model: None,
2807 additional_embedding_models: vec![],
2808 gate: Arc::new(AllowAllGate),
2809 packs: vec!["kg".to_string()],
2810 backend_id: BackendId::main(),
2811 brain_profile: None,
2812 visible_namespaces: vec![],
2813 allowed_outbound_namespaces: vec![],
2814 actor_id: None,
2815 events_split: None,
2816 exec: Default::default(),
2817 };
2818 KhiveRuntime::new_for_test(config.clone()).expect("create migrated database");
2819 #[cfg(unix)]
2820 khive_storage::test_support::freeze_snapshot_sidecars(&path);
2821 let runtime = KhiveRuntime::new_readonly_for_test(config).expect("read-only boot");
2822 assert!(runtime.is_read_only());
2823 (dir, runtime)
2824 }
2825
2826 #[tokio::test]
2827 async fn install_blob_store_on_read_only_runtime_refuses_mutators() {
2828 let (_dir, runtime) = make_read_only_runtime();
2829
2830 use khive_storage::BlobStore as _;
2833 let blob_root = tempfile::tempdir().unwrap();
2834 let writable = Arc::new(
2835 khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
2836 .expect("fs blob store"),
2837 );
2838 let seeded = writable
2839 .put(b"seed".to_vec())
2840 .await
2841 .expect("seed put through the raw store");
2842 runtime
2843 .install_blob_store(writable.clone())
2844 .expect("read-only install wraps rather than refusing");
2845
2846 let installed = runtime.blob_store().expect("installed store");
2847 assert!(
2848 installed.exists(&seeded).await.expect("exists"),
2849 "bounded read surface must stay available"
2850 );
2851 let err = installed
2852 .put(b"post-boot".to_vec())
2853 .await
2854 .expect_err("put must refuse on a read-only runtime");
2855 assert!(
2856 err.to_string().contains("read-only"),
2857 "refusal must name the mode: {err}"
2858 );
2859
2860 runtime
2864 .install_blob_store(writable.clone())
2865 .expect("reinstalling the same raw store must be idempotent");
2866
2867 let bypass_root = tempfile::tempdir().unwrap();
2871 let bypass_store = Arc::new(
2872 khive_db::stores::blob::FsBlobStore::new(bypass_root.path().to_path_buf(), 0)
2873 .expect("fs blob store"),
2874 );
2875 let writable_hydrator = Arc::new(
2876 crate::BlobHydrator::new(bypass_store, crate::DEFAULT_BLOB_HYDRATION_BYTES)
2877 .expect("construct writable hydrator"),
2878 );
2879 let err = runtime
2880 .install_blob_hydrator(writable_hydrator)
2881 .expect_err("a writable hydrator must be refused on a read-only runtime");
2882 assert!(
2883 err.to_string().contains("read-only"),
2884 "refusal must name the mode: {err}"
2885 );
2886 }
2887
2888 #[tokio::test]
2889 async fn mode_aware_hydrator_constructor_satisfies_read_only_install() {
2890 let (_dir, runtime) = make_read_only_runtime();
2896 let blob_root = tempfile::tempdir().unwrap();
2897 let raw = Arc::new(
2898 khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
2899 .expect("fs blob store"),
2900 ) as Arc<dyn khive_storage::BlobStore>;
2901 let seeded = raw.put(b"seed".to_vec()).await.expect("seed put");
2902 let hydrator = Arc::new(
2903 crate::BlobHydrator::for_mode(
2904 Arc::clone(&raw),
2905 crate::DEFAULT_BLOB_HYDRATION_BYTES,
2906 true,
2907 )
2908 .expect("mode-aware read-only construction"),
2909 );
2910 runtime
2911 .install_blob_hydrator(hydrator)
2912 .expect("a for_mode(read_only) hydrator must pass the read-only gate");
2913 let installed = runtime.blob_store().expect("installed store");
2914 assert!(installed.exists(&seeded).await.expect("exists"));
2915 installed
2916 .put(b"post".to_vec())
2917 .await
2918 .expect_err("mutation must refuse through the wrapped store");
2919
2920 let twin = Arc::new(
2924 crate::BlobHydrator::for_mode(
2925 Arc::clone(&raw),
2926 crate::DEFAULT_BLOB_HYDRATION_BYTES,
2927 true,
2928 )
2929 .expect("twin construction"),
2930 );
2931 runtime
2932 .install_blob_hydrator(twin)
2933 .expect("an equivalent pairing must read as idempotent");
2934 }
2935
2936 #[tokio::test]
2937 async fn shared_install_permits_governed_writable_hydrator_on_read_only_handle() {
2938 let (_dir, runtime) = make_read_only_runtime();
2947 let blob_root = tempfile::tempdir().unwrap();
2948 let raw = Arc::new(
2949 khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
2950 .expect("fs blob store"),
2951 ) as Arc<dyn khive_storage::BlobStore>;
2952 let hand_paired = Arc::new(
2953 crate::BlobHydrator::for_mode(raw, crate::DEFAULT_BLOB_HYDRATION_BYTES, false)
2954 .expect("writable construction"),
2955 );
2956 runtime
2957 .install_blob_hydrator(Arc::clone(&hand_paired))
2958 .expect_err("the plain seam must refuse a writable hydrator on a read-only handle");
2959 runtime
2960 .install_shared_blob_hydrator(hand_paired)
2961 .expect_err("the shared seam must refuse a hand-paired (ungoverned) hydrator");
2962
2963 let governing = khive_db::StorageBackend::memory().expect("memory backend");
2967 let cfg = crate::KhiveConfig {
2968 storage: crate::engine_config::StorageSectionConfig {
2969 blob: Some(crate::engine_config::BlobConfig::Fs {
2970 root: Some(blob_root.path().to_string_lossy().into_owned()),
2971 floor_bytes: Some(0),
2972 }),
2973 },
2974 ..crate::KhiveConfig::default()
2975 };
2976 let governed = Arc::new(
2977 crate::BlobHydrator::resolve_for_governing_backend(
2978 &cfg,
2979 &governing,
2980 &governing,
2981 crate::DEFAULT_BLOB_HYDRATION_BYTES,
2982 )
2983 .expect("governed construction"),
2984 );
2985 runtime
2986 .install_shared_blob_hydrator(governed)
2987 .expect("the shared seam accepts a governed hydrator");
2988 let installed = runtime.blob_store().expect("installed store");
2989 let put = installed
2990 .put(b"shared-write".to_vec())
2991 .await
2992 .expect("the governed writable capability must actually mutate");
2993 assert!(installed.exists(&put).await.expect("exists"));
2994 }
2995
2996 #[test]
3006 #[serial]
3007 fn tilde_prefixed_db_override_resolves_and_boots_like_the_absolute_equivalent() {
3008 if crate::test_process::run_in_child() {
3009 return;
3010 }
3011
3012 let original_home = std::env::var_os("HOME");
3013 let original_cwd = std::env::current_dir().expect("read cwd");
3014 let home_dir = tempfile::tempdir().expect("home tempdir");
3015 let work_dir = tempfile::tempdir().expect("work tempdir");
3016 std::env::set_var("HOME", home_dir.path());
3017 std::env::set_current_dir(work_dir.path()).expect("chdir into isolated work dir");
3018
3019 let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3020 let tilde_anchor = crate::config::resolve_db_anchor(Some("~/data.db"))
3021 .expect("an explicit path always anchors");
3022 let expected = home_dir.path().join("data.db");
3023 assert_eq!(
3024 tilde_anchor, expected,
3025 "resolve_db_anchor must expand a leading ~ to $HOME before it ever \
3026 reaches RuntimeConfig.db_path"
3027 );
3028
3029 let absolute_anchor = crate::config::resolve_db_anchor(Some(
3030 expected.to_str().expect("utf8 tempdir path"),
3031 ))
3032 .expect("an explicit path always anchors");
3033 assert_eq!(
3034 tilde_anchor, absolute_anchor,
3035 "a ~-prefixed override and its equivalent absolute path must resolve to \
3036 the identical anchor"
3037 );
3038
3039 let make_config = |db_path: std::path::PathBuf| RuntimeConfig {
3040 web: Default::default(),
3041 telemetry: Default::default(),
3042 mounts: Vec::new(),
3043 brain: Default::default(),
3044 git_write: Default::default(),
3045 display_timezone: chrono_tz::Tz::UTC,
3046 events_split: None,
3047 db_path: Some(db_path),
3048 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3049 default_namespace: Namespace::local(),
3050 embedding_model: None,
3051 additional_embedding_models: vec![],
3052 gate: Arc::new(AllowAllGate),
3053 packs: vec!["kg".to_string()],
3054 backend_id: BackendId::main(),
3055 brain_profile: None,
3056 visible_namespaces: vec![],
3057 allowed_outbound_namespaces: vec![],
3058 actor_id: None,
3059 exec: Default::default(),
3060 };
3061
3062 let tilde_cfg = make_config(tilde_anchor.clone());
3063
3064 let rt =
3065 KhiveRuntime::new_for_test(tilde_cfg).expect("boot must open the expanded path");
3066 assert_eq!(
3067 rt.backend_data_dir().expect("file backend"),
3068 home_dir.path(),
3069 "single-backend boot must open the file under the expanded $HOME \
3070 directory, not a literal ~ path relative to cwd"
3071 );
3072 assert!(
3073 expected.exists(),
3074 "the database file must be created at the expanded $HOME path"
3075 );
3076 assert!(
3077 !work_dir.path().join("~").exists(),
3078 "boot must never create a literal '~' directory under the process cwd"
3079 );
3080 }));
3081
3082 match &original_home {
3083 Some(h) => std::env::set_var("HOME", h),
3084 None => std::env::remove_var("HOME"),
3085 }
3086 let _ = std::env::set_current_dir(&original_cwd);
3087 outcome.expect("test body panicked");
3088 }
3089
3090 #[test]
3091 fn from_backend_uses_provided_backend() {
3092 let backend = Arc::new(StorageBackend::memory().expect("memory backend"));
3093 let config = RuntimeConfig {
3094 web: Default::default(),
3095 telemetry: Default::default(),
3096 mounts: Vec::new(),
3097 brain: Default::default(),
3098 git_write: Default::default(),
3099 display_timezone: chrono_tz::Tz::UTC,
3100 events_split: None,
3101 db_path: None,
3102 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3103 default_namespace: Namespace::local(),
3104 embedding_model: None,
3105 additional_embedding_models: vec![],
3106 gate: Arc::new(AllowAllGate),
3107 packs: vec!["kg".to_string()],
3108 backend_id: BackendId::parse("lore").expect("valid backend id"),
3109 brain_profile: None,
3110 visible_namespaces: vec![],
3111 allowed_outbound_namespaces: vec![],
3112 actor_id: None,
3113 exec: Default::default(),
3114 };
3115 let rt = KhiveRuntime::from_backend(backend, config);
3116 assert_eq!(rt.backend_id().as_str(), "lore");
3117 assert!(rt.config().db_path.is_none());
3118 }
3119
3120 #[test]
3121 fn backend_id_defaults_to_main() {
3122 let rt = KhiveRuntime::memory().unwrap();
3123 assert_eq!(rt.backend_id().as_str(), BackendId::MAIN);
3124 }
3125
3126 #[test]
3127 fn store_accessors_return_ok() {
3128 let rt = KhiveRuntime::memory().unwrap();
3129 let tok = NamespaceToken::local();
3130 assert!(rt.entities(&tok).is_ok());
3131 assert!(rt.graph(&tok).is_ok());
3132 assert!(rt.notes(&tok).is_ok());
3133 assert!(rt.events(&tok).is_ok());
3134 }
3135
3136 fn attributed_event_runtime() -> (KhiveRuntime, NamespaceToken) {
3137 let runtime = KhiveRuntime::new(RuntimeConfig {
3138 db_path: None,
3139 actor_id: Some("lambda:enrolled".to_string()),
3140 ..RuntimeConfig::no_embeddings()
3141 })
3142 .expect("memory runtime");
3143 let token = runtime
3144 .authorize(Namespace::local())
3145 .expect("configured actor is allowed by the default gate");
3146 (runtime, token)
3147 }
3148
3149 fn forged_event(verb: &str) -> Event {
3150 Event::new(
3151 "caller-selected-namespace",
3152 verb,
3153 EventKind::Audit,
3154 SubstrateKind::Event,
3155 "caller-selected-actor",
3156 )
3157 }
3158
3159 fn assert_token_attribution(event: &Event) {
3160 assert_eq!(event.namespace, "local");
3161 assert_eq!(event.actor, "actor:lambda:enrolled");
3162 }
3163
3164 #[tokio::test]
3165 async fn token_scoped_event_store_stamps_resolved_attribution_on_single_append() {
3166 let (runtime, token) = attributed_event_runtime();
3167 let store = runtime.events(&token).expect("event store");
3168 let event = forged_event("single");
3169 let id = event.id;
3170
3171 store.append_event(event).await.expect("append");
3172
3173 let stored = store
3174 .get_event(id)
3175 .await
3176 .expect("read")
3177 .expect("the token-stamped event remains visible to the token");
3178 assert_token_attribution(&stored);
3179 assert_eq!(stored.verb, "single", "non-attribution fields survive");
3180 }
3181
3182 #[tokio::test]
3183 async fn token_scoped_event_store_stamps_resolved_attribution_on_batch_paths() {
3184 let (runtime, token) = attributed_event_runtime();
3185 let store = runtime.events(&token).expect("event store");
3186 let ordinary = forged_event("batch");
3187 let ordinary_id = ordinary.id;
3188
3189 store
3190 .append_events(vec![ordinary])
3191 .await
3192 .expect("ordinary batch append");
3193 let stored = store
3194 .get_event(ordinary_id)
3195 .await
3196 .expect("read")
3197 .expect("ordinary batch event remains visible to the token");
3198 assert_token_attribution(&stored);
3199
3200 let idempotent = forged_event("idempotent_batch");
3201 let idempotent_id = idempotent.id;
3202 let outcome = store
3203 .append_events_idempotent(vec![idempotent])
3204 .await
3205 .expect("idempotent batch append");
3206 assert_eq!(
3207 outcome.rows,
3208 vec![khive_storage::event::EventAppendDisposition::Inserted]
3209 );
3210 let stored = store
3211 .get_event(idempotent_id)
3212 .await
3213 .expect("read")
3214 .expect("idempotent batch event remains visible to the token");
3215 assert_token_attribution(&stored);
3216 }
3217
3218 #[test]
3219 fn vectors_returns_unconfigured_without_model() {
3220 let rt = KhiveRuntime::memory().unwrap();
3221 let tok = NamespaceToken::local();
3222 match rt.vectors(&tok) {
3223 Err(crate::RuntimeError::Unconfigured(s)) => assert_eq!(s, "embedding_model"),
3224 Err(other) => panic!("expected Unconfigured, got {:?}", other),
3225 Ok(_) => panic!("expected Err, got Ok"),
3226 }
3227 }
3228
3229 #[test]
3230 fn vec_model_key_sanitizes_dots_and_dashes() {
3231 assert_eq!(
3232 vec_model_key(EmbeddingModel::BgeSmallEnV15),
3233 "bge_small_en_v1_5"
3234 );
3235 assert_eq!(
3236 vec_model_key(EmbeddingModel::BgeBaseEnV15),
3237 "bge_base_en_v1_5"
3238 );
3239 assert_eq!(
3240 vec_model_key(EmbeddingModel::AllMiniLmL6V2),
3241 "all_minilm_l6_v2"
3242 );
3243 }
3244
3245 #[test]
3246 fn default_config_uses_allow_all_gate() {
3247 let cfg = RuntimeConfig::default();
3248 assert_eq!(cfg.default_namespace.as_str(), "local");
3249 let _: GateRef = cfg.gate.clone();
3250 }
3251
3252 #[test]
3253 fn parse_pack_list_handles_comma_and_whitespace() {
3254 assert_eq!(parse_pack_list("kg"), vec!["kg".to_string()]);
3255 assert_eq!(
3256 parse_pack_list("kg,gtd"),
3257 vec!["kg".to_string(), "gtd".to_string()]
3258 );
3259 assert_eq!(
3260 parse_pack_list(" kg , gtd "),
3261 vec!["kg".to_string(), "gtd".to_string()]
3262 );
3263 assert_eq!(
3264 parse_pack_list("kg gtd"),
3265 vec!["kg".to_string(), "gtd".to_string()]
3266 );
3267 assert_eq!(parse_pack_list(",,"), Vec::<String>::new());
3268 assert_eq!(parse_pack_list(""), Vec::<String>::new());
3269 }
3270
3271 #[test]
3272 fn default_config_packs_loads_production_set() {
3273 let prior = std::env::var("KHIVE_PACKS").ok();
3274 unsafe {
3276 std::env::remove_var("KHIVE_PACKS");
3277 }
3278 let cfg = RuntimeConfig::default();
3280 assert_eq!(cfg.packs, RuntimeConfig::built_in_packs());
3281 assert!(cfg.packs.contains(&"kg".to_string()));
3282 assert!(cfg.packs.contains(&"gtd".to_string()));
3283 assert!(cfg.packs.contains(&"memory".to_string()));
3284 assert!(cfg.packs.contains(&"brain".to_string()));
3285 assert!(cfg.packs.contains(&"comm".to_string()));
3286 assert!(cfg.packs.contains(&"schedule".to_string()));
3287 assert!(cfg.packs.contains(&"knowledge".to_string()));
3288 assert!(cfg.packs.contains(&"session".to_string()));
3291 assert!(cfg.packs.contains(&"git".to_string()));
3292 assert!(cfg.packs.contains(&"code".to_string()));
3293 assert!(cfg.packs.contains(&"workspace".to_string()));
3294 assert!(cfg.packs.contains(&"blob".to_string()));
3299 assert!(cfg.packs.contains(&"tool".to_string()));
3302 assert!(cfg.packs.contains(&"exec".to_string()));
3303 assert_eq!(cfg.packs.len(), 14);
3304 if let Some(v) = prior {
3305 unsafe {
3307 std::env::set_var("KHIVE_PACKS", v);
3308 }
3309 }
3310 }
3311
3312 #[test]
3313 fn default_config_uses_minilm_when_env_unset() {
3314 let prior = std::env::var("KHIVE_EMBEDDING_MODEL").ok();
3315 unsafe {
3318 std::env::remove_var("KHIVE_EMBEDDING_MODEL");
3319 }
3320 let cfg = RuntimeConfig::default();
3321 assert_eq!(cfg.embedding_model, Some(EmbeddingModel::AllMiniLmL6V2));
3322 if let Some(v) = prior {
3323 unsafe {
3325 std::env::set_var("KHIVE_EMBEDDING_MODEL", v);
3326 }
3327 }
3328 }
3329
3330 use crate::engine_config::{ActorConfig, KhiveConfig};
3333
3334 fn khive_cfg_with_actor(id: &str) -> KhiveConfig {
3335 KhiveConfig {
3336 engines: vec![],
3337 actor: ActorConfig {
3338 id: Some(id.to_string()),
3339 display_name: None,
3340 ..Default::default()
3341 },
3342 ..KhiveConfig::default()
3343 }
3344 }
3345
3346 #[test]
3347 fn runtime_config_from_khive_config_actor_id_does_not_override_default_namespace() {
3348 let base = RuntimeConfig {
3353 web: Default::default(),
3354 telemetry: Default::default(),
3355 mounts: Vec::new(),
3356 brain: Default::default(),
3357 git_write: Default::default(),
3358 display_timezone: chrono_tz::Tz::UTC,
3359 events_split: None,
3360 db_path: None,
3361 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3362 default_namespace: Namespace::local(),
3363 embedding_model: None,
3364 additional_embedding_models: vec![],
3365 gate: Arc::new(AllowAllGate),
3366 packs: vec!["kg".to_string()],
3367 backend_id: BackendId::main(),
3368 brain_profile: None,
3369 visible_namespaces: vec![],
3370 allowed_outbound_namespaces: vec![],
3371 actor_id: None,
3372 exec: Default::default(),
3373 };
3374 let cfg = khive_cfg_with_actor("lambda:khive");
3375 let result = runtime_config_from_khive_config(&cfg, base);
3376 assert_eq!(
3377 result.default_namespace.as_str(),
3378 "local",
3379 "actor.id must not become default_namespace (ADR-007 Rev 4 Rule 0); writes pin to local"
3380 );
3381 }
3382
3383 #[test]
3384 fn runtime_config_from_khive_config_empty_actor_id_keeps_base_namespace() {
3385 let base = RuntimeConfig {
3386 web: Default::default(),
3387 telemetry: Default::default(),
3388 mounts: Vec::new(),
3389 brain: Default::default(),
3390 git_write: Default::default(),
3391 display_timezone: chrono_tz::Tz::UTC,
3392 events_split: None,
3393 db_path: None,
3394 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3395 default_namespace: Namespace::parse("lambda:base").unwrap(),
3396 embedding_model: None,
3397 additional_embedding_models: vec![],
3398 gate: Arc::new(AllowAllGate),
3399 packs: vec!["kg".to_string()],
3400 backend_id: BackendId::main(),
3401 brain_profile: None,
3402 visible_namespaces: vec![],
3403 allowed_outbound_namespaces: vec![],
3404 actor_id: None,
3405 exec: Default::default(),
3406 };
3407 let cfg = KhiveConfig {
3408 engines: vec![],
3409 actor: ActorConfig {
3410 id: Some(String::new()),
3411 display_name: None,
3412 ..Default::default()
3413 },
3414 ..KhiveConfig::default()
3415 };
3416 let result = runtime_config_from_khive_config(&cfg, base);
3417 assert_eq!(
3418 result.default_namespace.as_str(),
3419 "lambda:base",
3420 "empty actor.id must not override base namespace"
3421 );
3422 }
3423
3424 #[test]
3425 fn runtime_config_from_khive_config_absent_actor_id_keeps_base_namespace() {
3426 let base = RuntimeConfig {
3427 web: Default::default(),
3428 telemetry: Default::default(),
3429 mounts: Vec::new(),
3430 brain: Default::default(),
3431 git_write: Default::default(),
3432 display_timezone: chrono_tz::Tz::UTC,
3433 events_split: None,
3434 db_path: None,
3435 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3436 default_namespace: Namespace::parse("lambda:base").unwrap(),
3437 embedding_model: None,
3438 additional_embedding_models: vec![],
3439 gate: Arc::new(AllowAllGate),
3440 packs: vec!["kg".to_string()],
3441 backend_id: BackendId::main(),
3442 brain_profile: None,
3443 visible_namespaces: vec![],
3444 allowed_outbound_namespaces: vec![],
3445 actor_id: None,
3446 exec: Default::default(),
3447 };
3448 let cfg = KhiveConfig::default(); let result = runtime_config_from_khive_config(&cfg, base);
3450 assert_eq!(
3451 result.default_namespace.as_str(),
3452 "lambda:base",
3453 "absent actor.id must not override base namespace"
3454 );
3455 }
3456
3457 #[test]
3458 fn runtime_config_from_khive_config_actor_id_with_engines() {
3459 let base = RuntimeConfig {
3460 web: Default::default(),
3461 telemetry: Default::default(),
3462 mounts: Vec::new(),
3463 brain: Default::default(),
3464 git_write: Default::default(),
3465 display_timezone: chrono_tz::Tz::UTC,
3466 events_split: None,
3467 db_path: None,
3468 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3469 default_namespace: Namespace::local(),
3470 embedding_model: None,
3471 additional_embedding_models: vec![],
3472 gate: Arc::new(AllowAllGate),
3473 packs: vec!["kg".to_string()],
3474 backend_id: BackendId::main(),
3475 brain_profile: None,
3476 visible_namespaces: vec![],
3477 allowed_outbound_namespaces: vec![],
3478 actor_id: None,
3479 exec: Default::default(),
3480 };
3481 let cfg = KhiveConfig {
3482 engines: vec![crate::engine_config::EngineConfig {
3483 name: "default".to_string(),
3484 model: "all-minilm-l6-v2".to_string(),
3485 default: true,
3486 fusion_weight: None,
3487 dims: None,
3488 }],
3489 actor: ActorConfig {
3490 id: Some("lambda:test".to_string()),
3491 display_name: None,
3492 ..Default::default()
3493 },
3494 ..KhiveConfig::default()
3495 };
3496 let result = runtime_config_from_khive_config(&cfg, base);
3497 assert_eq!(
3498 result.default_namespace.as_str(),
3499 "local",
3500 "actor.id must not override default_namespace (ADR-007 Rev 4 Rule 0); \
3501 writes pin to local; engine config is still applied"
3502 );
3503 assert!(result.embedding_model.is_some());
3504 }
3505
3506 #[test]
3509 fn runtime_config_from_khive_config_display_timezone_overrides_base() {
3510 let base = RuntimeConfig {
3511 web: Default::default(),
3512 telemetry: Default::default(),
3513 mounts: Vec::new(),
3514 brain: Default::default(),
3515 git_write: Default::default(),
3516 display_timezone: chrono_tz::Tz::UTC,
3517 events_split: None,
3518 db_path: None,
3519 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3520 default_namespace: Namespace::local(),
3521 embedding_model: None,
3522 additional_embedding_models: vec![],
3523 gate: Arc::new(AllowAllGate),
3524 packs: vec!["kg".to_string()],
3525 backend_id: BackendId::main(),
3526 brain_profile: None,
3527 visible_namespaces: vec![],
3528 allowed_outbound_namespaces: vec![],
3529 actor_id: None,
3530 exec: Default::default(),
3531 };
3532 let cfg = KhiveConfig {
3533 display: crate::engine_config::DisplaySectionConfig {
3534 timezone: Some("America/New_York".to_string()),
3535 },
3536 ..KhiveConfig::default()
3537 };
3538 let result = runtime_config_from_khive_config(&cfg, base);
3539 assert_eq!(
3540 result.display_timezone,
3541 "America/New_York".parse::<chrono_tz::Tz>().unwrap(),
3542 "[display] timezone in khive.toml must override base.display_timezone"
3543 );
3544 }
3545
3546 #[test]
3547 fn runtime_config_from_khive_config_absent_display_timezone_keeps_base() {
3548 let base = RuntimeConfig {
3549 web: Default::default(),
3550 telemetry: Default::default(),
3551 mounts: Vec::new(),
3552 brain: Default::default(),
3553 git_write: Default::default(),
3554 display_timezone: "Asia/Tokyo".parse().unwrap(),
3555 events_split: None,
3556 db_path: None,
3557 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3558 default_namespace: Namespace::local(),
3559 embedding_model: None,
3560 additional_embedding_models: vec![],
3561 gate: Arc::new(AllowAllGate),
3562 packs: vec!["kg".to_string()],
3563 backend_id: BackendId::main(),
3564 brain_profile: None,
3565 visible_namespaces: vec![],
3566 allowed_outbound_namespaces: vec![],
3567 actor_id: None,
3568 exec: Default::default(),
3569 };
3570 let cfg = KhiveConfig::default(); let result = runtime_config_from_khive_config(&cfg, base);
3572 assert_eq!(
3573 result.display_timezone,
3574 "Asia/Tokyo".parse::<chrono_tz::Tz>().unwrap(),
3575 "absent [display] timezone must preserve base.display_timezone unchanged"
3576 );
3577 }
3578
3579 #[test]
3588 #[serial]
3589 fn runtime_config_from_khive_config_engines_present_preserves_env_actor_when_toml_has_none() {
3590 let prior = std::env::var("KHIVE_ACTOR").ok();
3591 unsafe {
3593 std::env::set_var("KHIVE_ACTOR", "lambda:test-env-actor");
3594 }
3595 let base = RuntimeConfig::default();
3596 assert_eq!(base.actor_id.as_deref(), Some("lambda:test-env-actor"));
3597
3598 let cfg = KhiveConfig {
3599 engines: vec![crate::engine_config::EngineConfig {
3600 name: "default".to_string(),
3601 model: "all-minilm-l6-v2".to_string(),
3602 default: true,
3603 fusion_weight: None,
3604 dims: None,
3605 }],
3606 actor: ActorConfig::default(), ..KhiveConfig::default()
3608 };
3609 let result = runtime_config_from_khive_config(&cfg, base);
3610 assert_eq!(
3611 result.actor_id.as_deref(),
3612 Some("lambda:test-env-actor"),
3613 "engines-present arm must preserve base.actor_id (env actor) when TOML has no [actor] id"
3614 );
3615
3616 unsafe {
3618 match prior {
3619 Some(v) => std::env::set_var("KHIVE_ACTOR", v),
3620 None => std::env::remove_var("KHIVE_ACTOR"),
3621 }
3622 }
3623 }
3624
3625 #[test]
3626 #[serial]
3627 fn runtime_config_from_khive_config_engines_empty_preserves_env_actor_when_toml_has_none() {
3628 let prior = std::env::var("KHIVE_ACTOR").ok();
3629 unsafe {
3631 std::env::set_var("KHIVE_ACTOR", "lambda:test-env-actor");
3632 }
3633 let base = RuntimeConfig::default();
3634 assert_eq!(base.actor_id.as_deref(), Some("lambda:test-env-actor"));
3635
3636 let cfg = KhiveConfig {
3637 engines: vec![],
3638 actor: ActorConfig::default(), ..KhiveConfig::default()
3640 };
3641 let result = runtime_config_from_khive_config(&cfg, base);
3642 assert_eq!(
3643 result.actor_id.as_deref(),
3644 Some("lambda:test-env-actor"),
3645 "engines-empty early-return arm must preserve base.actor_id (env actor) when TOML has no [actor] id"
3646 );
3647
3648 unsafe {
3650 match prior {
3651 Some(v) => std::env::set_var("KHIVE_ACTOR", v),
3652 None => std::env::remove_var("KHIVE_ACTOR"),
3653 }
3654 }
3655 }
3656
3657 #[test]
3658 #[serial]
3659 fn runtime_config_from_khive_config_toml_actor_wins_over_env_actor() {
3660 let prior = std::env::var("KHIVE_ACTOR").ok();
3661 unsafe {
3663 std::env::set_var("KHIVE_ACTOR", "lambda:test-env-actor");
3664 }
3665 let base = RuntimeConfig::default();
3666 assert_eq!(base.actor_id.as_deref(), Some("lambda:test-env-actor"));
3667
3668 let cfg = khive_cfg_with_actor("lambda:toml-actor");
3669 let result = runtime_config_from_khive_config(&cfg, base);
3670 assert_eq!(
3671 result.actor_id.as_deref(),
3672 Some("lambda:toml-actor"),
3673 "TOML [actor] id must win over the env-resolved base.actor_id"
3674 );
3675
3676 unsafe {
3678 match prior {
3679 Some(v) => std::env::set_var("KHIVE_ACTOR", v),
3680 None => std::env::remove_var("KHIVE_ACTOR"),
3681 }
3682 }
3683 }
3684
3685 fn migrated_memory_backend() -> Arc<StorageBackend> {
3691 let backend = StorageBackend::memory().expect("memory backend");
3692 {
3693 let mut writer = backend.pool().try_writer().expect("writer");
3694 khive_db::run_migrations(writer.conn_mut()).expect("migrations");
3695 }
3696 Arc::new(backend)
3697 }
3698
3699 fn secondary_config() -> RuntimeConfig {
3700 RuntimeConfig {
3701 web: Default::default(),
3702 telemetry: Default::default(),
3703 mounts: Vec::new(),
3704 brain: Default::default(),
3705 git_write: Default::default(),
3706 exec: Default::default(),
3707 display_timezone: chrono_tz::Tz::UTC,
3708 events_split: None,
3709 db_path: None,
3710 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3711 default_namespace: Namespace::local(),
3712 embedding_model: None,
3713 additional_embedding_models: vec![],
3714 gate: Arc::new(AllowAllGate),
3715 packs: vec!["kg".to_string()],
3716 backend_id: BackendId::parse("lore").expect("valid backend id"),
3717 brain_profile: None,
3718 visible_namespaces: vec![],
3719 allowed_outbound_namespaces: vec![],
3720 actor_id: None,
3721 }
3722 }
3723
3724 #[test]
3725 fn core_on_main_runtime_returns_same_backend_id() {
3726 let rt = KhiveRuntime::memory().unwrap();
3728 assert_eq!(rt.backend_id().as_str(), BackendId::MAIN);
3729 let core_rt = rt.core();
3730 assert_eq!(core_rt.backend_id().as_str(), BackendId::MAIN);
3731 }
3732
3733 #[tokio::test]
3734 async fn core_on_main_runtime_round_trips_note() {
3735 let rt = KhiveRuntime::memory().unwrap();
3738 let tok = NamespaceToken::local();
3739
3740 let note = rt
3741 .core()
3742 .create_note(
3743 &tok,
3744 "observation",
3745 None,
3746 "adr073-main-round-trip",
3747 None,
3748 None,
3749 vec![],
3750 )
3751 .await
3752 .expect("create_note via core()");
3753
3754 let found = rt
3755 .notes(&tok)
3756 .expect("notes store")
3757 .get_note(note.id)
3758 .await
3759 .expect("get_note");
3760
3761 assert!(
3762 found.is_some(),
3763 "note written via core() must be visible through original rt"
3764 );
3765 }
3766
3767 #[tokio::test]
3779 async fn cross_backend_split_note_to_main_aux_to_secondary() {
3780 use khive_storage::{SqlStatement, SqlValue};
3781
3782 let main_arc = migrated_memory_backend();
3783 let secondary_arc = migrated_memory_backend();
3784
3785 let main_config = RuntimeConfig {
3786 web: Default::default(),
3787 telemetry: Default::default(),
3788 mounts: Vec::new(),
3789 brain: Default::default(),
3790 git_write: Default::default(),
3791 display_timezone: chrono_tz::Tz::UTC,
3792 events_split: None,
3793 db_path: None,
3794 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3795 default_namespace: Namespace::local(),
3796 embedding_model: None,
3797 additional_embedding_models: vec![],
3798 gate: Arc::new(AllowAllGate),
3799 packs: vec!["kg".to_string()],
3800 backend_id: BackendId::main(),
3801 brain_profile: None,
3802 visible_namespaces: vec![],
3803 allowed_outbound_namespaces: vec![],
3804 actor_id: None,
3805 exec: Default::default(),
3806 };
3807
3808 let rt_main = KhiveRuntime::from_backend(main_arc.clone(), main_config);
3809 let rt_secondary = KhiveRuntime::from_backend(secondary_arc, secondary_config())
3810 .with_core_backend(main_arc.clone());
3811
3812 let tok = NamespaceToken::local();
3813
3814 let note = rt_secondary
3817 .core()
3818 .create_note(
3819 &tok,
3820 "observation",
3821 None,
3822 "adr073-split-test",
3823 None,
3824 None,
3825 vec![],
3826 )
3827 .await
3828 .expect("create_note via core()");
3829 let note_id = note.id;
3830
3831 let in_main = rt_main
3833 .notes(&tok)
3834 .expect("main notes store")
3835 .get_note(note_id)
3836 .await
3837 .expect("get_note from main");
3838 assert!(
3839 in_main.is_some(),
3840 "note written via core() must appear in main backend A"
3841 );
3842
3843 let in_secondary = rt_secondary
3845 .notes(&tok)
3846 .expect("secondary notes store")
3847 .get_note(note_id)
3848 .await
3849 .expect("get_note from secondary");
3850 assert!(
3851 in_secondary.is_none(),
3852 "note written to main via core() must NOT appear in secondary backend B"
3853 );
3854
3855 {
3858 let mut writer = rt_secondary.sql().writer().await.expect("secondary writer");
3859 writer
3860 .execute(SqlStatement {
3861 sql: "CREATE TABLE IF NOT EXISTS _test_adr073_aux \
3862 (marker TEXT PRIMARY KEY)"
3863 .into(),
3864 params: vec![],
3865 label: None,
3866 })
3867 .await
3868 .expect("create aux table in B");
3869 writer
3870 .execute(SqlStatement {
3871 sql: "INSERT INTO _test_adr073_aux VALUES (?1)".into(),
3872 params: vec![SqlValue::Text("b-side-sentinel".into())],
3873 label: None,
3874 })
3875 .await
3876 .expect("insert into aux table in B");
3877 }
3878
3879 let mut reader_b = rt_secondary.sql().reader().await.expect("secondary reader");
3881 let rows_b = reader_b
3882 .query_all(SqlStatement {
3883 sql: "SELECT marker FROM _test_adr073_aux".into(),
3884 params: vec![],
3885 label: None,
3886 })
3887 .await
3888 .expect("select from B");
3889 assert_eq!(rows_b.len(), 1, "aux row must exist in B");
3890 match rows_b[0].get("marker") {
3891 Some(SqlValue::Text(s)) => {
3892 assert_eq!(s, "b-side-sentinel", "sentinel value must match")
3893 }
3894 other => panic!("expected Text('b-side-sentinel'), got {other:?}"),
3895 }
3896
3897 let mut reader_a = rt_main.sql().reader().await.expect("main reader");
3899 let result_a = reader_a
3900 .query_all(SqlStatement {
3901 sql: "SELECT marker FROM _test_adr073_aux".into(),
3902 params: vec![],
3903 label: None,
3904 })
3905 .await;
3906 match result_a {
3908 Err(e) => assert!(
3909 e.to_string().contains("no such table"),
3910 "expected 'no such table' error from A, got: {e}"
3911 ),
3912 Ok(rows) => assert!(
3913 rows.is_empty(),
3914 "aux table must not have rows in A, got {} rows",
3915 rows.len()
3916 ),
3917 }
3918 }
3919
3920 #[test]
3921 fn constructors_leave_core_backend_none_by_behavior() {
3922 let rt_mem = KhiveRuntime::memory().unwrap();
3925 assert_eq!(rt_mem.core().backend_id().as_str(), BackendId::MAIN);
3926
3927 let backend = migrated_memory_backend();
3928 let rt_from = KhiveRuntime::from_backend(
3929 backend,
3930 RuntimeConfig {
3931 web: Default::default(),
3932 telemetry: Default::default(),
3933 mounts: Vec::new(),
3934 brain: Default::default(),
3935 git_write: Default::default(),
3936 display_timezone: chrono_tz::Tz::UTC,
3937 events_split: None,
3938 db_path: None,
3939 blob_hydration_bytes: crate::DEFAULT_BLOB_HYDRATION_BYTES,
3940 default_namespace: Namespace::local(),
3941 embedding_model: None,
3942 additional_embedding_models: vec![],
3943 gate: Arc::new(AllowAllGate),
3944 packs: vec!["kg".to_string()],
3945 backend_id: BackendId::parse("lore").expect("valid backend id"),
3946 brain_profile: None,
3947 visible_namespaces: vec![],
3948 allowed_outbound_namespaces: vec![],
3949 actor_id: None,
3950 exec: Default::default(),
3951 },
3952 );
3953 assert_eq!(rt_from.core().backend_id().as_str(), "lore");
3956 }
3957
3958 #[test]
3959 fn with_core_backend_sets_core_then_core_returns_main_id() {
3960 let main_arc = migrated_memory_backend();
3962 let secondary_arc = migrated_memory_backend();
3963
3964 let rt_secondary = KhiveRuntime::from_backend(secondary_arc, secondary_config())
3965 .with_core_backend(main_arc);
3966
3967 assert_eq!(rt_secondary.backend_id().as_str(), "lore");
3968 assert_eq!(
3969 rt_secondary.core().backend_id().as_str(),
3970 BackendId::MAIN,
3971 "core() on a secondary runtime must return a main-bound handle"
3972 );
3973 }
3974
3975 #[test]
3976 fn attachment_store_rejects_secondary_handle_and_accepts_its_core_projection() {
3977 let main_arc = migrated_memory_backend();
3978 let secondary_arc = migrated_memory_backend();
3979 let rt_secondary = KhiveRuntime::from_backend(secondary_arc, secondary_config())
3980 .with_core_backend(main_arc);
3981
3982 let error = match rt_secondary.attachments() {
3983 Ok(_) => panic!("a secondary runtime must not expose attachment mutation"),
3984 Err(error) => error,
3985 };
3986 assert!(matches!(error, RuntimeError::InvalidInput(_)));
3987 assert!(
3988 error.to_string().contains("canonical main backend"),
3989 "secondary refusal must explain the liveness authority: {error}"
3990 );
3991 rt_secondary
3992 .core()
3993 .attachments()
3994 .expect("core projection must expose the main attachment store");
3995 }
3996
3997 #[tokio::test]
3998 async fn record_plus_attachment_publication_rejects_a_secondary_runtime() {
3999 use khive_storage::{BlobStore as _, NewAttachment};
4000
4001 let main_arc = migrated_memory_backend();
4002 let secondary_arc = migrated_memory_backend();
4003 let rt_secondary =
4004 KhiveRuntime::from_backend(Arc::clone(&secondary_arc), secondary_config())
4005 .with_core_backend(Arc::clone(&main_arc));
4006 let blob_root = tempfile::tempdir().expect("blob root");
4007 let blob_store = Arc::new(
4008 khive_db::stores::blob::FsBlobStore::new(blob_root.path().to_path_buf(), 0)
4009 .expect("blob store"),
4010 );
4011 let content_ref = blob_store.put(b"secondary-ref".to_vec()).await.unwrap();
4012 rt_secondary
4013 .install_blob_store(blob_store.clone())
4014 .expect("shared blob store");
4015 let token = rt_secondary.authorize(Namespace::local()).unwrap();
4016
4017 let error = rt_secondary
4018 .create_entity_with_attachments(
4019 &token,
4020 "artifact",
4021 Some("visual_asset"),
4022 "must route through core",
4023 None,
4024 None,
4025 vec![],
4026 vec![NewAttachment {
4027 role: "content".to_string(),
4028 content_ref: content_ref.clone(),
4029 media_type: None,
4030 size_bytes: Some(13),
4031 }],
4032 )
4033 .await
4034 .expect_err("secondary attachment publication must fail closed");
4035 assert!(error.to_string().contains("canonical main backend"));
4036 assert!(rt_secondary
4037 .list_entities(&token, None, None, 10, 0)
4038 .await
4039 .unwrap()
4040 .is_empty());
4041 assert!(rt_secondary
4042 .core()
4043 .list_entities(&token, None, None, 10, 0)
4044 .await
4045 .unwrap()
4046 .is_empty());
4047 assert!(
4048 blob_store.exists(&content_ref).await.unwrap(),
4049 "refusal must not mutate the already-published object"
4050 );
4051 }
4052
4053 #[tokio::test]
4054 async fn list_embedding_models_returns_empty_when_table_absent() {
4055 let rt = KhiveRuntime::memory().expect("memory runtime");
4058 let records = rt
4059 .list_embedding_models(None)
4060 .await
4061 .expect("list ok on empty table");
4062 assert!(records.is_empty());
4063 }
4064
4065 #[tokio::test]
4066 async fn list_embedding_models_returns_row_after_insert() {
4067 use khive_storage::{SqlStatement, SqlValue};
4068
4069 let rt = KhiveRuntime::memory().expect("memory runtime");
4070 let sql = rt.sql();
4071
4072 let now = 1_000_000i64;
4073 let id = uuid::Uuid::new_v4();
4074 let canonical_key = b"test_engine:test-model-v1:v1:384".to_vec();
4075
4076 let mut writer = sql.writer().await.expect("writer");
4077 writer
4078 .execute(SqlStatement {
4079 sql: "INSERT INTO _embedding_models \
4080 (id, engine_name, model_id, key_version, dim, output_dim, status, \
4081 activated_at, superseded_at, superseded_by, canonical_key, created_at) \
4082 VALUES (?1, ?2, ?3, ?4, ?5, NULL, ?6, ?7, NULL, NULL, ?8, ?9)"
4083 .into(),
4084 params: vec![
4085 SqlValue::Blob(id.as_bytes().to_vec()),
4086 SqlValue::Text("test_engine".into()),
4087 SqlValue::Text("test-model-v1".into()),
4088 SqlValue::Text("v1".into()),
4089 SqlValue::Integer(384),
4090 SqlValue::Text("active".into()),
4091 SqlValue::Integer(now),
4092 SqlValue::Blob(canonical_key),
4093 SqlValue::Integer(now),
4094 ],
4095 label: None,
4096 })
4097 .await
4098 .expect("insert row");
4099 drop(writer);
4100
4101 let records = rt.list_embedding_models(None).await.expect("list ok");
4102 assert_eq!(records.len(), 1);
4103 assert_eq!(records[0].engine_name, "test_engine");
4104 assert_eq!(records[0].model_id, "test-model-v1");
4105 assert_eq!(records[0].key_version, "v1");
4106 assert_eq!(records[0].dimensions, 384);
4107 assert_eq!(records[0].status, "active");
4108
4109 let filtered = rt
4111 .list_embedding_models(Some("test_engine"))
4112 .await
4113 .expect("filter ok");
4114 assert_eq!(filtered.len(), 1);
4115
4116 let no_match = rt
4118 .list_embedding_models(Some("other_engine"))
4119 .await
4120 .expect("no-match ok");
4121 assert!(no_match.is_empty());
4122 }
4123
4124 #[test]
4125 fn named_vector_identity_rejects_ambiguous_or_unsafe_values() {
4126 assert!(NamedVectorIdentity::new("", "model", 4).is_err());
4127 assert!(NamedVectorIdentity::new("bad-key", "model", 4).is_err());
4128 assert!(NamedVectorIdentity::new("valid_key", " model", 4).is_err());
4129 assert!(NamedVectorIdentity::new("valid_key", "model", 0).is_err());
4130 assert!(NamedVectorIdentity::new("valid_key", "model", 8193).is_err());
4131 assert!(NamedVectorIdentity::new("k".repeat(128), "m".repeat(512), 4).is_ok());
4132 assert!(NamedVectorIdentity::new("k".repeat(129), "model", 4).is_err());
4133 assert!(NamedVectorIdentity::new("valid_key", "m".repeat(513), 4).is_err());
4134 assert_eq!(
4135 NamedVectorIdentity::new("valid_key", "model", 4)
4136 .expect("valid identity")
4137 .dimensions(),
4138 4
4139 );
4140 }
4141
4142 #[tokio::test]
4143 async fn named_vector_store_rejects_dimension_or_model_key_rebinding() {
4144 let rt = KhiveRuntime::memory().expect("memory runtime");
4145 let token = rt.authorize(Namespace::local()).expect("authorize");
4146 let original = NamedVectorIdentity::new("visual_contract", "model-a", 4).unwrap();
4147 rt.vectors_for_named_identity(&token, &original)
4148 .await
4149 .expect("create named vector store");
4150 let registered = rt
4151 .list_embedding_models(Some("visual_contract"))
4152 .await
4153 .expect("list model registry");
4154 assert!(registered.iter().any(|record| {
4155 record.model_id == "model-a"
4156 && record.key_version == "visual_contract"
4157 && record.dimensions == 4
4158 }));
4159 let wrong_dimensions = NamedVectorIdentity::new("visual_contract", "model-a", 5).unwrap();
4160 let Err(dimension_error) = rt
4161 .vectors_for_named_identity(&token, &wrong_dimensions)
4162 .await
4163 else {
4164 panic!("same key cannot change dimensions");
4165 };
4166 assert!(dimension_error.to_string().contains("dimensions"));
4167
4168 let wrong_model = NamedVectorIdentity::new("visual_contract", "model-b", 4).unwrap();
4169 let Err(model_error) = rt.vectors_for_named_identity(&token, &wrong_model).await else {
4170 panic!("same key cannot change model identity");
4171 };
4172 assert!(model_error.to_string().contains("already bound"));
4173 }
4174
4175 #[tokio::test]
4176 async fn repeated_named_vector_lookup_avoids_writer_acquisition() {
4177 let runtime = KhiveRuntime::memory().expect("memory runtime");
4178 let token = runtime.authorize(Namespace::local()).expect("authorize");
4179 let identity = NamedVectorIdentity::new("visual_cached", "model-a", 4).unwrap();
4180 runtime
4181 .vectors_for_named_identity(&token, &identity)
4182 .await
4183 .expect("first lookup validates and registers");
4184 let writer_before = runtime.backend().pool().writer_acquisition_snapshot();
4185
4186 runtime
4187 .clone()
4188 .vectors_for_named_identity(&token, &identity)
4189 .await
4190 .expect("clone reuses verified store");
4191 assert_eq!(
4192 runtime.backend().pool().writer_acquisition_snapshot(),
4193 writer_before,
4194 "repeated reads must not reach vector-table setup or model registration"
4195 );
4196 }
4197
4198 #[tokio::test]
4199 async fn core_projection_reuses_main_named_vector_cache() {
4200 let main_backend = migrated_memory_backend();
4201 let main =
4202 KhiveRuntime::from_backend(Arc::clone(&main_backend), RuntimeConfig::no_embeddings());
4203 let secondary = KhiveRuntime::from_backend(migrated_memory_backend(), secondary_config())
4204 .with_core_embedders_from(&main)
4205 .with_core_backend(Arc::clone(&main_backend));
4206 let core = secondary.core();
4207 let token = core.authorize(Namespace::local()).expect("authorize");
4208 let identity = NamedVectorIdentity::new("core_visual_cached", "model-a", 4).unwrap();
4209 core.vectors_for_named_identity(&token, &identity)
4210 .await
4211 .expect("first lookup validates on main");
4212 let writer_before = main_backend.pool().writer_acquisition_snapshot();
4213
4214 secondary
4215 .core()
4216 .vectors_for_named_identity(&token, &identity)
4217 .await
4218 .expect("new core projection reuses main store");
4219 assert_eq!(
4220 main_backend.pool().writer_acquisition_snapshot(),
4221 writer_before
4222 );
4223 }
4224
4225 #[tokio::test]
4226 async fn concurrent_named_vector_first_bind_has_one_immutable_winner() {
4227 let rt = KhiveRuntime::memory().expect("memory runtime");
4228 let token = rt.authorize(Namespace::local()).expect("authorize");
4229 let first = NamedVectorIdentity::new("visual_race", "model-a", 4).unwrap();
4230 let second = NamedVectorIdentity::new("visual_race", "model-b", 4).unwrap();
4231
4232 let (first_result, second_result) = tokio::join!(
4233 rt.vectors_for_named_identity(&token, &first),
4234 rt.vectors_for_named_identity(&token, &second),
4235 );
4236 assert_ne!(
4237 first_result.is_ok(),
4238 second_result.is_ok(),
4239 "the active engine_name uniqueness rule must select exactly one first binding"
4240 );
4241
4242 let (winner, loser) = if first_result.is_ok() {
4243 (&first, &second)
4244 } else {
4245 (&second, &first)
4246 };
4247 rt.vectors_for_named_identity(&token, winner)
4248 .await
4249 .expect("winning identity remains idempotent");
4250 let error = match rt.vectors_for_named_identity(&token, loser).await {
4251 Ok(_) => panic!("losing identity cannot rebind the empty table"),
4252 Err(error) => error,
4253 };
4254 assert!(error.to_string().contains("already bound"));
4255
4256 let registered = rt
4257 .list_embedding_models(Some("visual_race"))
4258 .await
4259 .expect("list race registry");
4260 assert_eq!(registered.len(), 1);
4261 assert_eq!(registered[0].model_id, winner.model_name());
4262 }
4263
4264 #[tokio::test]
4265 async fn named_vector_registry_keeps_immutable_revisions_active_together() {
4266 let rt = KhiveRuntime::memory().expect("memory runtime");
4267 let token = rt.authorize(Namespace::local()).expect("authorize");
4268 let first = NamedVectorIdentity::new("visual_revision_a", "visual-model", 4).unwrap();
4269 let second = NamedVectorIdentity::new("visual_revision_b", "visual-model", 4).unwrap();
4270
4271 rt.vectors_for_named_identity(&token, &first)
4272 .await
4273 .expect("open first immutable space");
4274 rt.vectors_for_named_identity(&token, &second)
4275 .await
4276 .expect("open second immutable space");
4277
4278 let registered = rt.list_embedding_models(None).await.expect("list registry");
4279 assert!(registered.iter().any(|record| {
4280 record.engine_name == "visual_revision_a"
4281 && record.model_id == "visual-model"
4282 && record.key_version == "visual_revision_a"
4283 && record.status == "active"
4284 }));
4285 assert!(registered.iter().any(|record| {
4286 record.engine_name == "visual_revision_b"
4287 && record.model_id == "visual-model"
4288 && record.key_version == "visual_revision_b"
4289 && record.status == "active"
4290 }));
4291 }
4292}